import 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 => (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 => { 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 { 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)); } }