import 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((resolve) => { releaseSlow = resolve; }); const fastStarted = new Promise((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((_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((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, timeoutMs = 2_000): Promise { 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"); }