import assert from "node:assert/strict"; import http from "node:http"; import test from "node:test"; 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 { 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: { $quixosRef: "obj:target" }, tags: ["one", "two"], }); assert.deepEqual(protoValueToJs(value), { title: "A task", target: "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("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(); } });