ce6ae8f662
Enable evolution by default for source-backed workspaces. Add stable conformance ownership, semantic-major review, candidate typechecking, and durable fenced cutover with explicit migrations and forward recovery. Independently supervise package runtimes so unchanged resource owners keep their processes and connections across cutover. Add scoped invocation authority, resource sessions, and typed callback rebinding. Wire opaque object references through generated bindings and RPCs. Add canonical relationship sets, keyed maps, and ordered lists with scoped transactional mutations, revision checks, and inverse consistency. Support planned cascade deletion, protection, tombstones, and lifecycle foundations. Add journaled structural edits, package/function/migration scaffolding, managed repository creation, and resumable bottom-up dependency pin publication. Document lifetime boundaries, revision pinning, prototype compatibility policy, commands, and deferred work. Validate with 210 tests, user-systemd process/connection continuity, generated-package TypeScript checks, and Nix host/protocol checks. TTL handoff, physical reclamation, general multi-step migrations, and root-systemd migration isolation acceptance remain deferred.
567 lines
31 KiB
JavaScript
567 lines
31 KiB
JavaScript
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 { 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 ports = new Map();
|
|
for (const dependency of request.dependencies) {
|
|
switch (dependency.binding.case) {
|
|
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;
|
|
const invoke = async (operationId, input) => {
|
|
const response = await orch.invokeCapability({
|
|
capability: create(CapabilityRefSchema, { interfaceRevisionId, operationId }),
|
|
objectId: dependencyObjectId,
|
|
input: Object.fromEntries(Object.entries(input).map(([key, value]) => [key, jsToProtoValue(value)])),
|
|
});
|
|
if (!response.ok)
|
|
throw new Error(response.error || `Capability ${operationId} failed`);
|
|
for (const dependency of response.dependencies) {
|
|
if (dependency.kind === "state") {
|
|
await recordDependency({
|
|
kind: "state",
|
|
objectId: dependency.objectId,
|
|
attachmentId: dependency.attachmentId,
|
|
});
|
|
}
|
|
else if (dependency.kind === "edge") {
|
|
await recordDependency({
|
|
kind: "edge",
|
|
objectId: dependency.objectId,
|
|
attachmentId: dependency.attachmentId,
|
|
projectionId: dependency.projectionId,
|
|
});
|
|
}
|
|
}
|
|
return response.result;
|
|
};
|
|
const capability = {
|
|
objectId: referenceFromWire(dependencyObjectId),
|
|
interfaceRevisionId,
|
|
async invoke(operationId, input = {}) {
|
|
return protoValueToJs(await invoke(operationId, input));
|
|
},
|
|
async live(operationId, input = {}) {
|
|
const value = await invoke(operationId, input);
|
|
if (!value)
|
|
throw new Error(`Capability ${operationId} returned no value`);
|
|
return liveValue(value);
|
|
},
|
|
};
|
|
ports.set(dependency.portId, capability);
|
|
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 {
|
|
objectId: 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"),
|
|
};
|
|
};
|
|
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;
|
|
};
|