Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
15 kB · 366 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { buildContextPacket } from "../src/agents/context.js";import { loadAgentDeclarations } from "../src/agents/declarations.js";import { auditRunEvidence } from "../src/agents/evidence.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import type { AgentRunner, ThoughtAgentDeclaration } from "../src/agents/types.js";import { FilesystemConnector } from "../src/connectors/filesystem.js";import { stableKey } from "../src/core/ids.js";import { sha256 } from "../src/core/json.js";import type { EventCandidate } from "../src/events/types.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject, testDeclarationEnvironment, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];
afterEach(async () => { await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("Jazz-native producer and consumer topology", () => { test("settles source-local sequence, events, and cursor in one durable producer batch", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const cursor = { id: "cursor:rss:fixture", source: "rss:fixture", cursor: { etag: "v1" }, lastSuccessAt: "2026-07-14T00:00:02.000Z", updatedAt: "2026-07-14T00:00:02.000Z", }; const candidates = [candidate("one"), candidate("two")];
const first = await store.appendProducerBatch(candidates, cursor); expect(first.inserted.map((event) => event.sourceSequence)).toEqual([1, 2]); expect((await store.listSources()).find((source) => source.id === "rss:fixture")?.lastSequence).toBe(2); expect(await store.getSourceCursor(cursor.id)).toEqual(cursor);
const replay = await store.appendProducerBatch(candidates, cursor); expect(replay.inserted).toHaveLength(0); expect(replay.unchanged.map((event) => event.sourceSequence)).toEqual([1, 2]); expect((await store.listSources()).find((source) => source.id === "rss:fixture")?.lastSequence).toBe(2);
const narrow = await store.queryConsumerEvents({ consumerId: "fixture-consumer", consumerVersion: 1, source: "rss:fixture", eventTypes: ["stream.thought.source.rss.item"], acceptedPrivacy: ["private"], afterSequence: 1, }); expect(narrow.map((event) => event.externalId)).toEqual(["two"]); });
test("starts a new consumer at the durable source head when replay is now", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const first = (await store.appendEvent(candidate("one"))).event; const declaration = { ...subscriptionDeclaration(), sourcePatterns: ["rss:fixture"], initialReplay: "now" as const, }; const runner: AgentRunner = { mode: "deterministic", run: async () => ({ summary: "Observed a post-activation item", tags: ["rss"], importance: "normal", confidence: 1, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]);
expect(await runtime.consumeBacklog([declaration])).toHaveLength(0); expect((await store.listConsumerProgress())[0]).toMatchObject({ lastSequence: 1, lastEventId: first.id });
await store.appendEvent(candidate("two")); const [processed] = await runtime.consumeBacklog([declaration]); expect(processed?.output?.summary).toBe("Observed a post-activation item"); expect((await store.listConsumerProgress())[0]).toMatchObject({ lastSequence: 2 }); });
test("replays from durable consumer progress and abandons interrupted execution without leases", async () => { const project = await temporaryProject(); roots.push(project); const vault = path.join(project, "vault"); await fs.cp(path.join(process.cwd(), "fixtures", "vault"), vault, { recursive: true }); let store = await testStore(project); stores.push(store); const connector = new FilesystemConnector({ id: "filesystem:fixture", root: vault }); const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const declaration = declarations.find((item) => item.id === "document-structure")!; let runtime = new ThoughtAgentRuntime(store);
const initial = await connector.scan(store); expect(initial.added).toBe(3); expect(await runtime.consumeBacklog(declarations)).toHaveLength(3); expect(await runtime.consumeBacklog(declarations)).toHaveLength(0);
const brief = path.join(vault, "spec", "brief.md"); await fs.appendFile(brief, "\n## Linked heading\n\nSee [[decision]].\n"); const changed = await connector.scan(store); const changedEvent = changed.events[0]!; expect(changedEvent.sourceSequence).toBeGreaterThan(0); const executionKey = stableKey("execution", declaration.id, String(declaration.version), changedEvent.id); const interruptedRunId = stableKey("run", executionKey, "1"); const interruptedAt = "2026-07-14T00:00:00.000Z"; await store.upsertRun({ id: interruptedRunId, executionKey, triggerEventId: changedEvent.id, agentId: declaration.id, agentVersion: declaration.version, status: "running", inputEventIds: [changedEvent.id], outputEventIds: [], attempt: 1, provider: declaration.mode, model: declaration.mode, promptHash: sha256(declaration.systemPrompt), contextManifest: buildContextPacket(declaration, changedEvent).manifest, createdAt: interruptedAt, startedAt: interruptedAt, updatedAt: interruptedAt, });
stores.splice(stores.indexOf(store), 1); await store.close(); store = await testStore(project); stores.push(store); runtime = new ThoughtAgentRuntime(store); const [recovered] = await runtime.consumeBacklog(declarations);
expect(recovered?.output && "recommendation" in recovered.output ? recovered.output.recommendation?.target : undefined).toBe("charter"); expect((await store.getRun(interruptedRunId))?.status).toBe("abandoned"); const latest = await store.latestRunForExecution(executionKey); expect(latest).toMatchObject({ status: "completed", attempt: 2, triggerEventId: changedEvent.id }); expect((await store.listEvents()).filter((event) => ( event.parentEventId === changedEvent.id && event.type === "stream.thought.derived.document.structure" ))).toHaveLength(1);
const progress = (await store.listConsumerProgress()).find((item) => ( item.consumerId === declaration.id && item.source === changedEvent.source )); expect(progress).toMatchObject({ consumerVersion: declaration.version, lastSequence: changedEvent.sourceSequence, lastEventId: changedEvent.id, }); expect(await runtime.consumeBacklog(declarations)).toHaveLength(0); expect(auditRunEvidence(await store.listRuns(), await store.listEvents()).every((report) => report.consistent)).toBe(true); });
test("discovers a wildcard source after startup and consumes it through a live Jazz subscription", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const declaration = subscriptionDeclaration(); const runner: AgentRunner = { mode: "deterministic", run: async () => ({ summary: "Observed a subscribed RSS item", tags: ["rss"], importance: "normal", confidence: 1, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); const consumers = await runtime.startConsumers([declaration]);
await store.appendProducerBatch([candidate("one"), { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: "cursor:rss:fixture", idempotencyKey: "cursor-one", occurredAt: "2026-07-14T00:00:00.500Z", actor: "rss:fixture", correlationId: "poll-fixture", privacy: "private", payload: { source: "rss:fixture", cursorId: "cursor:rss:fixture", cursor: { etag: "v1" } }, }]); await waitFor(async () => (await store.listConsumerProgress()).some((progress) => progress.lastSequence === 1)); await store.appendProducerBatch([candidate("two")]); await waitFor(async () => (await store.listConsumerProgress()).some((progress) => progress.lastSequence === 3)); await consumers.stop();
expect((await store.listRuns()).map((run) => run.triggerEventId)).toHaveLength(2); expect((await store.listRuns()).every((run) => run.status === "completed")).toBe(true); expect((await store.listConsumerProgress())[0]).toMatchObject({ consumerId: declaration.id, source: "rss:fixture", lastSequence: 3, }); expect((await store.listSources()).find((source) => source.id === "agent:rss-fixture-consumer")?.lastSequence).toBe(6); });
test("runs independent consumer queues concurrently without losing per-source ordering", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const slow = { ...subscriptionDeclaration(), id: "slow-consumer", name: "Slow consumer", sourcePatterns: ["rss:slow"], }; const fast = { ...subscriptionDeclaration(), id: "fast-consumer", name: "Fast consumer", sourcePatterns: ["rss:fast"], }; let releaseSlow!: () => void; let markFastStarted!: () => void; const slowCanFinish = new Promise<void>((resolve) => { releaseSlow = resolve; }); const fastStarted = new Promise<void>((resolve) => { markFastStarted = resolve; }); const executionOrder: string[] = []; const runner: AgentRunner = { mode: "deterministic", run: async ({ declaration, event }) => { executionOrder.push(`${declaration.id}:${event.externalId}:start`); if (declaration.id === slow.id) await slowCanFinish; else { markFastStarted(); releaseSlow(); } executionOrder.push(`${declaration.id}:${event.externalId}:finish`); return { summary: `Observed ${event.externalId}`, tags: ["rss"], importance: "normal", confidence: 1, }; }, }; await store.appendProducerBatch([ { ...candidate("one"), source: "rss:slow", externalId: "slow-one", idempotencyKey: "slow-one" }, { ...candidate("two"), source: "rss:slow", externalId: "slow-two", idempotencyKey: "slow-two" }, ]); await store.appendProducerBatch([ { ...candidate("one"), source: "rss:fast", externalId: "fast-one", idempotencyKey: "fast-one" }, ]);
const consumers = await new ThoughtAgentRuntime(store, [runner]).startConsumers([slow, fast]); await Promise.race([ fastStarted, new Promise<never>((_resolve, reject) => setTimeout(() => reject(new Error("Fast consumer was starved by the slow queue")), 1_000)), ]); await waitFor(async () => (await store.listConsumerProgress()).filter((progress) => ( [slow.id, fast.id].includes(progress.consumerId) )).length === 2); await consumers.stop();
expect(executionOrder.indexOf("fast-consumer:fast-one:start")) .toBeLessThan(executionOrder.indexOf("slow-consumer:slow-one:finish")); expect(executionOrder.indexOf("slow-consumer:slow-one:finish")) .toBeLessThan(executionOrder.indexOf("slow-consumer:slow-two:start")); expect((await store.listRuns()).filter((run) => run.status === "completed")).toHaveLength(3); });
test("keeps a live consumer ordered and recoverable after one subscription cycle throws", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const declaration = subscriptionDeclaration(); const runner: AgentRunner = { mode: "deterministic", run: async (input) => ({ summary: `Observed ${String(input.event.payload.title)}`, tags: ["rss"], importance: "normal", confidence: 1, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); const consumers = await runtime.startConsumers([declaration]); const originalGetProgress = store.getConsumerProgress.bind(store); let failed = false; let observeFailure!: () => void; const failureObserved = new Promise<void>((resolve) => { observeFailure = resolve; }); store.getConsumerProgress = async (id) => { if (!failed) { failed = true; observeFailure(); throw new Error("fixture transient progress read failure"); } return originalGetProgress(id); };
await store.appendProducerBatch([candidate("one")]); await failureObserved; await store.appendProducerBatch([candidate("two")]); await waitFor(async () => (await store.listConsumerProgress()).some((progress) => progress.lastSequence === 2)); await consumers.drain(); await consumers.stop();
expect((await store.listRuns()).map((run) => run.triggerEventId)).toHaveLength(2); expect((await store.listRuns()).every((run) => run.status === "completed")).toBe(true); });});
function candidate(id: string): EventCandidate { return { type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: id, idempotencyKey: id, occurredAt: `2026-07-14T00:00:0${id === "one" ? "0" : "1"}.000Z`, actor: "rss:fixture", correlationId: "poll-fixture", privacy: "private", payload: { title: id }, };}
function subscriptionDeclaration(): ThoughtAgentDeclaration { return { id: "rss-fixture-consumer", version: 1, name: "RSS fixture consumer", description: "Proves Jazz subscription delivery without producer handoff", mode: "deterministic", eventTypes: ["stream.thought.source.rss.item"], compiledEventTypes: ["stream.thought.source.rss.item"], sourcePatterns: ["rss:*"], acceptedPrivacy: ["private"], outputEventType: "stream.thought.derived.document.read", emit: ["stream.thought.derived.document.read"], promptRef: "fixture", systemPrompt: "Observe the fixture.", enabled: true, maxEvents: 1, maxInputChars: 10_000, maxOutputTokens: 100, timeoutMs: 1_000, tools: [], externalActions: false, };}
async function waitFor(predicate: () => Promise<boolean>, timeoutMs = 2_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) return; await new Promise((resolve) => setTimeout(resolve, 10)); } throw new Error("Timed out waiting for Jazz subscription delivery");}