import assert from "node:assert/strict"; import http from "node:http"; import test from "node:test"; import { referenceFromWire } from "../dist/references.js"; import { create } from "@bufbuild/protobuf"; import { createClient } from "@connectrpc/connect"; import { connectNodeAdapter, createConnectTransport } from "@connectrpc/connect-node"; import { createPackageRuntimeRoutes, createRuntimeContext, derived, jsToProtoValue, liveValue, protoValueToJs, } from "../dist/index.js"; import { CaminoObjectSchema, CaminoService, CrdtValueSchema, StateValueSourceSchema, ValueSchema, ValueSourceSchema, } from "../dist/camino/api_pb.js"; import { EdgeDependencySchema, InjectedDependencySchema, PackageExportRefSchema } from "../dist/quixos/refs_pb.js"; import { InvokeCapabilityResponseSchema, OrchestratorRuntime } from "../dist/quixos/orch_pb.js"; import { DerivedDependencySchema, PackageRuntime } from "../dist/quixos/runtime_pb.js"; const listen = async (routes) => { const server = http.createServer((request, response) => void connectNodeAdapter({ routes })(request, response)); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const address = server.address(); assert.ok(address && typeof address === "object"); return { url: `http://127.0.0.1:${address.port}`, close: async () => { server.closeAllConnections(); await new Promise((resolve) => server.close(resolve)); }, }; }; test("generic values preserve nested values and object references", () => { const value = jsToProtoValue({ title: "A task", target: referenceFromWire("obj:target"), tags: ["one", "two"], }); assert.deepEqual(protoValueToJs(value), { title: "A task", target: referenceFromWire("obj:target"), tags: ["one", "two"], }); }); test("live values preserve writable CRDT source identity through derivation", () => { const value = create(ValueSchema, { kind: { case: "stringValue", value: "notes" }, source: create(ValueSourceSchema, { state: create(StateValueSourceSchema, { objectId: "obj:task", slotId: "slot:task:notes", revision: 7n, crdtSnapshot: create(CrdtValueSchema, { type: "quixos.automerge-document.v1", encoding: "automerge-snapshot-v1", payload: new Uint8Array([1, 2, 3]), }), }), }), }); const forwarded = jsToProtoValue(liveValue(value)); assert.equal(forwarded.source?.state?.objectId, "obj:task"); assert.equal(forwarded.source?.state?.slotId, "slot:task:notes"); assert.equal(forwarded.source?.state?.revision, 7n); assert.deepEqual(forwarded.source?.state?.crdtSnapshot?.payload, new Uint8Array([1, 2, 3])); }); test("runtime context exposes only explicitly injected ports", async () => { const writes = []; const camino = { readState: async ({ objectId, slotId }) => ({ value: jsToProtoValue(`${objectId}/${slotId}`), }), writeState: async (request) => { writes.push(request); return {}; }, resolveEdge: async () => ({ edges: [] }), connectEdge: async () => ({}), }; const context = createRuntimeContext( camino, {}, { objectId: "obj:task", input: {}, dependencies: [ create(InjectedDependencySchema, { portId: "port:title", binding: { case: "stateSlotId", value: "slot:task:title" }, }), create(InjectedDependencySchema, { portId: "port:children", binding: { case: "edge", value: create(EdgeDependencySchema, { edgeTypeId: "edge:task:children", projectionId: "projection:task:children", }), }, }), ], }, ); assert.equal(await context.state("port:title").get(), "obj:task/slot:task:title"); await context.state("port:title").set("Changed"); assert.equal(writes.length, 1); assert.equal("camino" in context, false); assert.equal("orch" in context, false); assert.throws(() => context.state("port:not-injected"), /Missing slotId dependency/); assert.throws(() => context.interface("port:title"), /Missing interfaceRevisionId dependency/); assert.deepEqual(await context.edge("port:children").resolve(), []); }); test("injected ports preserve an explicitly traversed object through PackageRuntime RPC", async () => { const reads = []; const caminoServer = await listen((router) => router.service(CaminoService, { readState: (request) => { reads.push(request); return { value: jsToProtoValue("Project name") }; }, }), ); const runtimeServer = await listen( createPackageRuntimeRoutes({ packageRevisionId: "package:test@1", caminoUrl: caminoServer.url, exports: { "export:test:value": (context) => context.state("port:value").get(), }, }), ); const client = createClient( PackageRuntime, createConnectTransport({ baseUrl: runtimeServer.url, httpVersion: "1.1", }), ); try { const response = await client.invoke({ export: create(PackageExportRefSchema, { packageRevisionId: "package:test@1", exportId: "export:test:value", }), objectId: "obj:component", dependencies: [ create(InjectedDependencySchema, { portId: "port:value", objectId: "obj:project", binding: { case: "stateSlotId", value: "slot:project:name" }, }), ], }); assert.equal(response.ok, true); assert.equal(protoValueToJs(response.result), "Project name"); assert.equal(reads.length, 1); assert.equal(reads[0]?.objectId, "obj:project"); } finally { await runtimeServer.close(); await caminoServer.close(); } }); test("interface views invoke a related object and propagate transitive live dependencies", async () => { const invocations = []; const sourceValue = create(ValueSchema, { kind: { case: "stringValue", value: "Project name" }, source: create(ValueSourceSchema, { state: create(StateValueSourceSchema, { objectId: "obj:project", slotId: "slot:project:name", revision: 4n, }), }), }); const orchServer = await listen((router) => router.service(OrchestratorRuntime, { invokeCapability: (request) => { invocations.push(request); return create(InvokeCapabilityResponseSchema, { ok: true, result: sourceValue, dependencies: [ create(DerivedDependencySchema, { kind: "state", objectId: "obj:project", attachmentId: "slot:project:name", }), ], }); }, }), ); const runtimeServer = await listen( createPackageRuntimeRoutes({ packageRevisionId: "package:test@1", orchUrl: orchServer.url, exports: { "export:test:value": derived((context) => context.interface("port:named").live("operation:named:name:get")), }, }), ); const client = createClient( PackageRuntime, createConnectTransport({ baseUrl: runtimeServer.url, httpVersion: "1.1", }), ); try { const response = await client.invoke({ export: create(PackageExportRefSchema, { packageRevisionId: "package:test@1", exportId: "export:test:value", }), objectId: "obj:component", dependencies: [ create(InjectedDependencySchema, { portId: "port:named", objectId: "obj:project", binding: { case: "interfaceRevisionId", value: "interface:named@1" }, }), ], }); assert.equal(response.ok, true); assert.equal(invocations[0]?.objectId, "obj:project"); assert.equal(invocations[0]?.capability?.interfaceRevisionId, "interface:named@1"); assert.equal(response.dependencies[0]?.objectId, "obj:project"); assert.equal(response.dependencies[0]?.attachmentId, "slot:project:name"); assert.equal(response.result?.source?.state?.revision, 4n); } finally { await runtimeServer.close(); await orchServer.close(); } }); test("derived interface views reread after their transitive subscriptions become live", async () => { let current = "before subscription"; let invocationCount = 0; const caminoServer = await listen((router) => router.service(CaminoService, { watchObject: async function* (request, context) { // Model a write racing the first nested capability read. Camino makes // the subscription live before yielding this snapshot. current = "after subscription"; yield { objectId: request.objectId, snapshot: create(CaminoObjectSchema, { id: request.objectId, atomId: "atom:project", workspaceRevisionId: "workspace:test@1", }), }; await new Promise((resolve) => context.signal.addEventListener("abort", resolve, { once: true })); }, }), ); const orchServer = await listen((router) => router.service(OrchestratorRuntime, { invokeCapability: () => { invocationCount += 1; return create(InvokeCapabilityResponseSchema, { ok: true, result: jsToProtoValue(current), dependencies: [ create(DerivedDependencySchema, { kind: "state", objectId: "obj:project", attachmentId: "slot:project:name", }), ], }); }, }), ); const runtimeServer = await listen( createPackageRuntimeRoutes({ packageRevisionId: "package:test@1", caminoUrl: caminoServer.url, orchUrl: orchServer.url, exports: { "export:test:value": derived((context) => context.interface("port:named").invoke("operation:named:name:get")), }, }), ); const client = createClient( PackageRuntime, createConnectTransport({ baseUrl: runtimeServer.url, httpVersion: "1.1", }), ); const controller = new AbortController(); try { const stream = client .watch( { export: create(PackageExportRefSchema, { packageRevisionId: "package:test@1", exportId: "export:test:value", }), objectId: "obj:component", dependencies: [ create(InjectedDependencySchema, { portId: "port:named", objectId: "obj:project", binding: { case: "interfaceRevisionId", value: "interface:named@1" }, }), ], }, { signal: controller.signal }, ) [Symbol.asyncIterator](); const initial = await stream.next(); assert.equal(protoValueToJs(initial.value?.value), "after subscription"); assert.equal(invocationCount, 2); } finally { controller.abort(); await runtimeServer.close(); await orchServer.close(); await caminoServer.close(); } }); test("derived watches subscribe before reading and emit only real changes", async () => { let current = "before"; let subscribed = false; const changes = []; const caminoServer = await listen((router) => router.service(CaminoService, { readState: () => { assert.equal(subscribed, true, "the dependency stream must be live before the state read"); return { value: jsToProtoValue(current) }; }, watchObject: async function* (request, context) { subscribed = true; yield { objectId: request.objectId, snapshot: create(CaminoObjectSchema, { id: request.objectId, atomId: "atom:test", workspaceRevisionId: "workspace:test@1", }), }; while (!context.signal.aborted) { const changed = await new Promise((resolve) => { const finish = () => resolve(true); changes.push(finish); context.signal.addEventListener("abort", () => resolve(false), { once: true }); }); if (!changed) return; yield { objectId: request.objectId }; } }, }), ); const runtimeServer = await listen( createPackageRuntimeRoutes({ packageRevisionId: "package:test@1", caminoUrl: caminoServer.url, exports: { "export:test:value": derived((context) => context.state("port:value").get()), }, }), ); const client = createClient( PackageRuntime, createConnectTransport({ baseUrl: runtimeServer.url, httpVersion: "1.1", }), ); const controller = new AbortController(); try { const stream = client .watch( { export: create(PackageExportRefSchema, { packageRevisionId: "package:test@1", exportId: "export:test:value", }), objectId: "obj:test", dependencies: [ create(InjectedDependencySchema, { portId: "port:value", binding: { case: "stateSlotId", value: "slot:test:value" }, }), ], }, { signal: controller.signal }, ) [Symbol.asyncIterator](); const initial = await stream.next(); assert.equal(protoValueToJs(initial.value?.value), "before"); assert.equal(initial.value?.initial, true); current = "after"; changes.shift()?.(); const updated = await stream.next(); assert.equal(protoValueToJs(updated.value?.value), "after"); assert.equal(updated.value?.initial, false); } finally { controller.abort(); await runtimeServer.close(); await caminoServer.close(); } });