Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
13 kB · 318 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test, vi } from "vitest";import { declarationFingerprint, loadAgentDeclarations } from "../src/agents/declarations.js";import { buildTelegramConversationContextPacket } from "../src/agents/context.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { CONVERSATION_COMPACTION_OUTPUT_CONTRACT } from "../src/agents/output-contracts.js";import type { AgentRunner, ThoughtAgentDeclaration } from "../src/agents/types.js";import { FilesystemConnector } from "../src/connectors/filesystem.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { rebuildRootActivity } from "../src/projections/activity.js";import { temporaryProject, 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("ThoughtAgentRuntime", () => { test("turns a file event into a derived event, trace, run, and activity lineage", async () => { const project = await temporaryProject(); roots.push(project); const vault = path.join(project, "vault"); const agents = path.join(project, "agents"); await fs.mkdir(path.join(vault, "spec"), { recursive: true }); await fs.mkdir(agents); await fs.mkdir(path.join(project, "prompts")); await fs.writeFile(path.join(vault, "spec", "brief.md"), "# Brief\n"); await fs.writeFile(path.join(project, "prompts", "triage.md"), "Classify.\n"); await fs.writeFile(path.join(agents, "triage.yaml"), [ "id: triage", "version: 1", "name: Triage", "description: Test triage", "enabled: true", "subscribe:", " types: [stream.thought.source.file.*]", " sources: [filesystem:test]", " privacy: [sensitive]", "context:", " maxEvents: 1", " maxChars: 64000", "runner:", " kind: deterministic", " maxOutputTokens: 2000", " timeoutMs: 60000", "prompt: prompts/triage.md", "emit: [stream.thought.derived.document.structure]", "policy:", " tools: []", " externalActions: false", ].join("\n")); const store = await testStore(project); stores.push(store); const runtime = new ThoughtAgentRuntime(store); const declarations = await loadAgentDeclarations(agents); const consumers = await runtime.startConsumers(declarations); const scan = await new FilesystemConnector({ id: "filesystem:test", root: vault }).scan(store); const sourceEvent = scan.events[0]; expect(sourceEvent).toBeDefined(); await waitFor(async () => (await store.listRuns()).some((run) => run.status === "completed")); await consumers.stop(); const result = (await store.listRuns())[0]!; const derived = (await store.listEvents()).find((event) => ( event.parentEventId === sourceEvent?.id && event.type === "stream.thought.derived.document.structure" ));
expect(result.errorText).toBeUndefined(); expect(result.result?.importance).toBe("high"); expect(result.result?.recommendation).toMatchObject({ target: "charter" }); expect(derived?.parentEventId).toBe(sourceEvent?.id); expect(derived?.rootEventId).toBe(sourceEvent?.id); expect((await store.listTrace(result.id))).toHaveLength(1); expect((await store.getRun(result.id))?.status).toBe("completed");
const activity = await rebuildRootActivity(store); expect(activity.totalEvents).toBe(5); expect(activity.items).toHaveLength(1); expect(activity.items[0]).toMatchObject({ id: sourceEvent?.id, type: "stream.thought.source.file.added", summary: "spec/brief.md added", descendantEventCount: 4, consumerRuns: [{ agentId: "triage", status: "completed", kind: "rule", outputCount: 1, }], }); expect(activity.items[0]?.consumerRuns[0]?.description) .toBe("Classified spec/brief.md as a specification change (high importance). Proposed: Re-evaluate Charter work affected by this specification change; do not edit files automatically. The proposal was not executed.");
const replay = await runtime.consumeBacklog(declarations); expect(replay).toHaveLength(0); expect(await store.listEvents()).toHaveLength(5); });
test("reconciles durable backlog when subscription wakeups are unavailable", async () => { const project = await temporaryProject(); roots.push(project); const vault = path.join(project, "vault"); const agents = path.join(project, "agents"); await fs.mkdir(path.join(vault, "spec"), { recursive: true }); await fs.mkdir(agents); await fs.mkdir(path.join(project, "prompts")); await fs.writeFile(path.join(vault, "spec", "brief.md"), "# Brief\n"); await fs.writeFile(path.join(project, "prompts", "triage.md"), "Classify.\n"); await fs.writeFile(path.join(agents, "triage.yaml"), [ "id: triage", "version: 1", "name: Triage", "description: Test triage without subscription callbacks", "enabled: true", "subscribe:", " types: [stream.thought.source.file.*]", " sources: [filesystem:test]", " privacy: [sensitive]", "context:", " maxEvents: 1", " maxChars: 64000", "runner:", " kind: deterministic", " maxOutputTokens: 2000", " timeoutMs: 60000", "prompt: prompts/triage.md", "emit: [stream.thought.derived.document.structure]", "policy:", " tools: []", " externalActions: false", ].join("\n")); const store = await testStore(project); stores.push(store); vi.spyOn(store, "subscribeConsumerEvents").mockImplementation(() => () => undefined); const declarations = await loadAgentDeclarations(agents); const runtime = new ThoughtAgentRuntime(store, undefined, { reconcileIntervalMs: 10 }); const consumers = await runtime.startConsumers(declarations);
await new FilesystemConnector({ id: "filesystem:test", root: vault }).scan(store); await waitFor(async () => (await store.listRuns()).some((run) => run.status === "completed")); await consumers.stop();
expect(await store.listRuns()).toHaveLength(1); expect((await store.listConsumerProgress())[0]).toMatchObject({ consumerId: "triage", source: "filesystem:test", lastSequence: 1, }); });
test("runs a separate no-tool compactor only at the boundary and exposes it to the parent on the next turn", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const parent = runtimeCompactionParent(); const compactor = runtimeCompactor(); parent.declarationFingerprint = declarationFingerprint(parent); compactor.declarationFingerprint = declarationFingerprint(compactor); let compactorCalls = 0; const runner: AgentRunner = { mode: "pi", run: async ({ declaration }) => { if ((declaration.role ?? "standard") === "compactor") { compactorCalls += 1; return { summary: "Boundary preview", boundary: "Boundary from the specialized clone.", openLoops: ["Wait for the next exact turn."], decisions: [], exactReferences: [], unresolved: [], lookupHints: ["boundary-preview"], confidence: 1, model: { provider: "tinker", id: "thinkingmachines/Inkling-Small" }, usage: { inputTokens: 120, outputTokens: 40 }, }; } return { summary: "Parent output", tags: ["conversation"], importance: "normal", confidence: 1 }; }, }; const runtime = new ThoughtAgentRuntime(store, [runner]); for (let index = 1; index <= 4; index += 1) { await store.appendEvent(runtimeTelegramMessage(index)); }
const results = await runtime.consumeBacklog([parent, compactor]); expect(results).toHaveLength(8); expect(compactorCalls).toBe(1); const compactorRuns = (await store.listRuns()).filter((run) => run.agentId === compactor.id); expect(compactorRuns.map((run) => run.status)).toEqual(["skipped", "skipped", "skipped", "completed"]); const boundaries = await store.listEvents({ types: ["stream.thought.derived.conversation.compaction"] }); expect(boundaries).toHaveLength(1); expect(boundaries[0]).toMatchObject({ source: `agent:${compactor.id}`, actor: compactor.id, privacy: "sensitive", payload: { compactionPlan: { targetAgentId: parent.id, coveredSourceEvents: 2, previousBoundaryEventId: null, }, }, });
const current = (await store.appendEvent(runtimeTelegramMessage(5))).event; const packet = await buildTelegramConversationContextPacket(parent, current, store); expect(packet.messages?.[0]).toMatchObject({ role: "compaction", content: expect.stringContaining("specialized clone"), }); expect(packet.manifest.conversationCompaction).toMatchObject({ boundaryEventId: boundaries[0]!.id }); });});
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 consumer completion");}
const runtimeAccounting = { leaseMs: 10_000, onExhaustion: "advance" as const, reservation: { inputTokens: 10_000, outputTokens: 1_000 }, limits: [{ window: "rolling" as const, durationMs: 60_000, maxCalls: 100, maxInputTokens: 1_000_000, maxOutputTokens: 100_000, }],};
function runtimeCompactionParent(): ThoughtAgentDeclaration { return { id: "telegram-conversation", version: 20, name: "Stream", description: "Fixture parent", mode: "pi", provider: "tinker", providerProfile: "tinker-default", model: "thinkingmachines/Inkling-Small", outputMode: "conversation-text", eventTypes: ["stream.thought.source.telegram.message"], compiledEventTypes: ["stream.thought.source.telegram.message"], sourcePatterns: ["telegram:fixture"], acceptedPrivacy: ["sensitive"], initialReplay: "beginning", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], promptRef: "prompts/parent.md", systemPrompt: "Parent.", enabled: true, maxEvents: 100, maxInputChars: 12_000, contextStrategy: "telegram-conversation", conversationHistoryAgentIds: ["telegram-conversation"], conversationAssistantHistoryMaxTurns: 2, conversationCompaction: { mode: "consume", agentId: "telegram-conversation-compactor" }, maxOutputTokens: 1_000, timeoutMs: 1_000, accounting: runtimeAccounting, tools: [], proposals: [], externalActions: false, };}
function runtimeCompactor(): ThoughtAgentDeclaration { return { ...runtimeCompactionParent(), id: "telegram-conversation-compactor", version: 1, name: "Stream Compactor", description: "Fixture compactor", role: "compactor", outputMode: "compaction-text", outputContract: { ...CONVERSATION_COMPACTION_OUTPUT_CONTRACT.identity }, outputEventType: "stream.thought.derived.conversation.compaction", emit: ["stream.thought.derived.conversation.compaction"], promptRef: "prompts/compactor.md", systemPrompt: "Compact.", contextStrategy: "telegram-compaction", conversationAssistantHistoryMaxTurns: undefined, conversationCompaction: { mode: "produce", targetAgentId: "telegram-conversation", triggerInputChars: 24, retainInputChars: 12, }, };}
function runtimeTelegramMessage(index: number) { return { type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:fixture", sourceKind: "telegram" as const, externalId: `message-${index}`, idempotencyKey: `message-${index}`, occurredAt: `2026-08-03T20:00:${String(index).padStart(2, "0")}.000Z`, actor: "cameron", correlationId: `message-${index}`, privacy: "sensitive" as const, payload: { chatId: "chat-one", senderId: "cameron", text: `User ${index}` }, };}