diff --git a/spec/incidents.md b/spec/incidents.md index 9154c5f..d482ad0 100644 --- a/spec/incidents.md +++ b/spec/incidents.md @@ -32,7 +32,7 @@ The first slice normalizes: - exhausted `ConsumerScheduler` operations after their configured retries; - `stream.thought.action.telegram.send.failed` as a delivery incident. -Projection is deterministic and idempotent. It does not alter the originating event, retry count, consumer progress, model invocation, or delivery claim. A failed incident projection is an observability failure, not authority to replay an external action. +Projection is deterministic and idempotent. The first valid append-only incident row for an origin event is authoritative on later projection passes; current consumer progress cannot rewrite its `retryable` or `progress` classification. If an earlier failed attempt is first projected only after a later attempt under the same execution key exists, the earlier failure is classified as progress-unchanged and retryable even when the later attempt has since advanced source progress. It does not alter the originating event, retry count, consumer progress, model invocation, or delivery claim. A failed incident projection is an observability failure, not authority to replay an external action. ## Private ledger diff --git a/src/incidents/projector.ts b/src/incidents/projector.ts index e0bdcf2..1a04444 100644 --- a/src/incidents/projector.ts +++ b/src/incidents/projector.ts @@ -74,22 +74,41 @@ export interface SchedulerExhaustionEvidence { export class OperationalIncidentProjector { async project(store: JazzThoughtStore): Promise { - const [events, progress] = await Promise.all([ + const [events, progress, runs, existingIncidents] = await Promise.all([ store.listEvents({ types: PROJECTED_EVENT_TYPES }), store.listConsumerProgress(), + store.listRuns(), + listOperationalIncidents(store), ]); + const existingByOrigin = new Map(existingIncidents.flatMap(({ event, incident }) => { + const originEventId = stringField(incident.references.originEventId); + return originEventId ? [[originEventId, { event, incident }] as const] : []; + })); let offered = 0; let inserted = 0; let unchanged = 0; const incidents: ThoughtEvent[] = []; + const pending: Array<{ incident: OperationalIncident; origin: ThoughtEvent }> = []; for (const event of events) { - const incident = await incidentForEvent(store, event, progress); + const existing = existingByOrigin.get(event.id); + if (existing) { + offered += 1; + unchanged += 1; + incidents.push(existing.event); + continue; + } + const incident = await incidentForEvent(store, event, progress, runs); if (!incident) continue; offered += 1; - const result = await store.appendEvent(incidentCandidate(incident, event)); - if (result.inserted) inserted += 1; - else unchanged += 1; - incidents.push(result.event); + pending.push({ incident, origin: event }); + } + if (pending.length > 0) { + const result = await store.appendProducerBatch(pending.map(({ incident, origin }) => ( + incidentCandidate(incident, origin) + ))); + inserted += result.inserted.length; + unchanged += result.unchanged.length; + incidents.push(...result.events); } return { offered, inserted, unchanged, incidents }; } @@ -145,6 +164,7 @@ async function incidentForEvent( store: JazzThoughtStore, event: ThoughtEvent, progress: ConsumerProgress[], + runs: AgentRun[], ): Promise { if (event.type === "stream.thought.connector.failed") { return connectorIncident(event, "connector-failure", "warning", "open", "connector-failed"); @@ -176,7 +196,7 @@ async function incidentForEvent( }); } if (event.type.startsWith("stream.thought.agent.run.")) { - return agentRunIncident(store, event, progress); + return agentRunIncident(store, event, progress, runs); } return undefined; } @@ -208,10 +228,11 @@ async function agentRunIncident( store: JazzThoughtStore, event: ThoughtEvent, progress: ConsumerProgress[], + runs: AgentRun[], ): Promise { const runId = stringField(event.payload.runId); if (!runId) return undefined; - const run = await store.getRun(runId); + const run = runs.find((candidate) => candidate.id === runId); if (!run) return undefined; const trigger = await store.getEvent(run.triggerEventId); if (!trigger) return undefined; @@ -225,7 +246,14 @@ async function agentRunIncident( if (status !== "failed" && status !== "blocked" && status !== "abandoned") return undefined; const category = `agent-run-${status}` as Extract; - const progressAdvanced = status === "blocked" || advanced; + // A later attempt can exist only when this terminal attempt left the trigger + // unconsumed. Current progress may already include that later attempt, so it + // cannot be used alone when rebuilding the earlier incident. + const retried = runs.some((candidate) => ( + candidate.executionKey === run.executionKey + && candidate.attempt > run.attempt + )); + const progressAdvanced = status === "blocked" || (!retried && advanced); const diagnostic = objectField(run.result?.failureDiagnostic); const code = agentCode(diagnostic?.code) ?? `agent-run-${status}`; const stage = agentStage(diagnostic?.stage); diff --git a/src/jazz/store.ts b/src/jazz/store.ts index 6290dc7..fddd15c 100644 --- a/src/jazz/store.ts +++ b/src/jazz/store.ts @@ -63,6 +63,7 @@ export class JazzThoughtStore { private readonly durabilityTier: "local" | "edge"; private readonly registry: EventRegistry; private readonly runtimeRevision: string; + private producerStorageRevision = 3; constructor(private readonly options: JazzThoughtStoreOptions) { this.registry = options.registry ?? createDefaultRegistry(); @@ -118,71 +119,87 @@ export class JazzThoughtStore { } if (cursor && cursor.source !== source) throw new Error("Producer cursor source does not match event source"); const sourceKind = candidates[0]?.sourceKind ?? "system"; - const observedAt = new Date().toISOString(); - await this.ensureTransactionReady(source); - const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ key: source }).limit(1), { - tier: this.durabilityTier, - }); - const eventSnapshots = new Map>(); - for (const candidate of candidates) { - const id = eventId(candidate); - if (eventSnapshots.has(id)) continue; - const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ key: id }).limit(1), { + for (let attempt = 0; attempt < 8; attempt += 1) { + const storageRevision = this.producerStorageRevision; + const observedAt = new Date().toISOString(); + await this.ensureTransactionReady(source); + const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ key: source }).limit(1), { tier: this.durabilityTier, }); - if (eventSnapshot) eventSnapshots.set(id, eventSnapshot); - } - const result = await this.db.transaction(async (tx) => { - const sourceRow = sourceSnapshot; - if (sourceRow && candidates.length > 0 && String(sourceRow.kind) !== sourceKind) { - throw new Error(`Source kind changed for ${source}: ${String(sourceRow.kind)} -> ${sourceKind}`); - } - let sequence = Number(sourceRow?.lastSequence ?? 0); - const events: ThoughtEvent[] = []; - const inserted: ThoughtEvent[] = []; - const unchanged: ThoughtEvent[] = []; - const seen = new Map(); + const eventSnapshots = new Map>(); for (const candidate of candidates) { const id = eventId(candidate); - const repeated = seen.get(id); - if (repeated) { - assertSameEventIdentity(repeated, candidate); - events.push(repeated); - unchanged.push(repeated); - continue; - } - const existing = eventSnapshots.get(id); - if (existing) { - const event = eventFromJazz(existing); - assertSameEventIdentity(event, candidate); - events.push(event); - unchanged.push(event); - seen.set(id, event); - continue; - } - sequence += 1; - const event = buildEvent(candidate, id, sequence, observedAt, this.runtimeRevision, this.registry); - tx.insert(thoughtstreamApp.events, eventToJazz(event), { id: jazzRowId("producer-event-v2", id) }); - events.push(event); - inserted.push(event); - seen.set(id, event); + if (eventSnapshots.has(id)) continue; + 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: source, - kind: sourceRow ? String(sourceRow.kind) : sourceKind, - enabled: true, - configJson: sourceRow ? String(sourceRow.configJson) : "{}", - lastSequence: sequence, - createdAt: sourceRow ? String(sourceRow.createdAt) : observedAt, - updatedAt: observedAt, - }; - if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); - else tx.insert(thoughtstreamApp.sources, sourceData, { id: jazzRowId("producer-source-v2", source) }); - if (cursor) await upsertCursorInTransaction(tx, cursor); - return { events, inserted, unchanged }; - }); - await result.wait({ tier: this.durabilityTier }); - return result.value; + try { + const result = await this.db.transaction(async (tx) => { + const sourceRow = sourceSnapshot; + if (sourceRow && candidates.length > 0 && String(sourceRow.kind) !== sourceKind) { + throw new Error(`Source kind changed for ${source}: ${String(sourceRow.kind)} -> ${sourceKind}`); + } + let sequence = Number(sourceRow?.lastSequence ?? 0); + const events: ThoughtEvent[] = []; + const inserted: ThoughtEvent[] = []; + const unchanged: ThoughtEvent[] = []; + const seen = new Map(); + for (const candidate of candidates) { + const id = eventId(candidate); + const repeated = seen.get(id); + if (repeated) { + assertSameEventIdentity(repeated, candidate); + events.push(repeated); + unchanged.push(repeated); + continue; + } + const existing = eventSnapshots.get(id); + if (existing) { + const event = eventFromJazz(existing); + assertSameEventIdentity(event, candidate); + events.push(event); + unchanged.push(event); + seen.set(id, event); + continue; + } + sequence += 1; + const event = buildEvent(candidate, id, sequence, observedAt, this.runtimeRevision, this.registry); + tx.insert(thoughtstreamApp.events, eventToJazz(event), { + id: jazzRowId(`producer-event-v${storageRevision}`, id), + }); + events.push(event); + inserted.push(event); + seen.set(id, event); + } + const sourceData = { + key: source, + kind: sourceRow ? String(sourceRow.kind) : sourceKind, + enabled: true, + configJson: sourceRow ? String(sourceRow.configJson) : "{}", + lastSequence: sequence, + createdAt: sourceRow ? String(sourceRow.createdAt) : observedAt, + updatedAt: observedAt, + }; + if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); + else { + tx.insert(thoughtstreamApp.sources, sourceData, { + id: jazzRowId(`producer-source-v${storageRevision}`, source), + }); + } + if (cursor) await upsertCursorInTransaction(tx, cursor, storageRevision); + return { events, inserted, unchanged }; + }); + await result.wait({ tier: this.durabilityTier }); + return result.value; + } catch (error) { + if (!isRecoverableProducerStorageError(error) || attempt === 7) throw error; + this.producerStorageRevision = storageRevision + 1; + await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); + } + } + throw new Error("Producer transaction exhausted storage recovery attempts"); } private async withProducerSourceLock(source: string, operation: () => Promise): Promise { @@ -1540,11 +1557,25 @@ function cursorToJazz(cursor: SourceCursor): Record { }; } -async function upsertCursorInTransaction(tx: TransactionScope, cursor: SourceCursor): Promise { - const row = await tx.one(thoughtstreamApp.sourceCursors.where({ id: jazzRowId("cursor", cursor.id) }).limit(1)); +async function upsertCursorInTransaction( + tx: TransactionScope, + cursor: SourceCursor, + storageRevision = 2, +): Promise { + const row = await tx.one(thoughtstreamApp.sourceCursors.where({ key: cursor.id }).limit(1)); const data = cursorToJazz(cursor); if (row) tx.update(thoughtstreamApp.sourceCursors, String(row.id), data); - else tx.insert(thoughtstreamApp.sourceCursors, data, { id: jazzRowId("cursor", cursor.id) }); + else { + tx.insert(thoughtstreamApp.sourceCursors, data, { + id: jazzRowId(`cursor-v${storageRevision}`, cursor.id), + }); + } +} + +function isRecoverableProducerStorageError(error: unknown): boolean { + if (!(error instanceof Error)) return false; + return error.message.includes("object already exists") + || error.message.includes("database is locked"); } async function upsertProgressInTransaction( diff --git a/test/incidents.test.ts b/test/incidents.test.ts index 22d697a..1953b4b 100644 --- a/test/incidents.test.ts +++ b/test/incidents.test.ts @@ -113,6 +113,58 @@ describe("operational incidents", () => { await expect(IncidentLedger.open(project, "../escape.jsonl")).rejects.toThrow("escapes the runtime root"); }); + test("keeps the first projected progress classification authoritative after progress advances", async () => { + const { store } = await fixtureStore(); + const trigger = await appendTrigger(store, "projection-stability", SECRET); + await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); + + const projector = new OperationalIncidentProjector(); + expect(await projector.project(store)).toMatchObject({ offered: 1, inserted: 1, unchanged: 0 }); + const first = await listOperationalIncidents(store); + expect(first[0]?.incident).toMatchObject({ retryable: true, progress: "unchanged" }); + + await store.initializeConsumerProgress({ + id: "consumer-progress:resident-fixture:1:jetstream:cameron-bluesky", + consumerId: "resident-fixture", + consumerVersion: 1, + source: trigger.event.source, + lastSequence: trigger.event.sourceSequence, + lastEventId: trigger.event.id, + updatedAt: "2026-07-22T01:00:10.000Z", + }); + + expect(await projector.project(store)).toMatchObject({ offered: 1, inserted: 0, unchanged: 1 }); + const replayed = await listOperationalIncidents(store); + expect(replayed).toHaveLength(1); + expect(replayed[0]?.incident).toEqual(first[0]?.incident); + }); + + test("classifies an unprojected failed attempt before its later successful retry", async () => { + const { store } = await fixtureStore(); + const trigger = await appendTrigger(store, "later-attempt", SECRET); + await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); + await appendCompletedRun(store, trigger.event, 2); + await store.initializeConsumerProgress({ + id: "consumer-progress:resident-fixture:1:jetstream:cameron-bluesky", + consumerId: "resident-fixture", + consumerVersion: 1, + source: trigger.event.source, + lastSequence: trigger.event.sourceSequence, + lastEventId: trigger.event.id, + updatedAt: "2026-07-22T01:00:10.000Z", + }); + + expect(await new OperationalIncidentProjector().project(store)) + .toMatchObject({ offered: 1, inserted: 1, unchanged: 0 }); + const incidents = await listOperationalIncidents(store); + expect(incidents).toHaveLength(1); + expect(incidents[0]?.incident).toMatchObject({ + attempt: 1, + retryable: true, + progress: "unchanged", + }); + }); + test("alerts on a Jetstream-rooted failure without treating successful Bluesky work as notification input", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); @@ -383,10 +435,10 @@ async function appendTerminalRun( return run; } -async function appendCompletedRun(store: JazzThoughtStore, trigger: ThoughtEvent): Promise { +async function appendCompletedRun(store: JazzThoughtStore, trigger: ThoughtEvent, attempt = 1): Promise { const at = new Date(Date.parse(trigger.occurredAt) + 1_000).toISOString(); await store.upsertRun({ - id: `run_${trigger.externalId}_completed`, + id: `run_${trigger.externalId}_completed_${attempt}`, executionKey: `execution_${trigger.externalId}`, triggerEventId: trigger.id, agentId: "resident-fixture", @@ -394,7 +446,7 @@ async function appendCompletedRun(store: JazzThoughtStore, trigger: ThoughtEvent status: "completed", inputEventIds: [trigger.id], outputEventIds: [], - attempt: 1, + attempt, provider: "letta-cloud", model: "fixture-model", promptHash: "fixture-prompt",