diff --git a/src/jazz/store.ts b/src/jazz/store.ts index f009515..ba36e18 100644 --- a/src/jazz/store.ts +++ b/src/jazz/store.ts @@ -66,6 +66,7 @@ export class JazzThoughtStore { private readonly registry: EventRegistry; private readonly runtimeRevision: string; private producerStorageRevision = 3; + private consumerStorageRevision = 5; constructor(private readonly options: JazzThoughtStoreOptions) { this.registry = options.registry ?? createDefaultRegistry(); @@ -196,7 +197,7 @@ export class JazzThoughtStore { await result.wait({ tier: this.durabilityTier }); return result.value; } catch (error) { - if (!isRecoverableProducerStorageError(error) || attempt === 7) throw error; + if (!isRecoverableStorageError(error) || attempt === 7) throw error; this.producerStorageRevision = storageRevision + 1; await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); } @@ -834,76 +835,89 @@ export class JazzThoughtStore { if (!outputSource || candidates.some((candidate) => candidate.source !== outputSource)) { throw new Error("Consumer settlement receipts must share one output source"); } - const settledAt = new Date().toISOString(); - await this.ensureTransactionReady(outputSource); - // Persistent/edge Jazz transactions do not inherit the readiness barrier of - // ordinary Db reads. Prime every pre-existing row the transaction may update - // or deduplicate, not only the output source. Otherwise a restarted consumer - // can miss its durable run/progress rows and leave a partially allocated - // lifecycle object when the transaction fails. - const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ - key: outputSource, - }).limit(1), { tier: this.durabilityTier }); - const [runSnapshot] = await this.db.all(thoughtstreamApp.runs.where({ key: run.id }).limit(1), { - tier: this.durabilityTier, - }); - if (!runSnapshot) throw new Error(`Cannot settle missing run ${run.id}`); - const eventSnapshots = new Map>(); - for (const candidate of candidates) { - const id = eventId(candidate); - const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ - key: id, - }).limit(1), { tier: this.durabilityTier }); - if (eventSnapshot) eventSnapshots.set(id, eventSnapshot); - } - let progressSnapshot: Record | undefined; - if (progress) { - [progressSnapshot] = await this.db.all(thoughtstreamApp.consumerProgress.where({ - key: progress.id, + for (let attempt = 0; attempt < 8; attempt += 1) { + const storageRevision = this.consumerStorageRevision; + const settledAt = new Date().toISOString(); + await this.ensureTransactionReady(outputSource); + // Persistent/edge Jazz transactions do not inherit the readiness barrier of + // ordinary Db reads. Prime every pre-existing row the transaction may update + // or deduplicate, not only the output source. Otherwise a restarted consumer + // can miss its durable run/progress rows and leave a partially allocated + // lifecycle object when the transaction fails. + const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ + key: outputSource, }).limit(1), { tier: this.durabilityTier }); - } - const result = await this.db.transaction(async (tx) => { - const sourceRow = sourceSnapshot; - if (sourceRow && String(sourceRow.kind) !== candidates[0]!.sourceKind) { - throw new Error(`Source kind changed for ${outputSource}`); - } - let sequence = Number(sourceRow?.lastSequence ?? 0); - const events: ThoughtEvent[] = []; + const [runSnapshot] = await this.db.all(thoughtstreamApp.runs.where({ key: run.id }).limit(1), { + tier: this.durabilityTier, + }); + if (!runSnapshot) throw new Error(`Cannot settle missing run ${run.id}`); + const eventSnapshots = new Map>(); for (const candidate of candidates) { const id = eventId(candidate); - const existing = eventSnapshots.get(id); - if (existing) { - const event = eventFromJazz(existing); - assertSameEventIdentity(event, candidate); - events.push(event); - continue; - } - sequence += 1; - const event = buildEvent(candidate, id, sequence, settledAt, this.runtimeRevision, this.registry); - // v2 and v3 may remain allocated but unlinked after a failed Jazz - // transaction. A failed insert also poisons the active batch, so do not - // probe old namespaces inside this transaction before selecting v4. - tx.insert(thoughtstreamApp.events, eventToJazz(event), { id: jazzRowId("consumer-event-v4", id) }); - events.push(event); + const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ + key: id, + }).limit(1), { tier: this.durabilityTier }); + if (eventSnapshot) eventSnapshots.set(id, eventSnapshot); } - const sourceData = { - key: outputSource, - kind: candidates[0]!.sourceKind, - enabled: true, - configJson: sourceRow ? String(sourceRow.configJson) : "{}", - lastSequence: sequence, - createdAt: sourceRow ? String(sourceRow.createdAt) : settledAt, - updatedAt: settledAt, - }; - if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); - else tx.insert(thoughtstreamApp.sources, sourceData, { id: jazzRowId("consumer-source-v3", outputSource) }); + let progressSnapshot: Record | undefined; + if (progress) { + [progressSnapshot] = await this.db.all(thoughtstreamApp.consumerProgress.where({ + key: progress.id, + }).limit(1), { tier: this.durabilityTier }); + } + try { + const result = await this.db.transaction(async (tx) => { + const sourceRow = sourceSnapshot; + if (sourceRow && String(sourceRow.kind) !== candidates[0]!.sourceKind) { + throw new Error(`Source kind changed for ${outputSource}`); + } + let sequence = Number(sourceRow?.lastSequence ?? 0); + const events: ThoughtEvent[] = []; + for (const candidate of candidates) { + const id = eventId(candidate); + const existing = eventSnapshots.get(id); + if (existing) { + const event = eventFromJazz(existing); + assertSameEventIdentity(event, candidate); + events.push(event); + continue; + } + sequence += 1; + const event = buildEvent(candidate, id, sequence, settledAt, this.runtimeRevision, this.registry); + tx.insert(thoughtstreamApp.events, eventToJazz(event), { + id: jazzRowId(`consumer-event-v${storageRevision}`, id), + }); + events.push(event); + } + const sourceData = { + key: outputSource, + kind: candidates[0]!.sourceKind, + enabled: true, + configJson: sourceRow ? String(sourceRow.configJson) : "{}", + lastSequence: sequence, + createdAt: sourceRow ? String(sourceRow.createdAt) : settledAt, + updatedAt: settledAt, + }; + if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); + else { + tx.insert(thoughtstreamApp.sources, sourceData, { + id: jazzRowId(`consumer-source-v${storageRevision}`, outputSource), + }); + } - tx.update(thoughtstreamApp.runs, String(runSnapshot.id), runToJazz(run)); - if (progress) await upsertProgressInTransaction(tx, progress, progressSnapshot); - return events; - }); - await result.wait({ tier: this.durabilityTier }); - return result.value; + tx.update(thoughtstreamApp.runs, String(runSnapshot.id), runToJazz(run)); + if (progress) await upsertProgressInTransaction(tx, progress, progressSnapshot, storageRevision); + return events; + }); + await result.wait({ tier: this.durabilityTier }); + return result.value; + } catch (error) { + if (!isRecoverableStorageError(error) || attempt === 7) throw error; + this.consumerStorageRevision = storageRevision + 1; + await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); + } + } + throw new Error("Consumer transaction exhausted storage recovery attempts"); } // Jazz transactions commit the account and reservation together but do not serialize @@ -1735,7 +1749,7 @@ async function upsertCursorInTransaction( } } -function isRecoverableProducerStorageError(error: unknown): boolean { +function isRecoverableStorageError(error: unknown): boolean { if (!(error instanceof Error)) return false; return error.message.includes("object already exists") || error.message.includes("database is locked"); @@ -1745,6 +1759,7 @@ async function upsertProgressInTransaction( tx: TransactionScope, progress: ConsumerProgress, row?: Record, + storageRevision = 3, ): Promise { const data = { key: progress.id, @@ -1767,6 +1782,8 @@ async function upsertProgressInTransaction( } tx.update(thoughtstreamApp.consumerProgress, String(row.id), data); } else { - tx.insert(thoughtstreamApp.consumerProgress, data, { id: jazzRowId("consumer-progress-v3", progress.id) }); + tx.insert(thoughtstreamApp.consumerProgress, data, { + id: jazzRowId(`consumer-progress-v${storageRevision}`, progress.id), + }); } } diff --git a/test/jazz-store.test.ts b/test/jazz-store.test.ts index a0ed1f9..bd630f9 100644 --- a/test/jazz-store.test.ts +++ b/test/jazz-store.test.ts @@ -1,7 +1,8 @@ import fs from "node:fs/promises"; import { afterEach, describe, expect, test } from "vitest"; -import type { EventCandidate } from "../src/events/types.js"; +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[] = []; @@ -146,6 +147,118 @@ describe("JazzThoughtStore", () => { }); }); + test("recovers consumer settlement after a failed transaction leaves allocated row ids", async () => { + const root = await temporaryProject(); + roots.push(root); + const store = 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, + }); + expect((store as unknown as { consumerStorageRevision: number }).consumerStorageRevision).toBeGreaterThan(5); + }); + test("delivers only newly added rows from a consumer subscription", async () => { const root = await temporaryProject(); roots.push(root);