Files
quixos/orchestrator/server.ts
T
2026-06-24 22:24:38 -07:00

1755 lines
52 KiB
TypeScript

import fs from "node:fs";
import http from "node:http";
import process from "node:process";
import { WebSocketServer, type WebSocket } from "ws";
import {
BootRequestParamsSchema,
CallRequestParamsSchema,
ClientMessageSchema,
EventSubscribeRequestParamsSchema,
EventUnsubscribeRequestParamsSchema,
RuntimeErrorRequestParamsSchema,
SubscriptionReadyRequestParamsSchema,
SubscriptionPublishRequestParamsSchema,
TargetContextOpenRequestParamsSchema,
ValueGetRequestParamsSchema,
ValueUnwatchRequestParamsSchema,
ValueWatchRequestParamsSchema,
parseJson,
sendServerMessage,
} from "@quixos/package-runtime";
import {
CONTROL_PLANE_ROUTES,
parseControlPlaneCheckRequest,
parseControlPlaneClearStateRequest,
parseControlPlaneCutoverRequest,
parseControlPlaneDescribeRequest,
parseControlPlaneDeviceAssignRequest,
parseControlPlaneDeviceIdentifyRequest,
parseControlPlaneDeviceRemoveRequest,
parseControlPlanePublicAuthGrantRequest,
parseControlPlanePublicAuthGrantRevokeRequest,
parseControlPlaneSurfaceSessionRequest,
parseControlPlaneErrorsRequest,
parseControlPlaneLogsRequest,
parseControlPlaneReplCloseRequest,
parseControlPlaneReplEvalRequest,
parseControlPlaneRunRequest,
type ControlPlaneCutoverResponse,
type ControlPlaneDeviceTarget,
type ControlPlaneDevicesResponse,
type ControlPlaneErrorsResponse,
type ControlPlaneLogsResponse,
type ControlPlanePublicAuthGrantRevokeResponse,
type ControlPlanePublicAuthGrantResponse,
type ControlPlaneReplLogEntry,
type ControlPlaneRunResponse,
type ControlPlaneSurfaceSessionResponse,
type ControlPlaneStatusResponse,
} from "@quixos/control-plane-protocol";
import {
CONTROL_PLANE_HOST,
CONTROL_PLANE_PORT,
CONTROL_PLANE_SECRET,
CONTROL_PLANE_SOCKET,
DEVICE_IDENTIFY_DURATION_MS,
ORCHESTRATOR_SHUTDOWN_TIMEOUT_MS,
} from "./src/config.js";
import { createDeviceRegistry } from "./src/device-registry.js";
import { attachBrowserControlSocketHandlers } from "./src/browser-control-socket.js";
import { attachDeviceControlSocketHandlers } from "./src/device-control-socket.js";
import {
applyControlPlaneCorsHeaders,
isAuthorizedControlPlaneRequest,
readJsonBody,
sendJson,
sendWebSocketJson,
} from "./src/http.js";
import {
createPackageResolver,
} from "./src/package-resolver.js";
import { createPackageOperations } from "./src/package-operations.js";
import {
evaluatePackageSchemaDescription,
listPackageCheckResults,
listPackageSchemaDescriptions,
runRemotePackageCheck,
warmWorkspaceHeadChecks,
} from "./src/package-inspection.js";
import {
createPackageLogStore,
} from "./src/package-logs.js";
import {
createPublicAuthGrant,
createPublicSessionToken,
loadPublicAuthGrants,
normalizePublicAccessTarget,
revokePublicAuthGrant,
} from "./src/public-auth-grants.js";
import {
createReplManager,
listReplSessionsForControlPlane,
} from "./src/repl-manager.js";
import { handlePublicControlPlaneRoute } from "./src/public-control-routes.js";
import { createRunningPackageManager } from "./src/running-packages.js";
import { createStateContextRegistry } from "./src/state-contexts.js";
import { createSubscriptionManager } from "./src/subscriptions.js";
import { listWorkspacePackagesForControlPlane } from "./src/workspace-packages.js";
import type {
PackageErrorBufferEntry,
RunningProcess,
StateContextId,
} from "./src/types.js";
import {
normalizePackageName,
nowIso,
sameJsonValue as sameDeviceTarget,
} from "./src/util.js";
const browserControlSockets = new Set<WebSocket>();
let browserControlWss: WebSocketServer | null = null;
let deviceControlWss: WebSocketServer | null = null;
let controlPlaneServer: http.Server | null = null;
const broadcastControlPlaneStatusInvalidated = (
reason:
| "state-context-opened"
| "state-context-closed"
| "running-packages-changed"
| "repl-sessions-changed"
| "runtime-errors-changed"
| "devices-changed",
) => {
for (const socket of browserControlSockets) {
if (socket.readyState !== 1) {
continue;
}
sendWebSocketJson(socket, {
type: "status-invalidated",
reason,
});
}
};
const {
acknowledgePackageErrors,
appendPackageLogEntry,
getPackageErrorSummary,
listPackageErrors,
listPackageLogs,
recordRuntimeError,
} = createPackageLogStore(() =>
broadcastControlPlaneStatusInvalidated("runtime-errors-changed"),
);
const {
performReplClose,
performReplEval,
stopAllReplSessions,
} = createReplManager(() =>
broadcastControlPlaneStatusInvalidated("repl-sessions-changed"),
);
const { resolvePinned, setDesiredPackageTarget } = createPackageResolver();
let getRunningProcessForStateContexts: (
packageName: string,
) => RunningProcess | undefined = () => undefined;
let listRunningProcessesForStateContexts: () => Iterable<RunningProcess> =
() => [];
const {
getActiveStateContextIds,
getPackageStateClearTargets,
getRunningPackageContextDepth,
getStateContext,
getStateContextDepth,
getOrCreateBootRootStateContext,
isRootStateContext,
listOpenStateContextsForControlPlane,
listRootStateContextIds,
markStateContextActive,
mintRootStateContext,
recordStateContextUsage,
resolveTargetStateContext,
updatePackageStateContextRefs,
} = createStateContextRegistry({
getRunningProcess: (packageName) => getRunningProcessForStateContexts(packageName),
listRunningProcesses: () => listRunningProcessesForStateContexts(),
});
const runningPackages = createRunningPackageManager({
appendPackageLogEntry,
getStateContext,
getOrCreateBootRootStateContext,
getStateContextDepth,
isRootStateContext,
markStateContextActive,
mintRootStateContext,
onRunningPackagesChanged: () =>
broadcastControlPlaneStatusInvalidated("running-packages-changed"),
onStateContextClosed: () =>
broadcastControlPlaneStatusInvalidated("state-context-closed"),
onStateContextOpened: () =>
broadcastControlPlaneStatusInvalidated("state-context-opened"),
recordStateContextUsage,
resolvePinned,
});
getRunningProcessForStateContexts = runningPackages.getRunningProcess;
listRunningProcessesForStateContexts = runningPackages.listRunningProcesses;
const {
attachSocket: attachPackageSocket,
detachSocket: detachPackageSocket,
ensureContextClose,
ensureContextOpen,
ensureRunning,
ensureRunningForTargetRequest,
findBySecret: findRunningPackageBySecret,
getRunningProcess,
getSocketRunningProcess,
listRunningProcesses,
openRootStateContext,
sendRequest,
stopAllRunningPackages,
stopPinnedPackage,
triggerRootBoot,
} = runningPackages;
const {
attachSubscriptionToPublisher,
detachPublisherSubscriptions,
detachSubscriberSocket,
getOrCreateLogicalSubscription,
getSubscription,
markSubscriptionReady,
publishSubscriptionData,
reattachSubscriptionsForPublisher,
removeSubscriberSocket,
removeSubscription,
setSubscriptionSubscriber,
} = createSubscriptionManager({
ensureContextOpen,
ensureRunning,
getRunningProcess,
sendRequest,
});
const validateDeviceTarget = (target: ControlPlaneDeviceTarget | null) => {
if (!target) {
return null;
}
const stateContext = getStateContext(target.stateContextId);
if (stateContext.packageName !== target.packageName) {
throw new Error(
`State context ${target.stateContextId} belongs to ${stateContext.packageName}, not ${target.packageName}`,
);
}
return target;
};
const {
applyConfiguredDeviceAssignments,
deleteDevice,
deviceSnapshotForControlPlane,
findDeviceForControlPlane,
getOrCreateDevice,
getSocketDeviceId,
isDeviceConnected,
listDevicesForControlPlane,
loadDeviceState,
markDeviceSeenThisSession,
registerDeviceSocket,
requireDevice,
saveDeviceAssignmentsConfig,
saveDeviceState,
sendDeviceRemoved,
sendDeviceState,
syncConfiguredDeviceAssignmentsFromDisk,
tryValidateDeviceTarget,
unregisterDeviceSocket,
} = createDeviceRegistry({
broadcastDevicesChanged: () =>
broadcastControlPlaneStatusInvalidated("devices-changed"),
validateDeviceTarget,
});
const {
getRunningPackageStatus,
performClearState,
performCutover,
} = createPackageOperations({
applyConfiguredDeviceAssignments,
ensureContextOpen,
ensureRunning,
getActiveStateContextIds,
getPackageErrorSummary,
getPackageStateClearTargets,
getRunningPackageContextDepth,
getRunningProcess,
getStateContextDepth,
listRootStateContextIds,
reattachSubscriptionsForPublisher,
resolvePinned,
setDesiredPackageTarget,
stopPinnedPackage,
updatePackageStateContextRefs,
});
const closeWebSocketServer = async (server: WebSocketServer | null) => {
if (!server) {
return;
}
for (const client of server.clients) {
try {
client.close();
client.terminate();
} catch {
// Ignore shutdown races for already-closed sockets.
}
}
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
};
const closeHttpServer = async (server: http.Server | null) => {
if (!server) {
return;
}
server.closeAllConnections?.();
server.closeIdleConnections?.();
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
};
let shutdownInFlight: Promise<void> | null = null;
const installShutdownHandlers = () => {
let shutdownTimer: NodeJS.Timeout | null = null;
const handleSignal = (signal: NodeJS.Signals) => {
if (shutdownInFlight) {
console.warn(`Received ${signal} during shutdown; forcing exit.`);
process.exit(signal === "SIGINT" ? 130 : 0);
}
shutdownTimer = setTimeout(() => {
console.warn(`Timed out shutting down after ${signal}; forcing exit.`);
process.exit(signal === "SIGINT" ? 130 : 0);
}, ORCHESTRATOR_SHUTDOWN_TIMEOUT_MS);
shutdownInFlight = Promise.all([
closeWebSocketServer(browserControlWss),
closeWebSocketServer(deviceControlWss),
closeWebSocketServer(wss),
closeHttpServer(controlPlaneServer),
])
.catch((error) => {
console.warn(
`Failed to close orchestrator servers during ${signal}: ${
error instanceof Error ? error.message : String(error)
}`,
);
})
.then(async () => {
await stopAllReplSessions();
await stopAllRunningPackages();
})
.catch((error) => {
console.warn(
`Failed to stop running processes during ${signal}: ${
error instanceof Error ? error.message : String(error)
}`,
);
})
.finally(() => {
if (shutdownTimer) {
clearTimeout(shutdownTimer);
}
process.exit(signal === "SIGINT" ? 130 : 0);
});
};
process.on("SIGINT", () => handleSignal("SIGINT"));
process.on("SIGTERM", () => handleSignal("SIGTERM"));
};
const cleanupSocketState = (socket: WebSocket) => {
detachSubscriberSocket(socket);
const entry = detachPackageSocket(socket);
if (entry) {
detachPublisherSubscriptions(entry);
}
};
const wss = new WebSocketServer({ port: 6245 });
installShutdownHandlers();
wss.on("connection", (socket) => {
const cleanupSocket = () => {
cleanupSocketState(socket);
};
socket.on("close", cleanupSocket);
socket.on("error", cleanupSocket);
socket.on("message", (data) => {
const parsed = parseJson(data.toString());
if (!parsed.ok) {
return;
}
const messageResult = ClientMessageSchema.safeParse(parsed.value);
if (!messageResult.success) {
return;
}
const message = messageResult.data;
if (message.type === "auth") {
const entry = findRunningPackageBySecret(message.secret);
if (!entry) {
socket.close();
return;
}
attachPackageSocket(socket, entry);
sendServerMessage(socket, { type: "auth-ack" });
console.log(`Server running for ${entry.label}`);
void reattachSubscriptionsForPublisher(entry).catch((error) => {
console.warn(
`Failed to reattach subscriptions for ${entry.packageName}: ${
error instanceof Error ? error.message : String(error)
}`,
);
});
triggerRootBoot();
return;
}
if (message.type === "request" && message.method === "call") {
const parsedParams = CallRequestParamsSchema.safeParse(message.params);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid call params",
});
return;
}
const { target, targetPackageName, functionName, params } =
parsedParams.data;
if (!target) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Missing target",
});
return;
}
const caller = getSocketRunningProcess(socket);
void (async () => {
try {
const entry = await ensureRunning(target);
const normalizedTargetPackageName = targetPackageName
? normalizePackageName(targetPackageName)
: entry.packageName;
if (normalizedTargetPackageName !== entry.packageName) {
throw new Error(
`Target package mismatch: expected ${entry.packageName}, got ${normalizedTargetPackageName}`,
);
}
const stateContext = resolveTargetStateContext({
targetPackageName: entry.packageName,
targetPackageRef: entry.flakeRef,
callerStateContextId: parsedParams.data.callerStateContextId,
explicitStateContextId: parsedParams.data.stateContextId,
contextNamespace: parsedParams.data.contextNamespace,
sourcePackageRef: caller?.flakeRef,
sourcePackageName: caller?.packageName,
operation: "call",
functionName,
});
await ensureContextOpen(entry, stateContext.id);
const result = await sendRequest(entry, "call", {
functionName,
params,
stateContextId: stateContext.id,
});
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
result,
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (message.type === "request" && message.method === "boot") {
const parsedParams = BootRequestParamsSchema.safeParse(message.params);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid boot params",
});
return;
}
void (async () => {
try {
await ensureRunning(parsedParams.data.target);
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (message.type === "request" && message.method === "target-context-open") {
const parsedParams = TargetContextOpenRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid target context open params",
});
return;
}
const { target, targetPackageName } = parsedParams.data;
if (!target) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Missing target",
});
return;
}
const caller = getSocketRunningProcess(socket);
void (async () => {
try {
const entry = await ensureRunning(target);
const normalizedTargetPackageName = targetPackageName
? normalizePackageName(targetPackageName)
: entry.packageName;
if (normalizedTargetPackageName !== entry.packageName) {
throw new Error(
`Target package mismatch: expected ${entry.packageName}, got ${normalizedTargetPackageName}`,
);
}
const stateContext = resolveTargetStateContext({
targetPackageName: entry.packageName,
targetPackageRef: entry.flakeRef,
callerStateContextId: parsedParams.data.callerStateContextId,
explicitStateContextId: parsedParams.data.stateContextId,
contextNamespace: parsedParams.data.contextNamespace,
sourcePackageRef: caller?.flakeRef,
sourcePackageName: caller?.packageName,
operation: "use-package",
});
await ensureContextOpen(entry, stateContext.id);
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
result: {
stateContextId: stateContext.id,
},
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (message.type === "request" && message.method === "runtime-error") {
const parsedParams = RuntimeErrorRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid runtime error params",
});
return;
}
const publisher = getSocketRunningProcess(socket);
if (!publisher) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Unknown publisher",
});
return;
}
recordRuntimeError(publisher, parsedParams.data);
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
return;
}
if (
message.type === "request" &&
message.method === "event-subscribe"
) {
const parsedParams = EventSubscribeRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid event subscribe params",
});
return;
}
const { target, targetPackageName, eventName, params } = parsedParams.data;
if (!target) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Missing target",
});
return;
}
const caller = getSocketRunningProcess(socket);
void (async () => {
try {
const entry = await ensureRunning(target);
const normalizedTargetPackageName = targetPackageName
? normalizePackageName(targetPackageName)
: entry.packageName;
if (normalizedTargetPackageName !== entry.packageName) {
throw new Error(
`Target package mismatch: expected ${entry.packageName}, got ${normalizedTargetPackageName}`,
);
}
const stateContext = resolveTargetStateContext({
targetPackageName: entry.packageName,
targetPackageRef: entry.flakeRef,
callerStateContextId: parsedParams.data.callerStateContextId,
explicitStateContextId: parsedParams.data.stateContextId,
contextNamespace: parsedParams.data.contextNamespace,
sourcePackageRef: caller?.flakeRef,
sourcePackageName: caller?.packageName,
operation: "event-subscribe",
resourceKind: "event",
resourceName: eventName,
});
const subscription = getOrCreateLogicalSubscription({
kind: "event",
targetPackageName: entry.packageName,
targetStateContextId: stateContext.id,
resourceName: eventName,
resourceParams: params,
subscriptionNamespace:
parsedParams.data.subscriptionNamespace ?? null,
subscriberPackageName: caller?.packageName,
subscriberStateContextId: parsedParams.data.callerStateContextId,
});
setSubscriptionSubscriber(subscription, socket);
await attachSubscriptionToPublisher(subscription, entry);
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
result: { subscriptionId: subscription.id },
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (
message.type === "request" &&
message.method === "subscription-ready"
) {
const parsedParams = SubscriptionReadyRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid subscription ready params",
});
return;
}
const subscription = getSubscription(
parsedParams.data.subscriptionId,
);
if (!subscription) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
return;
}
if (subscription.subscriber !== socket) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Not subscription owner",
});
return;
}
void markSubscriptionReady(subscription)
.then(() => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
})
.catch((error) => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
});
return;
}
if (
message.type === "request" &&
message.method === "event-unsubscribe"
) {
const parsedParams = EventUnsubscribeRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid event unsubscribe params",
});
return;
}
const { subscriptionId } = parsedParams.data;
const subscription = getSubscription(subscriptionId);
if (!subscription) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
return;
}
if (subscription.subscriber !== socket) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Not subscription owner",
});
return;
}
void removeSubscription(subscriptionId, true).then(() => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
});
return;
}
if (message.type === "request" && message.method === "value-get") {
const parsedParams = ValueGetRequestParamsSchema.safeParse(message.params);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid value get params",
});
return;
}
const { target, targetPackageName, valueName, params } = parsedParams.data;
if (!target) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Missing target",
});
return;
}
const caller = getSocketRunningProcess(socket);
void (async () => {
try {
const entry = await ensureRunning(target);
const normalizedTargetPackageName = targetPackageName
? normalizePackageName(targetPackageName)
: entry.packageName;
if (normalizedTargetPackageName !== entry.packageName) {
throw new Error(
`Target package mismatch: expected ${entry.packageName}, got ${normalizedTargetPackageName}`,
);
}
const stateContext = resolveTargetStateContext({
targetPackageName: entry.packageName,
targetPackageRef: entry.flakeRef,
callerStateContextId: parsedParams.data.callerStateContextId,
explicitStateContextId: parsedParams.data.stateContextId,
contextNamespace: parsedParams.data.contextNamespace,
sourcePackageRef: caller?.flakeRef,
sourcePackageName: caller?.packageName,
operation: "value-get",
resourceKind: "value",
resourceName: valueName,
});
await ensureContextOpen(entry, stateContext.id);
const result = await sendRequest(entry, "value-get", {
valueName,
params,
stateContextId: stateContext.id,
});
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
result,
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (message.type === "request" && message.method === "value-watch") {
const parsedParams = ValueWatchRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid value watch params",
});
return;
}
const { target, targetPackageName, valueName, params } = parsedParams.data;
if (!target) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Missing target",
});
return;
}
const caller = getSocketRunningProcess(socket);
void (async () => {
try {
const entry = await ensureRunning(target);
const normalizedTargetPackageName = targetPackageName
? normalizePackageName(targetPackageName)
: entry.packageName;
if (normalizedTargetPackageName !== entry.packageName) {
throw new Error(
`Target package mismatch: expected ${entry.packageName}, got ${normalizedTargetPackageName}`,
);
}
const stateContext = resolveTargetStateContext({
targetPackageName: entry.packageName,
targetPackageRef: entry.flakeRef,
callerStateContextId: parsedParams.data.callerStateContextId,
explicitStateContextId: parsedParams.data.stateContextId,
contextNamespace: parsedParams.data.contextNamespace,
sourcePackageRef: caller?.flakeRef,
sourcePackageName: caller?.packageName,
operation: "value-watch",
resourceKind: "value",
resourceName: valueName,
});
const subscription = getOrCreateLogicalSubscription({
kind: "value",
targetPackageName: entry.packageName,
targetStateContextId: stateContext.id,
resourceName: valueName,
resourceParams: params,
subscriptionNamespace:
parsedParams.data.subscriptionNamespace ?? null,
subscriberPackageName: caller?.packageName,
subscriberStateContextId: parsedParams.data.callerStateContextId,
});
setSubscriptionSubscriber(subscription, socket);
await attachSubscriptionToPublisher(subscription, entry);
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
result: { subscriptionId: subscription.id },
});
} catch (error) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: (error as Error).message ?? "Unknown error",
});
}
})();
return;
}
if (message.type === "request" && message.method === "value-unwatch") {
const parsedParams = ValueUnwatchRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid value unwatch params",
});
return;
}
const { subscriptionId } = parsedParams.data;
const subscription = getSubscription(subscriptionId);
if (!subscription) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
return;
}
if (subscription.subscriber !== socket) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Not subscription owner",
});
return;
}
void removeSubscription(subscriptionId, true).then(() => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
});
return;
}
if (
message.type === "request" &&
message.method === "subscription-publish"
) {
const parsedParams = SubscriptionPublishRequestParamsSchema.safeParse(
message.params,
);
if (!parsedParams.success) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Invalid subscription publish params",
});
return;
}
const publisher = getSocketRunningProcess(socket);
if (!publisher) {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: "Unknown publisher",
});
return;
}
const { subscriptionId: publisherSubscriptionId, data } = parsedParams.data;
void publishSubscriptionData(publisher, publisherSubscriptionId, data)
.then(() => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: true,
});
})
.catch((error) => {
sendServerMessage(socket, {
type: "response",
requestId: message.requestId,
ok: false,
error: error instanceof Error ? error.message : String(error),
});
});
return;
}
if (message.type === "response") {
const requestId = message.requestId;
for (const entry of listRunningProcesses()) {
const pending = entry.pendingRequests.get(requestId);
if (pending) {
entry.pendingRequests.delete(requestId);
clearTimeout(pending.timeout);
if (message.ok === false) {
pending.reject(new Error(message.error ?? "Unknown error"));
} else {
pending.resolve(message.result);
}
return;
}
}
}
});
});
loadDeviceState();
loadPublicAuthGrants();
syncConfiguredDeviceAssignmentsFromDisk();
if (CONTROL_PLANE_SECRET) {
browserControlWss = new WebSocketServer({ noServer: true });
deviceControlWss = new WebSocketServer({ noServer: true });
attachBrowserControlSocketHandlers({
attachSubscriptionToPublisher,
browserControlSockets,
browserControlWss,
controlSecret: CONTROL_PLANE_SECRET,
ensureContextClose,
ensureContextOpen,
ensureRunning,
ensureRunningForTargetRequest,
getOrCreateLogicalSubscription,
getStateContext,
getSubscription,
markSubscriptionReady,
openRootStateContext,
removeSubscriberSocket,
removeSubscription,
resolveTargetStateContext,
sendRequest,
setSubscriptionSubscriber,
});
attachDeviceControlSocketHandlers({
broadcastDevicesChanged: () =>
broadcastControlPlaneStatusInvalidated("devices-changed"),
deviceControlWss,
deviceSnapshotForControlPlane,
getOrCreateDevice,
getSocketDeviceId,
markDeviceSeenThisSession,
registerDeviceSocket,
requireDevice,
saveDeviceState,
sendDeviceState,
tryValidateDeviceTarget,
unregisterDeviceSocket,
});
controlPlaneServer = http.createServer(async (request, response) => {
applyControlPlaneCorsHeaders(response);
if (request.method === "OPTIONS") {
response.statusCode = 204;
response.end();
return;
}
if (!request.url) {
sendJson(response, 404, { error: "Not found" });
return;
}
if (await handlePublicControlPlaneRoute({
broadcastDevicesChanged: () =>
broadcastControlPlaneStatusInvalidated("devices-changed"),
findDeviceForControlPlane,
getOrCreateDevice,
isDeviceConnected,
markDeviceSeenThisSession,
request,
requireDevice,
response,
saveDeviceState,
sendDeviceState,
tryValidateDeviceTarget,
})) {
return;
}
if (!isAuthorizedControlPlaneRequest(request)) {
sendJson(response, 401, { error: "Unauthorized" });
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.publicAuthGrant
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlanePublicAuthGrantRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid public auth grant request" });
return;
}
const result: ControlPlanePublicAuthGrantResponse = createPublicAuthGrant(
params.target,
);
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.publicAuthGrantRevoke
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlanePublicAuthGrantRevokeRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid public auth grant revoke request" });
return;
}
const result: ControlPlanePublicAuthGrantRevokeResponse = {
revoked: revokePublicAuthGrant(params.authGrant),
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.surfaceSession
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneSurfaceSessionRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid surface session request" });
return;
}
const normalized = normalizePublicAccessTarget(params.target);
if (normalized.kind === "device") {
sendJson(response, 400, { error: "Invalid surface session request" });
return;
}
const session = createPublicSessionToken(normalized);
const result: ControlPlaneSurfaceSessionResponse = {
...session,
target: normalized,
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "GET" &&
request.url === CONTROL_PLANE_ROUTES.devices
) {
syncConfiguredDeviceAssignmentsFromDisk();
const result: ControlPlaneDevicesResponse = {
devices: listDevicesForControlPlane(),
};
sendJson(response, 200, result);
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.deviceAssign
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDeviceAssignRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid device assign request" });
return;
}
const device = requireDevice(params.deviceId);
const validatedTarget = tryValidateDeviceTarget(params.target);
if (!validatedTarget.ok) {
sendJson(response, 400, { error: validatedTarget.error });
return;
}
let changed = false;
if (!sameDeviceTarget(device.desiredAssignment, validatedTarget.target)) {
device.desiredAssignment = validatedTarget.target;
changed = true;
}
if (changed) {
saveDeviceAssignmentsConfig();
broadcastControlPlaneStatusInvalidated("devices-changed");
sendDeviceState(device.deviceId);
}
sendJson(response, 200, {
device: findDeviceForControlPlane(device.deviceId),
});
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.deviceAssignLive
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDeviceAssignRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid live device assign request" });
return;
}
const device = requireDevice(params.deviceId);
const validatedTarget = tryValidateDeviceTarget(params.target);
if (!validatedTarget.ok) {
sendJson(response, 400, { error: validatedTarget.error });
return;
}
let changed = false;
if (!sameDeviceTarget(device.assignment, validatedTarget.target)) {
device.assignment = validatedTarget.target;
changed = true;
}
if (!sameDeviceTarget(device.desiredAssignment, validatedTarget.target)) {
device.desiredAssignment = validatedTarget.target;
changed = true;
saveDeviceAssignmentsConfig();
}
if (changed) {
saveDeviceState();
broadcastControlPlaneStatusInvalidated("devices-changed");
sendDeviceState(device.deviceId);
}
sendJson(response, 200, {
device: findDeviceForControlPlane(device.deviceId),
});
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.deviceIdentify
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDeviceIdentifyRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid device identify request" });
return;
}
const device = requireDevice(params.deviceId);
device.identifyUntil = new Date(
Date.now() + (params.durationMs ?? DEVICE_IDENTIFY_DURATION_MS),
).toISOString();
saveDeviceState();
broadcastControlPlaneStatusInvalidated("devices-changed");
sendDeviceState(device.deviceId);
sendJson(response, 200, {
device: findDeviceForControlPlane(device.deviceId),
});
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.deviceStopIdentify
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDeviceRemoveRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid stop-identify request" });
return;
}
const device = requireDevice(params.deviceId);
if (device.identifyUntil !== null) {
device.identifyUntil = null;
saveDeviceState();
broadcastControlPlaneStatusInvalidated("devices-changed");
sendDeviceState(device.deviceId);
}
sendJson(response, 200, {
device: findDeviceForControlPlane(device.deviceId) ?? null,
});
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.deviceRemove
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDeviceRemoveRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid remove-device request" });
return;
}
requireDevice(params.deviceId);
sendDeviceRemoved(params.deviceId);
deleteDevice(params.deviceId);
sendJson(response, 200, { removed: true, deviceId: params.deviceId });
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.run
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneRunRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid run request" });
return;
}
const stateContextId = await openRootStateContext(params.target);
const entry = await ensureRunning(params.target);
const packageStatus = await getRunningPackageStatus(entry);
const result: ControlPlaneRunResponse = {
stateContextId,
package: packageStatus,
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "GET" &&
request.url === CONTROL_PLANE_ROUTES.status
) {
try {
syncConfiguredDeviceAssignmentsFromDisk();
warmWorkspaceHeadChecks(listWorkspacePackagesForControlPlane());
const packages = await Promise.all(
[...listRunningProcesses()]
.sort((left, right) => left.label.localeCompare(right.label))
.map((entry) => getRunningPackageStatus(entry)),
);
const result: ControlPlaneStatusResponse = {
generatedAt: nowIso(),
workspacePackages: listWorkspacePackagesForControlPlane(),
stateContexts: listOpenStateContextsForControlPlane(),
replSessions: listReplSessionsForControlPlane(),
devices: listDevicesForControlPlane(),
packages,
checks: listPackageCheckResults(),
descriptions: listPackageSchemaDescriptions(),
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.check
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneCheckRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid check request" });
return;
}
const result = await runRemotePackageCheck(
normalizePackageName(params.packageName),
params.rev,
);
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.describe
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneDescribeRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid describe request" });
return;
}
const result = await evaluatePackageSchemaDescription(
normalizePackageName(params.packageName),
params.rev,
);
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.logs
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneLogsRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid logs request" });
return;
}
const result: ControlPlaneLogsResponse = {
packageName: params.packageName,
entries: listPackageLogs(params),
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.errors
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneErrorsRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid errors request" });
return;
}
const entries = listPackageErrors(params);
const acknowledgedCount = params.acknowledge
? acknowledgePackageErrors(entries)
: 0;
const result: ControlPlaneErrorsResponse = {
entries,
acknowledgedCount,
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.replEval
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneReplEvalRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid repl eval request" });
return;
}
const result = await performReplEval({
sessionId: params.sessionId,
code: params.code,
autoClose: params.autoClose ?? false,
defaultStateContextId: params.defaultStateContextId ?? null,
});
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.replClose
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneReplCloseRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid repl close request" });
return;
}
const result = await performReplClose(params.sessionId);
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.cutover
) {
try {
syncConfiguredDeviceAssignmentsFromDisk();
const body = await readJsonBody(request);
const params = parseControlPlaneCutoverRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid cutover request" });
return;
}
const packages = await performCutover(params.targets);
const result: ControlPlaneCutoverResponse = {
packages,
};
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
if (
request.method === "POST" &&
request.url === CONTROL_PLANE_ROUTES.clearState
) {
try {
const body = await readJsonBody(request);
const params = parseControlPlaneClearStateRequest(body);
if (!params) {
sendJson(response, 400, { error: "Invalid clear-state request" });
return;
}
const result = await performClearState({
packageName: params.packageName,
recursive: params.recursive ?? false,
});
sendJson(response, 200, result);
} catch (error) {
sendJson(response, 500, {
error: error instanceof Error ? error.message : String(error),
});
}
return;
}
sendJson(response, 404, { error: "Not found" });
});
const onControlPlaneListening = () => {
if (CONTROL_PLANE_SOCKET) {
try {
fs.chmodSync(CONTROL_PLANE_SOCKET, 0o660);
} catch {
// Best effort; systemd/Caddy can still proxy if permissions already fit.
}
console.log(`Control plane listening on unix:${CONTROL_PLANE_SOCKET}`);
} else {
console.log(
`Control plane listening on http://${CONTROL_PLANE_HOST}:${CONTROL_PLANE_PORT}`,
);
}
warmWorkspaceHeadChecks(listWorkspacePackagesForControlPlane());
triggerRootBoot();
};
if (CONTROL_PLANE_SOCKET) {
try {
fs.rmSync(CONTROL_PLANE_SOCKET, { force: true });
} catch {
// If removal fails, listen() will surface the useful error.
}
controlPlaneServer.listen(CONTROL_PLANE_SOCKET, onControlPlaneListening);
} else {
controlPlaneServer.listen(
CONTROL_PLANE_PORT,
CONTROL_PLANE_HOST,
onControlPlaneListening,
);
}
controlPlaneServer.on("upgrade", (request, socket, head) => {
const activeBrowserControlWss = browserControlWss;
const activeDeviceControlWss = deviceControlWss;
if (!activeBrowserControlWss || !activeDeviceControlWss) {
socket.destroy();
return;
}
if (!request.url) {
socket.destroy();
return;
}
const url = new URL(request.url, "http://localhost");
if (url.pathname !== CONTROL_PLANE_ROUTES.browserSocket) {
if (url.pathname !== CONTROL_PLANE_ROUTES.deviceSocket) {
socket.destroy();
return;
}
activeDeviceControlWss.handleUpgrade(request, socket, head, (ws) => {
activeDeviceControlWss.emit("connection", ws, request);
});
return;
}
activeBrowserControlWss.handleUpgrade(request, socket, head, (ws) => {
activeBrowserControlWss.emit("connection", ws, request);
});
});
} else {
console.log("QUIXOS_CONTROL_SECRET is not set; control plane disabled");
}