Host workspace repositories on Central Gitea
This commit is contained in:
@@ -0,0 +1,164 @@
|
||||
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,
|
||||
protoValueToJs,
|
||||
} from "../dist/index.js";
|
||||
import {
|
||||
CaminoObjectSchema,
|
||||
CaminoService,
|
||||
} 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("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();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user