import http from "node:http"; import { readFileSync } from "node:fs"; import { createInvocationRegistry } from "./invocations.js"; import { isObjectReference, referenceFromWire, referenceToWire, assertReferenceFree, } from "./references.js"; export * from "./bindings.js"; export * from "./queries.js"; export { relationshipMap, relationshipList, relationshipSet } from "./relationships.js"; import { AsyncLocalStorage } from "node:async_hooks"; import { createHmac, randomBytes, randomUUID, timingSafeEqual } from "node:crypto"; export { createMigrationContext, migrationObjectId, serveMigration, } from "./migration.js"; import { create, equals } from "@bufbuild/protobuf"; import { Code, ConnectError, createClient } from "@connectrpc/connect"; import { connectNodeAdapter, createConnectTransport } from "@connectrpc/connect-node"; import { CaminoService, CrdtValueSchema, ListValueSchema, NullValueSchema, ObjectValueSchema, RefValueSchema, ValueSchema, } from "./camino/api_pb.js"; import { OrchestratorRuntime } from "./quixos/orch_pb.js"; import { DerivedDependencySchema, HandshakeResponseSchema, InvokeResponseSchema, PackageRuntime, WatchEventSchema, } from "./quixos/runtime_pb.js"; import { CapabilityRefSchema } from "./quixos/refs_pb.js"; const dependencyKey = (dependency) => `${dependency.kind}:${dependency.objectId}:${dependency.attachmentId}:${dependency.kind === "edge" ? dependency.projectionId : ""}`; const dependencyScope = new AsyncLocalStorage(); const recordDependency = async (dependency) => { const scope = dependencyScope.getStore(); if (!scope) return; const key = `${dependency.kind}:${dependency.objectId}:${dependency.attachmentId}:${dependency.kind === "edge" ? dependency.projectionId : ""}`; scope.dependencies.set(key, dependency); await scope.observe?.(dependency); }; const bytesToBase64 = (value) => Buffer.from(value).toString("base64"); const base64ToBytes = (value) => Buffer.from(value, "base64"); const isRecord = (value) => Boolean(value) && typeof value === "object" && !Array.isArray(value); const isWrappedValue = (value) => isRecord(value) && "$quixosValue" in value && isRecord(value.$quixosValue) && value.$quixosValue.$typeName === "camino.Value"; export const objectRef = (reference) => { referenceToWire(reference); return reference; }; export const liveValue = (value) => ({ $quixosValue: value }); export const jsToProtoValue = (value) => { if (isObjectReference(value)) return create(ValueSchema, { kind: { case: "refValue", value: create(RefValueSchema, { objectId: referenceToWire(value) }) }, }); if (isWrappedValue(value)) return value.$quixosValue; if (value === null || value === undefined) { return create(ValueSchema, { kind: { case: "nullValue", value: create(NullValueSchema, {}) } }); } if (typeof value === "boolean") return create(ValueSchema, { kind: { case: "boolValue", value } }); if (typeof value === "number") return create(ValueSchema, { kind: { case: "numberValue", value } }); if (typeof value === "bigint") return create(ValueSchema, { kind: { case: "integerValue", value: String(value) } }); if (typeof value === "string") return create(ValueSchema, { kind: { case: "stringValue", value } }); if (value instanceof Uint8Array) return create(ValueSchema, { kind: { case: "bytesValue", value } }); if (Array.isArray(value)) { return create(ValueSchema, { kind: { case: "listValue", value: create(ListValueSchema, { values: value.map(jsToProtoValue) }) }, }); } if (isRecord(value) && "$quixosRef" in value) throw new Error("Raw ID wrappers are not object references"); if (isRecord(value) && typeof value.$quixosCrdtType === "string" && typeof value.$quixosCrdtPayload === "string") { return create(ValueSchema, { kind: { case: "crdtValue", value: create(CrdtValueSchema, { type: value.$quixosCrdtType, encoding: typeof value.$quixosCrdtEncoding === "string" ? value.$quixosCrdtEncoding : "base64", payload: base64ToBytes(value.$quixosCrdtPayload), }), }, }); } if (!isRecord(value)) throw new Error(`Unsupported runtime value ${typeof value}`); return create(ValueSchema, { kind: { case: "objectValue", value: create(ObjectValueSchema, { fields: Object.fromEntries(Object.entries(value).map(([key, entry]) => [key, jsToProtoValue(entry)])), }), }, }); }; export const protoValueToJs = (value) => { switch (value?.kind.case) { case "nullValue": case undefined: return null; case "boolValue": case "numberValue": case "stringValue": case "integerValue": return value.kind.value; case "bytesValue": return bytesToBase64(value.kind.value); case "refValue": return referenceFromWire(value.kind.value.objectId); case "listValue": return value.kind.value.values.map(protoValueToJs); case "objectValue": return Object.fromEntries(Object.entries(value.kind.value.fields).map(([key, entry]) => [key, protoValueToJs(entry)])); case "crdtValue": return { $quixosCrdtType: value.kind.value.type, $quixosCrdtEncoding: value.kind.value.encoding, $quixosCrdtPayload: bytesToBase64(value.kind.value.payload), }; } }; export const protoFieldsToJs = (fields) => Object.fromEntries(Object.entries(fields).map(([key, value]) => [key, protoValueToJs(value)])); export class RuntimeAuthorityError extends Error { retryable; constructor(message) { super(message); this.name = "RuntimeAuthorityError"; this.retryable = /WORKSPACE_FENCED|STALE_EPOCH/.test(message); } } const targetForEdge = (edge, projectionId) => (edge.firstProjectionId === projectionId ? edge.secondObjectId : edge.firstObjectId); export const createRuntimeContext = (camino, orch, request) => { const acquiredPort = (objectId, interfaceRevisionId, conformance) => { const invoke = async (operationId, input = {}) => { const response = await orch.invokeCapability({ objectId, capability: create(CapabilityRefSchema, { interfaceRevisionId, operationId, conformance }), input: Object.fromEntries(Object.entries(input).map(([key, value]) => [key, jsToProtoValue(value)])), }); if (!response.ok) throw new Error(response.error || "Capability invocation failed"); for (const dependency of response.dependencies) { if (dependency.kind === "state" || dependency.kind === "edge") await recordDependency({ kind: dependency.kind, objectId: dependency.objectId, attachmentId: dependency.attachmentId, ...(dependency.kind === "edge" ? { projectionId: dependency.projectionId } : {}), }); } if (!response.result) throw new Error("Capability returned no value"); return response.result; }; return { objectId: referenceFromWire(objectId), interfaceRevisionId, invoke: async (operation, input) => protoValueToJs(await invoke(operation, input)), live: async (operation, input) => liveValue(await invoke(operation, input)), }; }; const ports = new Map(); for (const dependency of request.dependencies) { switch (dependency.binding.case) { case "queryId": { const queryId = dependency.binding.value, objectId = dependency.objectId || request.objectId; ports.set(dependency.portId, { queryId, execute: (variables, expectedDefinitionDigest) => orch.executeQuery({ queryId, objectId, variables, expectedDefinitionDigest }), watch: (variables, signal, expectedDefinitionDigest) => orch.watchQuery({ queryId, objectId, variables, expectedDefinitionDigest }, { signal }), }); break; } case "stateSlotId": { const slotId = dependency.binding.value; const dependencyObjectId = dependency.objectId || request.objectId; const state = { slotId, async get() { await recordDependency({ kind: "state", objectId: dependencyObjectId, attachmentId: slotId }); return protoValueToJs((await camino.readState({ objectId: dependencyObjectId, slotId })).value); }, async live() { await recordDependency({ kind: "state", objectId: dependencyObjectId, attachmentId: slotId }); const value = await camino.readState({ objectId: dependencyObjectId, slotId }); if (!value.value) throw new Error(`State ${slotId} returned no value`); return liveValue(value.value); }, async set(value) { assertReferenceFree(isWrappedValue(value) ? protoValueToJs(value.$quixosValue) : value); await camino.writeState({ objectId: dependencyObjectId, slotId, value: jsToProtoValue(value) }); }, }; ports.set(dependency.portId, state); break; } case "edge": { const { edgeTypeId, projectionId } = dependency.binding.value; const dependencyObjectId = dependency.objectId || request.objectId; const collectionResult = (response) => ({ revision: response.revision, entries: response.entries.map((entry) => ({ edgeId: entry.edgeId, target: referenceFromWire(entry.targetObjectId), ...(entry.key ? { key: entry.key.kind.case === "integerValue" ? BigInt(entry.key.kind.value) : protoValueToJs(entry.key), } : {}), })), }); const edge = { edgeTypeId, projectionId, async collection() { await recordDependency({ kind: "edge", objectId: dependencyObjectId, attachmentId: edgeTypeId, projectionId, }); return collectionResult(await camino.readCollection({ objectId: dependencyObjectId, edgeTypeId, projectionId })); }, async replace(entries, expectedRevision) { return collectionResult(await camino.replaceCollection({ objectId: dependencyObjectId, edgeTypeId, projectionId, expectedRevision, entries: entries.map((entry) => { assertReferenceFree(entry.key); return { edgeId: entry.edgeId ?? "", targetObjectId: referenceToWire(entry.target), key: entry.key === undefined ? undefined : jsToProtoValue(entry.key), }; }), })); }, async resolve() { await recordDependency({ kind: "edge", objectId: dependencyObjectId, attachmentId: edgeTypeId, projectionId, }); const result = await camino.resolveEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId }); return result.edges.map((entry) => referenceFromWire(targetForEdge(entry, projectionId))); }, async connect(target) { const targetObjectId = referenceToWire(target); await camino.connectEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId, targetObjectId }); }, async disconnect(target) { const targetObjectId = referenceToWire(target); const result = await camino.resolveEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId }); for (const entry of result.edges) { if (targetForEdge(entry, projectionId) === targetObjectId) await camino.disconnectEdge({ edgeId: entry.id }); } }, }; ports.set(dependency.portId, edge); break; } case "interfaceRevisionId": { const interfaceRevisionId = dependency.binding.value; const dependencyObjectId = dependency.objectId || request.objectId; ports.set(dependency.portId, acquiredPort(dependencyObjectId, interfaceRevisionId)); break; } case "constructorAtomId": { const atomId = dependency.binding.value; const constructor = { atomId, async construct(input = {}) { const response = await orch.constructObject({ atomId, input: Object.fromEntries(Object.entries(input).map(([key, value]) => [key, jsToProtoValue(value)])), }); if (!response.object) throw new Error(`Constructor ${atomId} returned no object`); return referenceFromWire(response.object.id); }, }; ports.set(dependency.portId, constructor); } } } const requirePort = (portId, kind) => { const port = ports.get(portId); if (!port || !(kind in port)) throw new Error(`Missing ${kind} dependency port ${portId}`); return port; }; return { async tryConform(object, interfaceRevisionId) { const objectId = referenceToWire(object); const { conformance } = await orch.tryConform({ objectId, interfaceRevisionId }); if (conformance && (conformance.objectId !== objectId || conformance.interfaceRevisionId !== interfaceRevisionId)) throw new Error("Conformance response does not match the requested view"); return conformance ? acquiredPort(objectId, interfaceRevisionId, conformance) : undefined; }, get objectId() { return referenceFromWire(request.objectId); }, input: protoFieldsToJs(request.input), inputProto: request.input, ports, state: (portId) => requirePort(portId, "slotId"), edge: (portId) => requirePort(portId, "edgeTypeId"), interface: (portId) => requirePort(portId, "interfaceRevisionId"), constructor: (portId) => requirePort(portId, "atomId"), query: (portId) => requirePort(portId, "queryId"), }; }; export const derived = (get) => ({ kind: "derived", get }); const isDerived = (handler) => typeof handler === "object" && handler.kind === "derived"; const evaluate = async (handler, context, observe) => { const dependencies = new Map(); const result = await dependencyScope.run({ dependencies, observe }, () => isDerived(handler) ? handler.get(context) : handler(context)); return { value: jsToProtoValue(result), dependencies: [...dependencies.values()] }; }; const protoDependencies = (dependencies) => dependencies.map((entry) => create(DerivedDependencySchema, { kind: entry.kind, objectId: entry.objectId, attachmentId: entry.attachmentId, projectionId: entry.kind === "edge" ? entry.projectionId : "", })); export const createPackageRuntimeRoutes = (config) => { const invocations = createInvocationRegistry(); const headers = {}; const processToken = process.env.CAMINO_RUNTIME_AUTH_TOKEN ?? (process.env.CAMINO_RUNTIME_AUTH_TOKEN_FILE ? readFileSync(process.env.CAMINO_RUNTIME_AUTH_TOKEN_FILE, "utf8").trim() : ""); if (processToken) { headers["x-camino-runtime-token"] = processToken; } else if (process.env.CAMINO_RUNTIME_AUTH_REQUIRED === "1") { throw new Error("CAMINO_RUNTIME_AUTH_TOKEN is required"); } const camino = createClient(CaminoService, createConnectTransport({ baseUrl: config.caminoUrl ?? process.env.CAMINO_URL ?? "http://127.0.0.1:7310", httpVersion: "1.1", interceptors: headers["x-camino-runtime-token"] ? [ (next) => async (request) => { request.header.set("x-camino-runtime-token", headers["x-camino-runtime-token"]); return await next(request); }, ] : [], })); const orch = createClient(OrchestratorRuntime, createConnectTransport({ baseUrl: config.orchUrl ?? process.env.QUIXOS_ORCH_URL ?? "http://127.0.0.1:7311", httpVersion: "1.1", })); const authenticateInstance = (header) => { if (!process.env.QUIXOS_RUNTIME_INSTANCE_ID) return; // standalone development ABI const supplied = Buffer.from(header.get("x-quixos-instance-token") ?? ""); const expected = Buffer.from(processToken); if (!expected.length || supplied.length !== expected.length || !timingSafeEqual(supplied, expected)) { throw new ConnectError("Invalid runtime instance credential", Code.Unauthenticated); } }; const clientsFor = (request) => { const context = request.context; if (process.env.QUIXOS_RUNTIME_INSTANCE_ID && (!context?.grant || context.instanceId !== process.env.QUIXOS_RUNTIME_INSTANCE_ID || !context.workspaceEpoch)) { throw new ConnectError("Managed invocation requires an exact instance and epoch grant", Code.Unauthenticated); } if (!context?.grant) return { camino, orch }; const transport = (url) => createConnectTransport({ baseUrl: url, httpVersion: "1.1", interceptors: [ (next) => async (call) => { call.header.set("x-quixos-invocation-grant", context.grant); call.header.set("x-camino-runtime-token", processToken); return next(call); }, ], }); return { camino: createClient(CaminoService, transport(config.caminoUrl ?? process.env.CAMINO_URL ?? "http://127.0.0.1:7310")), orch: createClient(OrchestratorRuntime, transport(config.orchUrl ?? process.env.QUIXOS_ORCH_URL ?? "http://127.0.0.1:7311")), }; }; const runtimeControl = async (operation, input) => { const response = await fetch(`${config.caminoUrl ?? process.env.CAMINO_URL ?? "http://127.0.0.1:7310"}/__runtime/${operation}`, { method: "POST", headers: { "content-type": "application/json", "x-camino-runtime-token": processToken }, body: JSON.stringify(input), signal: AbortSignal.timeout(10_000), }); const value = (await response.json()); if (!response.ok) throw new RuntimeAuthorityError(value.error ?? "Runtime authority request failed"); return value; }; const attachSessions = (runtimeContext, request) => { if (!request.context?.grant || !request.context.ownerConformanceId) return; runtimeContext.openSession = async () => { const ownerId = request.context.ownerConformanceId; const registration = { grant: request.context.grant, objectId: request.objectId, ownerId, sessionId: `session:${randomBytes(16).toString("hex")}`, token: randomBytes(32).toString("base64url"), }; const register = () => runtimeControl("register-session", registration); const registered = await register().catch((error) => { // Retry a transport/lost-response failure with exactly the same identity. // Admission/authority errors are definitive and must not be retried here. if (error instanceof RuntimeAuthorityError) throw error; return register(); }); let closed = false; return { id: registered.sessionId, async run(work) { if (closed) throw new RuntimeAuthorityError("SESSION_CLOSED"); // Acquisition happens before user code. A fence failure can be retried // by the caller without replaying a side-effecting callback. const grant = await runtimeControl("acquire-session", registered); const execution = invocations.begin(grant.invocationId); const sessionRequest = { ...request, context: { grant: grant.grant, instanceId: grant.instanceId, workspaceEpoch: grant.epoch }, }; const clients = clientsFor(sessionRequest); const context = createRuntimeContext(clients.camino, clients.orch, sessionRequest); context.signal = execution.signal; try { return await work(context); } finally { execution.finish(); await runtimeControl("complete-invocation", { invocationId: grant.invocationId }).catch((error) => console.error("Session completion will be reconciled by the host", error)); } }, async close() { await runtimeControl("close-session", registered); closed = true; }, }; }; }; return (router) => router.service(PackageRuntime, { handshake: (request) => create(HandshakeResponseSchema, { packageRevisionId: config.packageRevisionId, runtimeProtocolVersion: "quixos-capabilities-v1", exportIds: Object.keys(config.exports), capabilities: ["invocation-completion-v1", "instance-authentication-v1", "epoch-grants-v1"], instanceId: process.env.QUIXOS_RUNTIME_INSTANCE_ID ?? "", authenticationProof: request.nonce && processToken ? createHmac("sha256", processToken) .update(JSON.stringify([ request.nonce, process.env.QUIXOS_RUNTIME_INSTANCE_ID ?? "", config.packageRevisionId, ])) .digest("hex") : "", }), getInvocationStatus: (request, context) => { authenticateInstance(context.requestHeader); return invocations.status(request.invocationId); }, cancelInvocation: (request, context) => { authenticateInstance(context.requestHeader); return invocations.cancel(request.invocationId); }, invoke: async (request, context) => { authenticateInstance(context.requestHeader); const { camino, orch } = clientsFor(request); const exportId = request.export?.exportId; const handler = exportId ? config.exports[exportId] : undefined; if (!handler) throw new ConnectError(`Unknown export ${exportId ?? ""}`, Code.NotFound); const execution = invocations.begin(request.invocationId || (process.env.QUIXOS_RUNTIME_INSTANCE_ID ? "" : randomUUID())); try { const runtimeContext = createRuntimeContext(camino, orch, request); attachSessions(runtimeContext, request); runtimeContext.signal = AbortSignal.any([context.signal, execution.signal]); const result = await evaluate(handler, runtimeContext); execution.finish(); return create(InvokeResponseSchema, { ok: true, result: result.value, dependencies: protoDependencies(result.dependencies), }); } catch (error) { execution.finish(true); return create(InvokeResponseSchema, { ok: false, error: error instanceof Error ? error.message : String(error), }); } }, watch: async function* (request, context) { authenticateInstance(context.requestHeader); const { camino, orch } = clientsFor(request); const exportId = request.export?.exportId; const handler = exportId ? config.exports[exportId] : undefined; if (!handler || !isDerived(handler)) { throw new ConnectError(`Export ${exportId ?? ""} is not derived`, Code.FailedPrecondition); } const execution = invocations.begin(request.invocationId || (process.env.QUIXOS_RUNTIME_INSTANCE_ID ? "" : randomUUID())); const signal = AbortSignal.any([context.signal, execution.signal]); try { const watchId = `watch:${randomUUID()}`; const runtimeContext = createRuntimeContext(camino, orch, request); attachSessions(runtimeContext, request); runtimeContext.signal = signal; const subscriptions = new Map(); const establishing = new Map(); let subscriptionEpoch = 0; const ensureSubscription = async (dependency) => { const key = dependencyKey(dependency); if (subscriptions.has(key)) return; const pending = establishing.get(key); if (pending) return await pending; const establish = (async () => { const controller = new AbortController(); const stream = camino .watchObject({ objectId: dependency.objectId, includeSnapshot: true, attachmentIds: request.context?.grant ? [dependency.attachmentId] : [], }, { signal: controller.signal })[Symbol.asyncIterator](); try { // Camino subscribes before producing the snapshot, so once this // resolves the following state/edge read cannot race the stream. const snapshot = await stream.next(); if (snapshot.done) throw new Error(`Dependency stream ${key} ended during setup`); const waitNext = () => stream.next().then((result) => ({ key, done: Boolean(result.done) }), (error) => ({ key, done: true, error })); const subscription = { dependency, controller, waitNext, next: Promise.resolve({ key, done: false }), }; subscription.next = waitNext(); subscriptions.set(key, subscription); subscriptionEpoch += 1; } catch (error) { controller.abort(); throw error; } })(); establishing.set(key, establish); try { await establish; } finally { establishing.delete(key); } }; const abortAll = () => { for (const subscription of subscriptions.values()) { subscription.controller.abort(); } }; signal.addEventListener("abort", abortAll, { once: true }); const evaluateWithStableSubscriptions = async () => { // A direct state/edge port records its dependency before reading it, // but a nested interface invocation can only report its transitive // dependencies after that invocation returns. Once a new dependency // stream is established, evaluate again so every read contributing to // the emitted value happened after its stream became live. for (let pass = 0; pass < 32; pass += 1) { const before = subscriptionEpoch; const result = await evaluate(handler, runtimeContext, ensureSubscription); if (subscriptionEpoch === before) return result; } throw new Error("Derived dependency discovery did not stabilize after 32 passes"); }; try { let current = await evaluateWithStableSubscriptions(); yield create(WatchEventSchema, { watchId, value: current.value, dependencies: protoDependencies(current.dependencies), initial: true, }); const abort = new Promise((resolve) => { if (signal.aborted) resolve("abort"); else signal.addEventListener("abort", () => resolve("abort"), { once: true }); }); while (!signal.aborted) { if (subscriptions.size === 0) { await abort; break; } const outcome = await Promise.race([...[...subscriptions.values()].map((entry) => entry.next), abort]); if (outcome === "abort") break; const subscription = subscriptions.get(outcome.key); if (!subscription) continue; if (outcome.error) throw outcome.error; if (outcome.done) throw new Error(`Dependency stream ${outcome.key} ended unexpectedly`); subscription.next = subscription.waitNext(); const updated = await evaluateWithStableSubscriptions(); const active = new Set(updated.dependencies.map(dependencyKey)); for (const [key, entry] of subscriptions) { if (!active.has(key)) { entry.controller.abort(); subscriptions.delete(key); } } if (!equals(ValueSchema, current.value, updated.value)) { yield create(WatchEventSchema, { watchId, value: updated.value, dependencies: protoDependencies(updated.dependencies), }); } current = updated; } } finally { signal.removeEventListener("abort", abortAll); abortAll(); } } finally { execution.finish(); } }, }); }; export const servePackageRuntime = (config) => { const host = process.env.QUIXOS_RUNTIME_HOST ?? "127.0.0.1"; const port = Number(process.env.QUIXOS_RUNTIME_PORT ?? "0"); const handler = connectNodeAdapter({ routes: createPackageRuntimeRoutes(config) }); const server = http.createServer((request, response) => void handler(request, response)); server.listen(port, host, () => console.log(`${config.packageRevisionId} listening on ${host}:${port}`)); const shutdown = () => { server.close(() => process.exit(0)); server.closeAllConnections(); }; process.on("SIGTERM", shutdown); process.on("SIGINT", shutdown); return server; };