Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 301 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302import fs from "node:fs/promises";import { afterEach, describe, expect, test } from "vitest";import { auditRunEvidence } from "../src/agents/evidence.js";import { stableKey } from "../src/core/ids.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { AgentRunFailure, type AgentRunner, type ThoughtAgentDeclaration } from "../src/agents/types.js";import type { JazzThoughtStore } from "../src/jazz/store.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 failures", () => { test("rejects two enabled declarations that share one Letta Cloud agent", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const runtime = new ThoughtAgentRuntime(store, []);
await expect(runtime.registerDeclarations([ lettaDeclaration("resident-one"), lettaDeclaration("resident-two"), ])).rejects.toThrow( "Letta Cloud agent agent-shared-fixture is shared by enabled declarations resident-one and resident-two", ); });
test("validates semantic output separately from trusted execution metadata", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: "metadata-fixture", idempotencyKey: "metadata-fixture", occurredAt: "2026-07-14T00:00:00.000Z", actor: "rss:fixture", correlationId: "metadata-fixture", privacy: "private", payload: { title: "Metadata boundary fixture" }, }); const runner: AgentRunner = { mode: "deterministic", run: async ({ declaration }) => ({ summary: "Contract-valid semantic output", tags: ["fixture"], importance: "normal", confidence: 1, model: { provider: "fixture", id: "fixture-model", revision: "fixture-revision" }, enrichments: [{ tool: "fixture.read", status: "succeeded", requestKeys: ["target"], argumentsRedacted: true, resultPresent: true, }], ...(declaration.id === "unknown-semantic-field" ? { unexpected: "must remain invalid" } : {}), }), }; const runtime = new ThoughtAgentRuntime(store, [runner]);
const results = await runtime.consumeBacklog([ failureDeclaration("metadata-is-not-semantic-output"), failureDeclaration("unknown-semantic-field"), ]);
expect(results[0]).toMatchObject({ output: { summary: "Contract-valid semantic output", model: { provider: "fixture", id: "fixture-model", revision: "fixture-revision" }, enrichments: [expect.objectContaining({ tool: "fixture.read" })], }, derivedEvent: { type: "stream.thought.derived.document.read" }, }); expect(results[1]).toMatchObject({ error: "Agent final output rejected by canonical contract" }); expect((await store.listRuns()).find((run) => run.agentId === "unknown-semantic-field")?.result) .toMatchObject({ failureDiagnostic: { reason: "output-contract-invalid" } }); });
test("records consistent terminal evidence for provider, abort, invalid-output, and timeout failures", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const source = await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: "fixture-item", idempotencyKey: "fixture-item", occurredAt: "2026-07-14T00:00:00.000Z", actor: "rss:fixture", correlationId: "poll-fixture", privacy: "private", payload: { title: "Failure fixture" }, }); const failures = new Map([ ["provider-failure", { message: "Pi provider run failed", diagnostic: { code: "provider-run-failed", stage: "provider" } }], ["abort-failure", { message: "Pi provider run failed", diagnostic: { code: "provider-run-failed", stage: "provider" } }], ["invalid-output-failure", { message: "Pi final output rejected: expected one strict JSON object", diagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "invalid-json", textParts: 1, textChars: 8, textSha256: "fixture-hash", thinkingParts: 0, thinkingChars: 0, thinkingRedacted: true, toolCallParts: 0, otherParts: 0, assistantMessages: 1, stopReason: "stop", provider: "tinker", model: "Qwen/Qwen3.5-4B", checkpointRevision: "checkpoint-fixture", }, }], ["timeout-failure", { message: "Agent run timed out after 50ms", diagnostic: { code: "timeout", stage: "provider" } }], ]); const runner: AgentRunner = { mode: "deterministic", run: async ({ declaration }) => { const failure = failures.get(declaration.id)!; throw new AgentRunFailure(failure.message, { diagnostic: failure.diagnostic }); }, }; const runtime = new ThoughtAgentRuntime(store, [runner]); const results = await runtime.consumeBacklog([...failures.keys()].map(failureDeclaration));
expect(results).toHaveLength(4); expect(results.map((result) => result.error)).toEqual([...failures.values()].map((failure) => failure.message)); expect(results.every((result) => result.derivedEvent === undefined)).toBe(true); const runs = await store.listRuns(); expect(runs).toHaveLength(4); expect(runs.every((run) => run.status === "failed" && run.outputEventIds.length === 0)).toBe(true); const invalidRun = runs.find((run) => run.agentId === "invalid-output-failure"); expect(invalidRun?.result).toEqual({ failureDiagnostic: failures.get("invalid-output-failure")!.diagnostic, }); expect(invalidRun).toMatchObject({ provider: "tinker", model: "Qwen/Qwen3.5-4B", checkpointRevision: "checkpoint-fixture", executionAdapterRevision: "deterministic-triage-v1", }); const events = await store.listEvents(); const failedEvents = events.filter((event) => event.type === "stream.thought.agent.run.failed"); expect(failedEvents).toHaveLength(4); expect(new Set(failedEvents.map((event) => event.payload.error))) .toEqual(new Set([...failures.values()].map((failure) => failure.message))); expect(failedEvents.find((event) => event.actor === "invalid-output-failure")?.payload.failureDiagnostic) .toEqual(failures.get("invalid-output-failure")!.diagnostic); expect(JSON.stringify({ runs, failedEvents })).not.toContain("not json"); expect(auditRunEvidence(runs, events).every((report) => report.consistent)).toBe(true); });
test("does not advance past a sandbox failure and retries the same event on the next cycle", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); for (let index = 1; index <= 2; index += 1) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: `fixture-${index}`, idempotencyKey: `fixture-${index}`, occurredAt: `2026-07-14T00:00:0${index}.000Z`, actor: "rss:fixture", correlationId: "poll-fixture", privacy: "private", payload: { title: `Fixture ${index}` }, }); } let shouldFail = true; const runner: AgentRunner = { mode: "deterministic", run: async ({ event }) => { if (shouldFail) { shouldFail = false; throw new AgentRunFailure("Sandbox execution unavailable", { diagnostic: { code: "sandbox-unavailable", stage: "sandbox-execution" }, advanceProgress: false, }); } return { summary: `Processed ${event.externalId}`, tags: ["fixture"], importance: "normal", confidence: 1, }; }, }; const declaration = failureDeclaration("sandbox-retry"); const runtime = new ThoughtAgentRuntime(store, [runner]);
const firstCycle = await runtime.consumeBacklog([declaration]); expect(firstCycle).toEqual([ expect.objectContaining({ error: "Sandbox execution unavailable", retryable: true }), ]); const progressId = stableKey("consumer-progress", "sandbox-retry", "1", "rss:fixture"); expect(await store.getConsumerProgress(progressId)).toBeUndefined(); expect((await store.listRuns())).toEqual([ expect.objectContaining({ status: "failed", attempt: 1 }), ]);
const secondCycle = await runtime.consumeBacklog([declaration]); expect(secondCycle.map((result) => result.output?.summary)).toEqual([ "Processed fixture-1", "Processed fixture-2", ]); expect(await store.getConsumerProgress(progressId)).toMatchObject({ lastSequence: 2, lastEventId: expect.any(String), }); const runs = await store.listRuns(); expect(runs).toEqual(expect.arrayContaining([ expect.objectContaining({ status: "completed", attempt: 2, triggerEventId: expect.any(String) }), expect.objectContaining({ status: "completed", attempt: 1, triggerEventId: expect.any(String) }), ])); });});
function failureDeclaration(id: string): ThoughtAgentDeclaration { return { id, version: 1, name: id, description: `Fixture ${id}`, mode: "deterministic", eventTypes: ["stream.thought.source.rss.item"], compiledEventTypes: ["stream.thought.source.rss.item"], sourcePatterns: ["rss:fixture"], acceptedPrivacy: ["private"], outputEventType: "stream.thought.derived.document.read", emit: ["stream.thought.derived.document.read"], promptRef: "fixture", systemPrompt: "Fixture", enabled: true, maxEvents: 1, maxInputChars: 10_000, maxOutputTokens: 100, timeoutMs: 50, tools: [], externalActions: false, };}
function lettaDeclaration(id: string): ThoughtAgentDeclaration { return { id, version: 1, name: id, description: `Fixture ${id}`, mode: "letta-agent-sdk", provider: "letta-cloud", lettaAgent: { backend: "cloud", agentIdEnv: "THOUGHTSTREAM_LETTA_SHARED_AGENT_ID", agentId: "agent-shared-fixture", conversation: "main", responseMode: "conversation-text", outputOnly: false, permissionMode: "unrestricted", dreaming: { trigger: "off" }, sandbox: { ttlMinutes: 5, terminateOnClose: false }, }, eventTypes: ["stream.thought.source.telegram.message"], compiledEventTypes: ["stream.thought.source.telegram.message"], sourcePatterns: [`telegram:${id}`], acceptedPrivacy: ["sensitive"], outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], promptRef: "fixture", systemPrompt: "Fixture", enabled: true, maxEvents: 1, maxInputChars: 1_000, contextStrategy: "single-event", maxOutputTokens: 1_000, timeoutMs: 60_000, tools: [], externalActions: false, };}