662 lines
33 KiB
JavaScript
662 lines
33 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 * 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;
|
|
};
|