From 7e695c355c2648a40c5b09421944a9257d8646a6 Mon Sep 17 00:00:00 2001 From: Timothy Laurent Date: Sun, 3 May 2026 21:39:08 -0700 Subject: [PATCH] fix(core): fail-fast on corruption, preserve buffer on transient errors - Throw from init() when integrity check fails and recovery doesn't help, preventing writes to a known-corrupt database - Only drain buffer on successful flush; requeue valid entries on transient failures (busy/IO) so they aren't silently lost - Add spy on flushAgentLogBuffer in deleteTask test to prove flush-before-delete Co-Authored-By: Claude Opus 4.6 --- packages/core/src/__tests__/store.test.ts | 5 ++++- packages/core/src/db.ts | 13 +++++++++++++ packages/core/src/store.ts | 13 ++++++++++--- 3 files changed, 27 insertions(+), 4 deletions(-) diff --git a/packages/core/src/__tests__/store.test.ts b/packages/core/src/__tests__/store.test.ts index cbbf9accd..257fa4740 100644 --- a/packages/core/src/__tests__/store.test.ts +++ b/packages/core/src/__tests__/store.test.ts @@ -4917,8 +4917,11 @@ Task with acceptance criteria const task = await createTestTask(); await store.appendAgentLog(task.id, "to be cascaded", "text"); - // deleteTask should flush first, then cascade-delete the entry + // Prove flush happens before delete + const flushSpy = vi.spyOn(store as any, "flushAgentLogBuffer"); await store.deleteTask(task.id); + expect(flushSpy).toHaveBeenCalled(); + flushSpy.mockRestore(); const after = (store as any).db.prepare( "SELECT COUNT(*) as count FROM agentLogEntries WHERE taskId = ?", diff --git a/packages/core/src/db.ts b/packages/core/src/db.ts index 909bea7f9..73d69f6d1 100644 --- a/packages/core/src/db.ts +++ b/packages/core/src/db.ts @@ -924,17 +924,30 @@ export class Database { this.corruptionDetected = false; console.warn(`[fusion:db] Database recovered via WAL checkpoint: ${this.dbPath}`); } else { + const recheckMsg = ("errors" in recheck && Array.isArray(recheck.errors)) + ? recheck.errors.slice(0, 3).join(" | ") + : "unknown"; console.error( `[fusion:db] Database is corrupted and could not be auto-recovered. ` + `Run: sqlite3 ${this.dbPath} ".recover" | sqlite3 ${this.dbPath}.recovered`, ); + throw new Error( + `[fusion:db] Refusing to initialize corrupted database at ${this.dbPath}. Integrity errors: ${recheckMsg}`, + ); } } catch (err) { + // Re-throw our own abort error; wrap others + if (err instanceof Error && err.message.startsWith("[fusion:db] Refusing")) { + throw err; + } const errMsg = err instanceof Error ? err.message : String(err); console.error( `[fusion:db] Database corruption detected for ${this.dbPath} and checkpoint recovery failed: ${errMsg}. ` + "Manual recovery required.", ); + throw new Error( + `[fusion:db] Refusing to initialize corrupted database at ${this.dbPath}. Recovery error: ${errMsg}`, + ); } } diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 7e82fa586..b23646e06 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -4741,13 +4741,15 @@ export class TaskStore extends EventEmitter { const batch = this.agentLogBuffer.slice(); const flushCount = batch.length; + let validEntries = batch; + let flushSucceeded = false; try { // Filter out entries for deleted tasks to prevent FK violations // from poisoning the entire buffer. const liveTaskIds = new Set( (this.db.prepare("SELECT id FROM tasks").all() as Array<{ id: string }>).map((r) => r.id), ); - const validEntries = batch.filter((e) => liveTaskIds.has(e.taskId)); + validEntries = batch.filter((e) => liveTaskIds.has(e.taskId)); const dropped = batch.length - validEntries.length; if (dropped > 0) { console.warn( @@ -4767,10 +4769,15 @@ export class TaskStore extends EventEmitter { this.db.bumpLastModified(); }); } + flushSucceeded = true; } finally { - // Always drain the flushed slice so a failed transaction doesn't - // cause the same entries to block every future flush. + // Always drain the original slice from the buffer. this.agentLogBuffer.splice(0, flushCount); + // On transient failures (busy/IO), requeue valid entries for retry. + // Stale rows were already filtered out above. + if (!flushSucceeded && validEntries.length > 0) { + this.agentLogBuffer.unshift(...validEntries); + } } }