Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
19 kB · 499 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500import fs from "node:fs/promises";import { afterEach, describe, expect, test } from "vitest";import type { EventCandidate, ThoughtEvent } from "../src/events/types.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import type { AgentRun, ConsumerProgress } from "../src/store/types.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("JazzThoughtStore", () => { test("serializes concurrent producer appends to one source", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const candidate = (index: number): EventCandidate => ({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:concurrent-producer", sourceKind: "rss", externalId: `item-${index}`, idempotencyKey: `item-${index}`, occurredAt: `2026-07-27T00:00:0${index}.000Z`, actor: "rss:concurrent-producer", correlationId: "concurrent-producer", privacy: "public-source", payload: { index }, });
const results = await Promise.all([store.appendEvent(candidate(1)), store.appendEvent(candidate(2))]); expect(results.map((result) => result.event.sourceSequence).sort((left, right) => left - right)).toEqual([1, 2]); expect((await store.listEvents({ source: "rss:concurrent-producer" }))).toHaveLength(2); });
test("appends one durable event for repeated source identity", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const candidate = { type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:test", sourceKind: "rss" as const, externalId: "item-1", idempotencyKey: "feed:item-1:hash", occurredAt: "2026-07-14T08:00:00.000Z", observedAt: "2026-07-14T08:00:01.000Z", actor: "https://example.com/feed", correlationId: "poll-1", privacy: "public-source" as const, payload: { title: "One" }, };
const first = await store.appendEvent(candidate); const second = await store.appendEvent(candidate);
expect(first.inserted).toBe(true); expect(second.inserted).toBe(false); expect(second.event.id).toBe(first.event.id); expect(await store.listEvents()).toHaveLength(1); expect((await store.getEvent(first.event.id))?.payload).toEqual({ title: "One" }); await expect(store.appendProducerBatch([ candidate, { ...candidate, payload: { title: "Conflicting replay" } }, ])).rejects.toThrow("different payload"); expect((await store.listSources()).find((source) => source.id === "rss:test")?.lastSequence).toBe(1); });
test("limits event queries in storage while preserving chronological output", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); for (const index of [1, 2, 3]) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:bounded", sourceKind: "rss", externalId: `item-${index}`, idempotencyKey: `item-${index}`, occurredAt: `2026-08-22T10:00:0${index}.000Z`, observedAt: "2026-08-22T10:00:04.000Z", actor: "rss:bounded", correlationId: "bounded-query", privacy: "public-source", payload: { index }, }); }
const recent = await store.listEvents({ types: ["stream.thought.source.rss.item"], limit: 2, }); const all = await store.listEvents({ types: ["stream.thought.source.rss.item"] }); expect(recent.map((event) => event.id)).toEqual(all.slice(-2).map((event) => event.id));
const rootEvent = all[0]!; for (const index of [4, 5, 6]) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:bounded", sourceKind: "rss", externalId: `child-${index}`, idempotencyKey: `child-${index}`, occurredAt: `2026-08-22T10:00:0${index}.000Z`, observedAt: `2026-08-22T10:00:0${index}.000Z`, actor: "rss:bounded", rootEventId: rootEvent.id, parentEventId: rootEvent.id, correlationId: "bounded-query", privacy: "public-source", payload: { index }, }); } const incomplete = await store.listRecentRootEvents({ types: ["stream.thought.source.rss.item"], limit: 2, maxScanned: 2, }); expect(incomplete).toMatchObject({ events: [], complete: false, scannedEvents: 2 }); const recovered = await store.listRecentRootEvents({ types: ["stream.thought.source.rss.item"], limit: 2, maxScanned: 10, }); expect(recovered.complete).toBe(true); expect(recovered.events).toHaveLength(2); });
test("persists events across a process-shaped reopen", async () => { const root = await temporaryProject(); roots.push(root); const appId = "thoughtstream-reopen-test"; const firstStore = await (await import("../src/jazz/store.js")).JazzThoughtStore.open({ projectRoot: root, appId }); await firstStore.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:test", sourceKind: "rss", externalId: "item-1", idempotencyKey: "item-1", occurredAt: "2026-07-14T08:00:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source", payload: { title: "Persistent" }, }); await firstStore.close();
const secondStore = await (await import("../src/jazz/store.js")).JazzThoughtStore.open({ projectRoot: root, appId }); stores.push(secondStore); const appended = await secondStore.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:test", sourceKind: "rss", externalId: "item-2", idempotencyKey: "item-2", occurredAt: "2026-07-14T08:01:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source", payload: { title: "Appended after reopen" }, }); expect(appended.inserted).toBe(true); const events = await secondStore.listEvents(); expect(events).toHaveLength(2); expect(events.map((event) => event.payload.title)).toEqual(["Persistent", "Appended after reopen"]); expect((await secondStore.listSources()).find((source) => source.id === "rss:test")?.lastSequence).toBe(2); });
test("attaches a second concurrent store to the first process's local server", async () => { // alpha55 storage takes an exclusive RocksDB lock, so a second process // cannot open the same dataPath. A concurrent open must attach as a client // to the first process's local server through the discovery file and see // its writes, then the first store must remain usable until it closes. const root = await temporaryProject(); roots.push(root); const appId = "thoughtstream-attach-test"; const firstStore = await (await import("../src/jazz/store.js")).JazzThoughtStore.open({ projectRoot: root, appId }); await firstStore.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:attach", sourceKind: "rss", externalId: "item-1", idempotencyKey: "item-1", occurredAt: "2026-07-14T08:00:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source", payload: { title: "From owner" }, }); const secondStore = await (await import("../src/jazz/store.js")).JazzThoughtStore.open({ projectRoot: root, appId }); stores.push(secondStore); const events = await secondStore.listEvents(); expect(events.map((event) => event.payload.title)).toEqual(["From owner"]); await secondStore.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:attach", sourceKind: "rss", externalId: "item-2", idempotencyKey: "item-2", occurredAt: "2026-07-14T08:01:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source", payload: { title: "From attacher" }, }); const ownerEvents = await firstStore.listEvents(); expect(ownerEvents.map((event) => event.payload.title)).toEqual(["From owner", "From attacher"]); await firstStore.close(); });
test("persists checkpoint provenance without changing the Jazz schema", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); await store.upsertRun({ id: "run:test", executionKey: "execution:test", triggerEventId: "event:test", agentId: "agent:test", agentVersion: 1, status: "completed", inputEventIds: ["event:test"], outputEventIds: ["event:output"], attempt: 1, provider: "tinker", model: "openai/gpt-oss-120b", checkpointRevision: "checkpoint:test", promptHash: "prompt:test", contextManifest: { input: "test" }, result: { content: "ok" }, createdAt: "2026-07-15T08:00:00.000Z", completedAt: "2026-07-15T08:01:00.000Z", updatedAt: "2026-07-15T08:01:00.000Z", });
expect(await store.getRun("run:test")).toMatchObject({ checkpointRevision: "checkpoint:test", result: { content: "ok" }, }); });
test("recovers consumer settlement after a failed transaction leaves allocated row ids", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const input = async (index: number): Promise<ThoughtEvent> => (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:consumer-residue", sourceKind: "rss", externalId: `item-${index}`, idempotencyKey: `item-${index}`, occurredAt: `2026-07-27T00:00:0${index}.000Z`, actor: "rss:consumer-residue", correlationId: "consumer-residue", privacy: "sensitive", payload: { index }, })).event; const firstInput = await input(1); const secondInput = await input(2); const progressId = "progress:consumer-residue"; await store.initializeConsumerProgress({ id: progressId, consumerId: "consumer-residue", consumerVersion: 1, source: firstInput.source, lastSequence: firstInput.sourceSequence, lastEventId: "evt_conflicting_progress_owner", updatedAt: "2026-07-27T00:01:00.000Z", }); const run = async (inputEvent: ThoughtEvent, suffix: string): Promise<AgentRun> => { const value: AgentRun = { id: `run:consumer-residue:${suffix}`, executionKey: `execution:consumer-residue:${suffix}`, triggerEventId: inputEvent.id, agentId: "consumer-residue", agentVersion: 1, status: "skipped", inputEventIds: [inputEvent.id], outputEventIds: [], attempt: 1, provider: "deterministic", model: "deterministic", privacy: "sensitive", promptHash: "prompt:consumer-residue", contextManifest: { skip: { code: "fixture" } }, errorText: "fixture skip", createdAt: "2026-07-27T00:01:00.000Z", startedAt: "2026-07-27T00:01:00.000Z", completedAt: "2026-07-27T00:01:00.000Z", updatedAt: "2026-07-27T00:01:00.000Z", }; await store.upsertRun(value); return value; }; const progress = (inputEvent: ThoughtEvent): ConsumerProgress => ({ id: progressId, consumerId: "consumer-residue", consumerVersion: 1, source: inputEvent.source, lastSequence: inputEvent.sourceSequence, lastEventId: inputEvent.id, updatedAt: "2026-07-27T00:01:00.000Z", }); const terminal = (value: AgentRun, inputEvent: ThoughtEvent): EventCandidate => ({ type: "stream.thought.agent.run.skipped", schemaVersion: 1, source: "agent:consumer-residue", sourceKind: "agent", externalId: value.id, idempotencyKey: `${value.id}:skipped`, occurredAt: "2026-07-27T00:01:00.000Z", actor: "consumer-residue", correlationId: value.id, privacy: "sensitive", payload: { runId: value.id, agentId: value.agentId, agentVersion: value.agentVersion, inputEventIds: [inputEvent.id], attempt: value.attempt, status: "skipped", }, });
const firstRun = await run(firstInput, "first"); await expect(store.settleConsumerFailure({ run: firstRun, inputEvent: firstInput, terminal: terminal(firstRun, firstInput), progress: progress(firstInput), })).rejects.toThrow("Consumer progress sequence conflict");
const secondRun = await run(secondInput, "second"); const receipt = await store.settleConsumerFailure({ run: secondRun, inputEvent: secondInput, terminal: terminal(secondRun, secondInput), progress: progress(secondInput), });
expect(receipt.type).toBe("stream.thought.agent.run.skipped"); expect(await store.getConsumerProgress(progressId)).toMatchObject({ lastSequence: secondInput.sourceSequence, lastEventId: secondInput.id, }); expect((await store.listSources()).find((source) => source.id === "agent:consumer-residue")).toMatchObject({ lastSequence: 1, }); // alpha55: failed transactions roll back atomically, so no storage-revision // retry is required and the revision stays at its initial value. (alpha53 // left allocated row ids behind and required a revision bump to recover.) expect((store as unknown as { consumerStorageRevision: number }).consumerStorageRevision).toBe(5); });
test("delivers only newly added rows from a consumer subscription", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const deliveries: string[][] = []; const unsubscribe = store.subscribeConsumerEvents({ consumerId: "structure", consumerVersion: 1, source: "rss:test", eventTypes: ["stream.thought.source.rss.item"], acceptedPrivacy: ["public-source"], afterSequence: 0, }, (events) => { if (events.length > 0) deliveries.push(events.map((event) => event.externalId)); }); try { const base = { type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:test", sourceKind: "rss" as const, occurredAt: "2026-07-14T08:00:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source" as const, payload: { title: "item" }, }; await store.appendEvent({ ...base, externalId: "item-1", idempotencyKey: "item-1" }); await waitUntil(() => deliveries.length === 1); await store.appendEvent({ ...base, externalId: "item-2", idempotencyKey: "item-2" }); await waitUntil(() => deliveries.length === 2); expect(deliveries).toEqual([["item-1"], ["item-2"]]); } finally { unsubscribe(); } });
test("delivers the initial snapshot as the first consumer subscription delta", async () => { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const base = { type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:initial-snapshot", sourceKind: "rss" as const, occurredAt: "2026-07-14T08:00:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source" as const, payload: { title: "item" }, }; await store.appendEvent({ ...base, externalId: "existing-1", idempotencyKey: "existing-1" }); await store.appendEvent({ ...base, externalId: "existing-2", idempotencyKey: "existing-2" }); const deliveries: string[][] = []; const unsubscribe = store.subscribeConsumerEvents({ consumerId: "initial-snapshot", consumerVersion: 1, source: "rss:initial-snapshot", eventTypes: ["stream.thought.source.rss.item"], acceptedPrivacy: ["public-source"], afterSequence: 0, }, (events) => { if (events.length > 0) deliveries.push(events.map((event) => event.externalId)); }); try { await waitUntil(() => deliveries.length === 1); expect(deliveries).toEqual([["existing-1", "existing-2"]]); } finally { unsubscribe(); } });
test("does not redeliver updated rows that a consumer subscription already delivered", async () => { // Events are append-only, so consumer subscriptions never observe updates // to already-delivered ids in practice; this test pins the contract that a // full-result subscription diff keyed by row id treats an update to an // already-delivered id as not-new, while a later insert is still new. const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const deliveries: string[][] = []; const unsubscribe = store.subscribeConsumerEvents({ consumerId: "update-diff", consumerVersion: 1, source: "rss:update-diff", eventTypes: ["stream.thought.source.rss.item"], acceptedPrivacy: ["public-source"], afterSequence: 0, }, (events) => { if (events.length > 0) deliveries.push(events.map((event) => event.externalId)); }); try { const base = { type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:update-diff", sourceKind: "rss" as const, occurredAt: "2026-07-14T08:00:00.000Z", actor: "test", correlationId: "poll", privacy: "public-source" as const, payload: { title: "item" }, }; await store.appendEvent({ ...base, externalId: "item-1", idempotencyKey: "item-1" }); await waitUntil(() => deliveries.length === 1); // Re-appending the same idempotency key is a no-op append (unchanged), // which surfaces as a full-snapshot update without a new row id. await store.appendEvent({ ...base, externalId: "item-1", idempotencyKey: "item-1" }); await store.appendEvent({ ...base, externalId: "item-2", idempotencyKey: "item-2" }); await waitUntil(() => deliveries.length === 2); expect(deliveries).toEqual([["item-1"], ["item-2"]]); } finally { unsubscribe(); } });});
async function waitUntil(predicate: () => boolean, timeoutMs = 2_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (!predicate()) { if (Date.now() >= deadline) throw new Error("Timed out waiting for Jazz subscription delivery"); await new Promise((resolve) => setTimeout(resolve, 10)); }}