Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
22 kB · 563 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564import fs from "node:fs/promises";import { afterEach, describe, expect, test } from "vitest";import { outputContractForDeclaration } from "../src/agents/output-contracts.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { AgentRunFailure, type AgentRunner, type ThoughtAgentDeclaration } from "../src/agents/types.js";import { stableKey } from "../src/core/ids.js";import { JazzThoughtStore } from "../src/jazz/store.js";import type { InferenceReservationRequest } from "../src/store/types.js";import { temporaryProject, testInferenceAccountingPolicy, 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("durable inference accounting", () => { test("atomically reserves concurrent calls and isolates agent accounts", async () => { const { store } = await fixtureStore(); const at = "2026-07-16T00:00:00.000Z"; const [left, right] = await Promise.all([ store.reserveInference(reservation("reserve-left", "run-left", "agent-a", at, 1)), store.reserveInference(reservation("reserve-right", "run-right", "agent-a", at, 1)), ]);
expect([left.approved, right.approved].sort()).toEqual([false, true]); expect([left.record.status, right.record.status].sort()).toEqual(["denied", "reserved"]); expect([left, right].find((decision) => !decision.approved)?.record.charged).toEqual({ calls: 0, inputTokens: 0, outputTokens: 0, costMicrousd: 0, }); const isolated = await store.reserveInference(reservation("reserve-other", "run-other", "agent-b", at, 1)); expect(isolated).toMatchObject({ approved: true, acquired: true, record: { status: "reserved" } }); expect(await store.listInferenceAccounting()).toHaveLength(3); });
test("preserves cost-bearing settlement telemetry across a store restart", async () => { const project = await temporaryProject("thoughtstream-accounting-restart-"); roots.push(project); const options = { projectRoot: project, appId: "thoughtstream-accounting-restart", runtimeRevision: "test" }; const first = await JazzThoughtStore.open(options); stores.push(first); const decision = await first.reserveInference(reservation( "persistent-reservation", "persistent-run", "persistent-agent", "2026-07-16T01:00:00.000Z", 2, )); expect(decision.approved).toBe(true); await first.settleInferenceReservation("persistent-reservation", { inputTokens: 321, outputTokens: 45, costMicrousd: 6_789, }, "2026-07-16T01:00:01.000Z"); await first.close(); stores.splice(stores.indexOf(first), 1);
const restarted = await JazzThoughtStore.open(options); stores.push(restarted); expect(await restarted.getInferenceAccountingRecord("persistent-reservation")).toMatchObject({ status: "settled", usageStatus: "reported", actualUsage: { inputTokens: 321, outputTokens: 45, costMicrousd: 6_789 }, charged: { calls: 1, inputTokens: 321, outputTokens: 45, costMicrousd: 6_789 }, }); });
test("persists a fully reported token-only settlement without serializing cost", async () => { const project = await temporaryProject("thoughtstream-accounting-token-only-restart-"); roots.push(project); const options = { projectRoot: project, appId: "thoughtstream-accounting-token-only-restart", runtimeRevision: "test" }; const first = await JazzThoughtStore.open(options); stores.push(first); const decision = await first.reserveInference(tokenOnlyReservation( "token-only-persistent-reservation", "token-only-persistent-run", "token-only-persistent-agent", "2026-07-16T01:30:00.000Z", 2, )); expect(decision).toMatchObject({ approved: true, record: { estimate: { calls: 1, inputTokens: 1_000, outputTokens: 100 }, charged: { calls: 1, inputTokens: 1_000, outputTokens: 100 }, }, }); expect(decision.record.estimate).not.toHaveProperty("costMicrousd"); expect(decision.record.charged).not.toHaveProperty("costMicrousd"); await first.settleInferenceReservation("token-only-persistent-reservation", { inputTokens: 321, outputTokens: 45, }, "2026-07-16T01:30:01.000Z"); await first.close(); stores.splice(stores.indexOf(first), 1);
const restarted = await JazzThoughtStore.open(options); stores.push(restarted); const record = await restarted.getInferenceAccountingRecord("token-only-persistent-reservation"); expect(record).toMatchObject({ status: "settled", usageStatus: "reported", actualUsage: { inputTokens: 321, outputTokens: 45 }, estimate: { calls: 1, inputTokens: 1_000, outputTokens: 100 }, charged: { calls: 1, inputTokens: 321, outputTokens: 45 }, }); expect(record?.estimate).not.toHaveProperty("costMicrousd"); expect(record?.charged).not.toHaveProperty("costMicrousd"); expect(record?.actualUsage).not.toHaveProperty("costMicrousd"); expect(JSON.stringify(record)).not.toContain("costMicrousd"); });
test("rejects a cost limit when the reservation does not track cost", async () => { const { store } = await fixtureStore(); const request = tokenOnlyReservation( "invalid-cost-policy", "invalid-cost-policy-run", "invalid-cost-policy-agent", "2026-07-16T01:40:00.000Z", 2, ); request.policy = { ...request.policy, limits: request.policy.limits.map((limit) => ({ ...limit, maxCostMicrousd: 10_000 })), };
await expect(store.reserveInference(request)).rejects.toThrow("Inference cost limits require a cost reservation"); expect(await store.listInferenceAccounting()).toEqual([]); });
test("enforces token-only rolling windows while preserving omitted cost", async () => { const { store } = await fixtureStore(); const reserveAt = async (id: string, at: string) => { const decision = await store.reserveInference(tokenOnlyReservation( id, `run-${id}`, "token-only-sliding-agent", at, 2, )); if (decision.approved) { await store.settleInferenceReservation(id, { inputTokens: 500, outputTokens: 50 }, at); } return decision; };
expect((await reserveAt("token-sliding-one", "2026-07-16T02:10:00.000Z")).approved).toBe(true); expect((await reserveAt("token-sliding-two", "2026-07-16T02:10:59.000Z")).approved).toBe(true); expect((await reserveAt("token-sliding-three", "2026-07-16T02:11:01.000Z")).approved).toBe(true); expect((await reserveAt("token-sliding-four", "2026-07-16T02:11:02.000Z")).approved).toBe(false);
const records = await store.listInferenceAccounting({ agentId: "token-only-sliding-agent" }); expect(records.map((record) => record.status).sort()).toEqual(["denied", "settled", "settled", "settled"]); expect(records.filter((record) => record.status === "settled").every((record) => record.usageStatus === "reported")).toBe(true); for (const record of records) { expect(record.estimate).not.toHaveProperty("costMicrousd"); expect(record.charged).not.toHaveProperty("costMicrousd"); } });
test("resets elapsed windows while retaining conservative charges for expired leases", async () => { const { store } = await fixtureStore(); const firstAt = "2026-07-16T02:00:00.000Z"; const first = reservation("lease-one", "lease-run-one", "lease-agent", firstAt, 1); first.policy = { ...first.policy, leaseMs: 10_000, limits: [{ window: "rolling", durationMs: 60_000, maxCalls: 1, maxInputTokens: 1_000, maxOutputTokens: 100, maxCostMicrousd: 10_000, }], }; expect((await store.reserveInference(first)).approved).toBe(true);
const beforeReset = await store.reserveInference({ ...reservation("lease-two", "lease-run-two", "lease-agent", "2026-07-16T02:00:11.000Z", 1), policy: first.policy, }); expect(beforeReset.approved).toBe(false); expect(await store.getInferenceAccountingRecord("lease-one")).toMatchObject({ status: "expired", usageStatus: "unavailable", charged: first.estimate, });
const afterReset = await store.reserveInference({ ...reservation("lease-three", "lease-run-three", "lease-agent", "2026-07-16T02:01:01.000Z", 1), policy: first.policy, }); expect(afterReset.approved).toBe(true); });
test("enforces a true sliding rolling window instead of a first-call bucket", async () => { const { store } = await fixtureStore(); const policy = { ...testInferenceAccountingPolicy(2), limits: [{ window: "rolling" as const, durationMs: 60_000, maxCalls: 2, maxInputTokens: 2_000, maxOutputTokens: 200, maxCostMicrousd: 20_000, }], }; const reserveAt = (id: string, at: string) => store.reserveInference({ ...reservation(id, `run-${id}`, "sliding-agent", at, 2), policy, });
expect((await reserveAt("sliding-one", "2026-07-16T02:10:00.000Z")).approved).toBe(true); expect((await reserveAt("sliding-two", "2026-07-16T02:10:59.000Z")).approved).toBe(true); expect((await reserveAt("sliding-three", "2026-07-16T02:11:01.000Z")).approved).toBe(true); expect((await reserveAt("sliding-four", "2026-07-16T02:11:02.000Z")).approved).toBe(false); });
test("blocks exhausted runs before runner dispatch and advances progress once", async () => { const { store } = await fixtureStore(); for (let sequence = 1; sequence <= 2; sequence += 1) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:accounting-fixture", sourceKind: "rss", externalId: `private-item-${sequence}`, idempotencyKey: `accounting-item-${sequence}`, occurredAt: `2026-07-16T03:00:0${sequence}.000Z`, actor: "rss:accounting-fixture", correlationId: "accounting-fixture", privacy: "private", payload: { title: `PRIVATE_SOURCE_BODY_${sequence}` }, }); } let runnerCalls = 0; const runner: AgentRunner = { mode: "pi", run: async () => { runnerCalls += 1; return { summary: "Budgeted result", tags: ["accounting"], importance: "normal", confidence: 1, usage: { inputTokens: 200, outputTokens: 20, costMicrousd: 2_000 }, }; }, }; const declaration = accountingDeclaration(); const runtime = new ThoughtAgentRuntime(store, [runner]);
const results = await runtime.consumeBacklog([declaration]);
expect(results).toHaveLength(2); expect(results[0]).toMatchObject({ output: { summary: "Budgeted result" } }); expect(results[1]).toMatchObject({ error: "Inference budget exhausted before provider dispatch" }); expect(runnerCalls).toBe(1); const runs = await store.listRuns(); expect(runs.map((run) => run.status).sort()).toEqual(["blocked", "completed"]); expect(runs.every((run) => Boolean(run.accountingReservationId))).toBe(true); const progress = await store.getConsumerProgress(stableKey( "consumer-progress", declaration.id, String(declaration.version), "rss:accounting-fixture", )); expect(progress?.lastSequence).toBe(2); const events = await store.listEvents(); expect(events.filter((event) => event.type === "stream.thought.agent.run.blocked")).toHaveLength(1); expect(events.filter((event) => event.type === "stream.thought.agent.run.failed")).toHaveLength(0); expect(events.filter((event) => event.type === "stream.thought.agent.repair.requested")).toHaveLength(0); expect(await runtime.consumeBacklog([declaration])).toEqual([]); expect(runnerCalls).toBe(1);
const accounting = await store.listInferenceAccounting({ agentId: declaration.id }); expect(accounting.map((record) => record.status).sort()).toEqual(["denied", "settled"]); expect(accounting.find((record) => record.status === "settled")).toMatchObject({ usageStatus: "reported", charged: { calls: 1, inputTokens: 200, outputTokens: 20, costMicrousd: 2_000 }, }); expect(JSON.stringify(accounting)).not.toContain("PRIVATE_SOURCE_BODY"); expect(JSON.stringify(accounting)).not.toContain(declaration.systemPrompt); });
test("defers budget-blocked conversation work without advancing progress and retries after the window", async () => { const { store } = await fixtureStore(); for (let sequence = 1; sequence <= 2; sequence += 1) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:accounting-fixture", sourceKind: "rss", externalId: `deferred-item-${sequence}`, idempotencyKey: `deferred-item-${sequence}`, occurredAt: `2026-07-16T03:05:0${sequence}.000Z`, actor: "rss:accounting-fixture", correlationId: "deferred-accounting-fixture", privacy: "private", payload: { title: `DEFERRED_PRIVATE_BODY_${sequence}` }, }); } let runnerCalls = 0; const runner: AgentRunner = { mode: "pi", run: async () => { runnerCalls += 1; return { summary: "Deferred budget result", tags: ["accounting"], importance: "normal", confidence: 1, usage: { inputTokens: 200, outputTokens: 20, costMicrousd: 2_000 }, }; }, }; const base = accountingDeclaration(); const declaration: ThoughtAgentDeclaration = { ...base, id: "deferred-accounting-observer", accounting: { ...base.accounting!, onExhaustion: "defer", limits: [{ window: "rolling", durationMs: 1_000, maxCalls: 1, maxInputTokens: 1_000, maxOutputTokens: 100, maxCostMicrousd: 10_000, }], }, }; const runtime = new ThoughtAgentRuntime(store, [runner]);
const initial = await runtime.consumeBacklog([declaration]); expect(initial).toHaveLength(2); expect(initial[1]).toMatchObject({ error: "Inference budget exhausted before provider dispatch", retryable: true, }); expect(runnerCalls).toBe(1); const progressId = stableKey( "consumer-progress", declaration.id, String(declaration.version), "rss:accounting-fixture", ); expect((await store.getConsumerProgress(progressId))?.lastSequence).toBe(1); const blocked = (await store.listRuns()).find((run) => run.status === "blocked"); expect(blocked?.result?.failureDiagnostic).toMatchObject({ code: "inference-budget-exhausted", progressDisposition: "deferred", retryAt: expect.any(String), });
const stillDeferred = await runtime.consumeBacklog([declaration]); expect(stillDeferred).toEqual([expect.objectContaining({ runId: blocked?.id, retryable: true })]); expect(runnerCalls).toBe(1); expect(await store.listRuns()).toHaveLength(2);
await new Promise((resolve) => setTimeout(resolve, 1_100)); const recovered = await runtime.consumeBacklog([declaration]); expect(recovered).toEqual([expect.objectContaining({ output: expect.objectContaining({ summary: "Deferred budget result" }), })]); expect(runnerCalls).toBe(2); expect((await store.getConsumerProgress(progressId))?.lastSequence).toBe(2); expect((await store.listRuns()).map((run) => run.status).sort()).toEqual(["blocked", "completed", "completed"]); });
test("backs off retryable model failures through retry policy independently of budget exhaustion policy", async () => { const { store } = await fixtureStore(); await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:accounting-fixture", sourceKind: "rss", externalId: "deferred-failure-item", idempotencyKey: "deferred-failure-item", occurredAt: "2026-07-16T03:07:00.000Z", actor: "rss:accounting-fixture", correlationId: "deferred-failure-fixture", privacy: "private", payload: { title: "DEFERRED_FAILURE_PRIVATE_BODY" }, }); let runnerCalls = 0; const runner: AgentRunner = { mode: "pi", run: async () => { runnerCalls += 1; throw new AgentRunFailure("Retryable sandbox failure", { advanceProgress: false, diagnostic: { code: "sandbox-worker-model-failed", stage: "sandbox-execution" }, }); }, }; const base = accountingDeclaration(); const declaration: ThoughtAgentDeclaration = { ...base, id: "deferred-failure-observer", accounting: { ...base.accounting!, onExhaustion: "advance" }, retry: { initialDelayMs: 5_000, maxDelayMs: 300_000 }, }; const runtime = new ThoughtAgentRuntime(store, [runner]);
const [failed] = await runtime.consumeBacklog([declaration]); expect(failed).toMatchObject({ error: "Retryable sandbox failure", retryable: true }); expect(runnerCalls).toBe(1); const [run] = await store.listRuns(); expect(run?.result?.failureDiagnostic).toMatchObject({ code: "sandbox-worker-model-failed", progressDisposition: "retry-delayed", retryAt: expect.any(String), });
const [held] = await runtime.consumeBacklog([declaration]); expect(held).toMatchObject({ runId: run?.id, retryable: true }); expect(runnerCalls).toBe(1); expect(await store.listRuns()).toHaveLength(1); });
test("settles provider usage exposed by a failed run", async () => { const { store } = await fixtureStore(); await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:accounting-fixture", sourceKind: "rss", externalId: "failed-usage-item", idempotencyKey: "failed-usage-item", occurredAt: "2026-07-16T03:10:00.000Z", actor: "rss:accounting-fixture", correlationId: "accounting-fixture-failed", privacy: "private", payload: { title: "PRIVATE_FAILED_SOURCE_BODY" }, }); const declaration = accountingDeclaration(); const runner: AgentRunner = { mode: "pi", run: async () => { throw new AgentRunFailure("Contract-invalid provider output", { usage: { inputTokens: 250, outputTokens: 30 }, diagnostic: { code: "invalid-final-output", stage: "final-output-validation" }, }); }, };
const [result] = await new ThoughtAgentRuntime(store, [runner]).consumeBacklog([declaration]);
expect(result).toMatchObject({ error: "Contract-invalid provider output" }); const [record] = await store.listInferenceAccounting({ agentId: declaration.id }); expect(record).toMatchObject({ status: "settled", usageStatus: "partial", charged: { calls: 1, inputTokens: 250, outputTokens: 30, costMicrousd: 10_000 }, }); expect(JSON.stringify(record)).not.toContain("PRIVATE_FAILED_SOURCE_BODY"); });});
async function fixtureStore(): Promise<{ project: string; store: JazzThoughtStore }> { const project = await temporaryProject("thoughtstream-accounting-"); roots.push(project); const store = await testStore(project); stores.push(store); return { project, store };}
function reservation( reservationId: string, runId: string, agentId: string, reservedAt: string, maxCalls: number,): InferenceReservationRequest { const policy = testInferenceAccountingPolicy(maxCalls); return { reservationId, scopeType: "agent", scopeKey: agentId, runId, agentId, agentVersion: 1, provider: "tinker", model: "fixture-model", policy, estimate: { calls: 1, ...policy.reservation }, reservedAt, };}
function tokenOnlyReservation( reservationId: string, runId: string, agentId: string, reservedAt: string, maxCalls: number,): InferenceReservationRequest { const policy = { leaseMs: 180_000, reservation: { inputTokens: 1_000, outputTokens: 100 }, limits: [{ window: "rolling" as const, durationMs: 60_000, maxCalls, maxInputTokens: maxCalls * 1_000, maxOutputTokens: maxCalls * 100, }], }; return { reservationId, scopeType: "agent", scopeKey: agentId, runId, agentId, agentVersion: 1, provider: "letta-cloud", model: "fixture-oauth-model", policy, estimate: { calls: 1, ...policy.reservation }, reservedAt, };}
function accountingDeclaration(): ThoughtAgentDeclaration { return { id: "accounting-observer", version: 1, name: "Accounting observer", description: "Exercises pre-dispatch inference accounting", mode: "pi", role: "standard", outputContract: outputContractForDeclaration({} as ThoughtAgentDeclaration), provider: "tinker", providerProfile: "tinker-default", model: "fixture-model", eventTypes: ["stream.thought.source.rss.item"], compiledEventTypes: ["stream.thought.source.rss.item"], sourcePatterns: ["rss:accounting-fixture"], acceptedPrivacy: ["private"], initialReplay: "beginning", outputEventType: "stream.thought.derived.document.read", emit: ["stream.thought.derived.document.read"], promptRef: "prompts/accounting-fixture.md", systemPrompt: "PRIVATE_AGENT_PROMPT", enabled: true, maxEvents: 1, maxInputChars: 1_000, maxOutputTokens: 100, timeoutMs: 5_000, accounting: testInferenceAccountingPolicy(1), tools: [], externalActions: false, };}