Implement workspace evolution, migrations, and runtime continuity
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.
This commit is contained in:
Vendored
+265
-129
@@ -1,7 +1,12 @@
|
||||
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 { randomUUID } from "node:crypto";
|
||||
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";
|
||||
@@ -24,9 +29,11 @@ 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 = (objectId) => ({ $quixosRef: objectId });
|
||||
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) {
|
||||
@@ -47,11 +54,8 @@ export const jsToProtoValue = (value) => {
|
||||
kind: { case: "listValue", value: create(ListValueSchema, { values: value.map(jsToProtoValue) }) },
|
||||
});
|
||||
}
|
||||
if (isRecord(value) && typeof value.$quixosRef === "string") {
|
||||
return create(ValueSchema, {
|
||||
kind: { case: "refValue", value: create(RefValueSchema, { objectId: value.$quixosRef }) },
|
||||
});
|
||||
}
|
||||
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, {
|
||||
@@ -79,7 +83,7 @@ export const protoValueToJs = (value) => {
|
||||
case "stringValue":
|
||||
case "integerValue": return value.kind.value;
|
||||
case "bytesValue": return bytesToBase64(value.kind.value);
|
||||
case "refValue": return value.kind.value.objectId;
|
||||
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 {
|
||||
@@ -90,6 +94,10 @@ export const protoValueToJs = (value) => {
|
||||
}
|
||||
};
|
||||
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();
|
||||
@@ -112,6 +120,7 @@ export const createRuntimeContext = (camino, orch, request) => {
|
||||
return liveValue(value.value);
|
||||
},
|
||||
async set(value) {
|
||||
assertReferenceFree(isWrappedValue(value) ? protoValueToJs(value.$quixosValue) : value);
|
||||
await camino.writeState({ objectId: dependencyObjectId, slotId, value: jsToProtoValue(value) });
|
||||
},
|
||||
};
|
||||
@@ -121,18 +130,34 @@ export const createRuntimeContext = (camino, orch, request) => {
|
||||
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) => targetForEdge(entry, projectionId));
|
||||
return result.edges.map((entry) => referenceFromWire(targetForEdge(entry, projectionId)));
|
||||
},
|
||||
async connect(targetObjectId) {
|
||||
async connect(target) {
|
||||
const targetObjectId = referenceToWire(target);
|
||||
await camino.connectEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId, targetObjectId });
|
||||
},
|
||||
async disconnect(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)
|
||||
@@ -174,7 +199,7 @@ export const createRuntimeContext = (camino, orch, request) => {
|
||||
return response.result;
|
||||
};
|
||||
const capability = {
|
||||
objectId: dependencyObjectId,
|
||||
objectId: referenceFromWire(dependencyObjectId),
|
||||
interfaceRevisionId,
|
||||
async invoke(operationId, input = {}) {
|
||||
return protoValueToJs(await invoke(operationId, input));
|
||||
@@ -200,7 +225,7 @@ export const createRuntimeContext = (camino, orch, request) => {
|
||||
});
|
||||
if (!response.object)
|
||||
throw new Error(`Constructor ${atomId} returned no object`);
|
||||
return response.object.id;
|
||||
return referenceFromWire(response.object.id);
|
||||
},
|
||||
};
|
||||
ports.set(dependency.portId, constructor);
|
||||
@@ -214,7 +239,7 @@ export const createRuntimeContext = (camino, orch, request) => {
|
||||
return port;
|
||||
};
|
||||
return {
|
||||
objectId: request.objectId,
|
||||
objectId: referenceFromWire(request.objectId),
|
||||
input: protoFieldsToJs(request.input),
|
||||
inputProto: request.input,
|
||||
ports,
|
||||
@@ -238,9 +263,12 @@ const protoDependencies = (dependencies) => dependencies.map((entry) => create(D
|
||||
projectionId: entry.kind === "edge" ? entry.projectionId : "",
|
||||
}));
|
||||
export const createPackageRuntimeRoutes = (config) => {
|
||||
const invocations = createInvocationRegistry();
|
||||
const headers = {};
|
||||
if (process.env.CAMINO_RUNTIME_AUTH_TOKEN) {
|
||||
headers["x-camino-runtime-token"] = process.env.CAMINO_RUNTIME_AUTH_TOKEN;
|
||||
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");
|
||||
@@ -259,19 +287,115 @@ export const createPackageRuntimeRoutes = (config) => {
|
||||
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: () => create(HandshakeResponseSchema, {
|
||||
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") : "",
|
||||
}),
|
||||
invoke: async (request) => {
|
||||
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 result = await evaluate(handler, createRuntimeContext(camino, orch, request));
|
||||
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,
|
||||
@@ -279,6 +403,7 @@ export const createPackageRuntimeRoutes = (config) => {
|
||||
});
|
||||
}
|
||||
catch (error) {
|
||||
execution.finish(true);
|
||||
return create(InvokeResponseSchema, {
|
||||
ok: false,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
@@ -286,130 +411,141 @@ export const createPackageRuntimeRoutes = (config) => {
|
||||
}
|
||||
},
|
||||
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 watchId = `watch:${randomUUID()}`;
|
||||
const runtimeContext = createRuntimeContext(camino, orch, request);
|
||||
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 }, { signal: controller.signal })[Symbol.asyncIterator]();
|
||||
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 {
|
||||
// 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;
|
||||
await establish;
|
||||
}
|
||||
catch (error) {
|
||||
controller.abort();
|
||||
throw error;
|
||||
finally {
|
||||
establishing.delete(key);
|
||||
}
|
||||
})();
|
||||
establishing.set(key, establish);
|
||||
};
|
||||
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 {
|
||||
await establish;
|
||||
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 {
|
||||
establishing.delete(key);
|
||||
}
|
||||
};
|
||||
const abortAll = () => {
|
||||
for (const subscription of subscriptions.values()) {
|
||||
subscription.controller.abort();
|
||||
}
|
||||
};
|
||||
context.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");
|
||||
};
|
||||
let current = await evaluateWithStableSubscriptions();
|
||||
yield create(WatchEventSchema, {
|
||||
watchId,
|
||||
value: current.value,
|
||||
dependencies: protoDependencies(current.dependencies),
|
||||
initial: true,
|
||||
});
|
||||
const abort = new Promise((resolve) => {
|
||||
if (context.signal.aborted)
|
||||
resolve("abort");
|
||||
else
|
||||
context.signal.addEventListener("abort", () => resolve("abort"), { once: true });
|
||||
});
|
||||
try {
|
||||
while (!context.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;
|
||||
signal.removeEventListener("abort", abortAll);
|
||||
abortAll();
|
||||
}
|
||||
}
|
||||
finally {
|
||||
context.signal.removeEventListener("abort", abortAll);
|
||||
abortAll();
|
||||
execution.finish();
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user