Files
camino-package-runtime/test/runtime.test.mjs
T

422 lines
13 KiB
JavaScript

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();
}
});