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); + } } }