From cf32b767aa9d7027c3b96e3487706f62c410dc66 Mon Sep 17 00:00:00 2001 From: QAyong Date: Thu, 8 Oct 2026 00:44:34 +0800 Subject: [PATCH] =?UTF-8?q?perf(buddy):=20=E5=90=AF=E5=8A=A8=E6=97=B6?= =?UTF-8?q?=E8=B7=B3=E8=BF=87=E7=A8=B3=E5=AE=9A=E5=8E=86=E5=8F=B2=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E7=9A=84=E6=81=A2=E5=A4=8D=E9=87=8D=E5=BB=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../__tests__/buddyServiceEnvironment.spec.ts | 8 + .../main/runtime/buddyServiceEnvironment.ts | 4 + .../service/src/events/RunEventFailure.ts | 12 +- apps/buddy/service/src/events/RunEventLog.ts | 149 +++++- .../buddy/service/src/events/RunEventPorts.ts | 2 + .../service/src/events/RunEventProjector.ts | 15 +- .../service/src/events/RunEventQueries.ts | 37 ++ .../service/src/events/RunEventRecovery.ts | 49 ++ .../buddy/service/src/events/RunEventStore.ts | 31 +- .../src/events/__tests__/RunEventLog.spec.ts | 3 +- .../events/__tests__/RunEventRecovery.spec.ts | 445 ++++++++++++++++++ .../service/src/events/createRunEventLog.ts | 1 + apps/buddy/service/src/index.ts | 12 +- .../src/storage/__tests__/schema.spec.ts | 35 ++ .../usageAnalyticsRepository.spec.ts | 3 + apps/buddy/service/src/storage/database.ts | 5 + .../migrations/v23RunEventCheckpoints.ts | 38 ++ apps/buddy/service/src/storage/schema.ts | 4 +- 18 files changed, 827 insertions(+), 26 deletions(-) create mode 100644 apps/buddy/service/src/events/RunEventRecovery.ts create mode 100644 apps/buddy/service/src/events/__tests__/RunEventRecovery.spec.ts create mode 100644 apps/buddy/service/src/storage/migrations/v23RunEventCheckpoints.ts diff --git a/apps/buddy/electron/main/runtime/__tests__/buddyServiceEnvironment.spec.ts b/apps/buddy/electron/main/runtime/__tests__/buddyServiceEnvironment.spec.ts index b292c889..175f7f79 100644 --- a/apps/buddy/electron/main/runtime/__tests__/buddyServiceEnvironment.spec.ts +++ b/apps/buddy/electron/main/runtime/__tests__/buddyServiceEnvironment.spec.ts @@ -14,6 +14,14 @@ import { const executeFile = promisify(execFile) describe('buddyServiceEnvironment', () => { + it.each(['linux', 'win32'] as const)('passes only the explicit full-recovery diagnostic switch on %s', (platform) => { + const name = platform === 'win32' ? 'lexora_buddy_event_recovery' : 'LEXORA_BUDDY_EVENT_RECOVERY' + const source = { [name]: 'full', UNRELATED_SECRET: 'secret' } + const buddyHome = platform === 'win32' ? 'C:\\fixture\\buddy' : '/fixture/buddy' + expect(createBuddyServiceEnvironment(source, buddyHome, platform).LEXORA_BUDDY_EVENT_RECOVERY).toBe('full') + expect(createBuddyServiceEnvironment({ [name]: 'arbitrary' }, buddyHome, platform).LEXORA_BUDDY_EVENT_RECOVERY).toBeUndefined() + expect(createBuddyServiceEnvironment(source, buddyHome, platform).UNRELATED_SECRET).toBeUndefined() + }) it('replaces all inherited proxy and bypass settings with the application gateway', () => { const gateway = 'http://lexora:fixture@127.0.0.1:3128' const env = createBuddyServiceEnvironment({ http_proxy: 'http://old.invalid:1', HTTPS_PROXY: 'http://old.invalid:2', NO_PROXY: '*' }, '/tmp/buddy', 'linux', gateway) diff --git a/apps/buddy/electron/main/runtime/buddyServiceEnvironment.ts b/apps/buddy/electron/main/runtime/buddyServiceEnvironment.ts index afa6a81e..e127cffc 100644 --- a/apps/buddy/electron/main/runtime/buddyServiceEnvironment.ts +++ b/apps/buddy/electron/main/runtime/buddyServiceEnvironment.ts @@ -35,11 +35,15 @@ export function createBuddyServiceEnvironment( const environmentSource = !isWindows(targetPlatform.id) && !source.HOME ? { ...source, HOME: homedir() } : source + const recoveryMode = Object.entries(environmentSource).find(([name]) => ( + (targetPlatform.environment.caseSensitive ? name : name.toUpperCase()) === 'LEXORA_BUDDY_EVENT_RECOVERY' + ))?.[1] const environment = createChildProcessEnvironment({ source: environmentSource, platform: targetPlatform, additions: { LEXORA_BUDDY_HOME: buddyHome, + LEXORA_BUDDY_EVENT_RECOVERY: recoveryMode === 'full' ? 'full' : undefined, NODE_USE_ENV_PROXY: '1', PI_CODING_AGENT_DIR: filePathAdapters[targetPlatform.id].resolveInput('agent', buddyHome), }, diff --git a/apps/buddy/service/src/events/RunEventFailure.ts b/apps/buddy/service/src/events/RunEventFailure.ts index fe30ef16..fc33e7ea 100644 --- a/apps/buddy/service/src/events/RunEventFailure.ts +++ b/apps/buddy/service/src/events/RunEventFailure.ts @@ -35,6 +35,16 @@ export class RunEventCorruptionError extends RunEventLogFatalError { } } +export class RunEventCheckpointError extends RunEventLogFatalError { + readonly code = 'EVENT_PROJECTION_FAILED' + readonly commitState = 'not_applicable' + + constructor(runId: string, options?: ErrorOptions) { + super('Lexora Buddy run event checkpoint invalidation failed', runEventFailureScope(runId), options) + this.name = 'RunEventCheckpointError' + } +} + export class RunEventProjectionError extends RunEventLogFatalError { readonly code = 'EVENT_PROJECTION_FAILED' readonly commitState = 'committed' @@ -58,7 +68,7 @@ export class RunEventStorageError extends RunEventLogFatalError { readonly code = 'EVENT_STORAGE_FAILED' readonly commitState: 'committed' | 'not_applicable' | 'unknown' readonly operation: 'append' | 'compact' | 'read' | 'repair' | 'scan' - readonly stage: 'close' | 'directory' | 'mkdir' | 'open' | 'read' | 'readdir' | 'rename' | 'sync' | 'truncate' | 'unlink' | 'write' + readonly stage: 'close' | 'directory' | 'mkdir' | 'open' | 'read' | 'readdir' | 'rename' | 'stat' | 'sync' | 'truncate' | 'unlink' | 'write' constructor( scope: RunEventFailureScope, diff --git a/apps/buddy/service/src/events/RunEventLog.ts b/apps/buddy/service/src/events/RunEventLog.ts index 7f524f11..81804335 100644 --- a/apps/buddy/service/src/events/RunEventLog.ts +++ b/apps/buddy/service/src/events/RunEventLog.ts @@ -7,7 +7,9 @@ import type { import type { RunEventLogPort } from './RunEventPorts' import type { RunEventProjector } from './RunEventProjector' import type { RunEventQueries } from './RunEventQueries' +import type { RunEventCheckpoint, RunEventRecoverySummary } from './RunEventRecovery' import type { RunEventStore } from './RunEventStore' +import { performance } from 'node:perf_hooks' import { Emitter } from '../../../shared/events/Emitter' import { copyEventSnapshot } from '../../../shared/events/eventSnapshot' import { @@ -17,24 +19,29 @@ import { } from './BuddyRunEvent' import { createRunEventCompactionPlan } from './RunEventCompaction' import { + RunEventCheckpointError, RunEventLogFatalError, RunEventProjectionError, } from './RunEventFailure' +import { canReuseRunEventCheckpoint, RUN_EVENT_PROJECTION_VERSION } from './RunEventRecovery' export interface RunEventLogCallbacks { onObserverError?: (error: unknown) => void onFatalFailure?: (error: RunEventLogFatalError) => void + onRecovery?: (summary: Readonly) => void } export type RunEventProjectorPort = Pick< RunEventProjector, - 'project' | 'rebuild' | 'removeEventRows' | 'validateNewFacts' + 'invalidateCheckpoint' | 'project' | 'rebuild' | 'removeEventRows' | 'validateNewFacts' > export type RunEventQueryPort = Pick< RunEventQueries, + | 'hasCheckpoint' | 'isTerminalRun' | 'list' + | 'listRecoveryRuns' | 'listCompactableTerminalRunIds' | 'listForConversation' | 'listForRuns' @@ -43,7 +50,7 @@ export type RunEventQueryPort = Pick< export type RunEventStorePort = Pick< RunEventStore, - 'append' | 'listPersistedRunIds' | 'readAndRepair' | 'replace' + 'append' | 'fingerprint' | 'listPersistedRunIds' | 'readAndRepair' | 'replace' > export interface RunEventLogOptions extends RunEventLogCallbacks { @@ -61,9 +68,11 @@ export class RunEventLogClosedError extends Error { export class RunEventLog implements RunEventLogPort { readonly #nextSequences = new Map() + readonly #compactionCheckedRuns = new Map() readonly #committed: Emitter> readonly onDidCommit: RunEventLogPort['onDidCommit'] readonly #onFatalFailure?: (error: RunEventLogFatalError) => void + readonly #onRecovery?: RunEventLogCallbacks['onRecovery'] readonly #projector: RunEventProjectorPort readonly #queries: RunEventQueryPort readonly #store: RunEventStorePort @@ -76,6 +85,7 @@ export class RunEventLog implements RunEventLogPort { this.#committed = new Emitter(options.onObserverError ?? (() => {})) this.onDidCommit = this.#committed.event this.#onFatalFailure = options.onFatalFailure + this.#onRecovery = options.onRecovery this.#projector = options.projector this.#queries = options.queries this.#store = options.store @@ -116,6 +126,7 @@ export class RunEventLog implements RunEventLogPort { })))) const committed = events.map(event => copyEventSnapshot(event)) this.#projector.validateNewFacts(events) + this.#invalidateCheckpoint(runId) await this.#runStoreOperation(() => this.#store.append(events)) this.#nextSequences.set(runId, nextSequence + events.length) await this.#projectCommitted(events) @@ -147,11 +158,65 @@ export class RunEventLog implements RunEventLogPort { } replay(runId: string): Promise { - return this.#enqueue(runId, async () => { - const events = await this.#readAndRepair(runId) - this.#nextSequences.set(runId, nextEventSequence(events)) - return this.#rebuildProjection(runId, events) - }) + return this.#enqueue(runId, () => this.#replay(runId)) + } + + async recoverAll(options: { force?: boolean } = {}): Promise> { + this.#assertOpen() + const started = performance.now() + const summary: RunEventRecoverySummary = { + durationMs: 0, + enumerationMs: 0, + inspectionMs: 0, + readMs: 0, + projectionMs: 0, + scanned: 0, + skipped: 0, + rebuilt: 0, + events: 0, + failed: 0, + } + try { + const enumerationStarted = performance.now() + const runs = this.#queries.listRecoveryRuns() + const conversations = new Map(runs.map(run => [run.runId, run.conversationId])) + const persisted = new Set(await this.#runStoreOperation(() => ( + this.#store.listPersistedRunIds(runs.map(run => run.runId), conversations) + ))) + summary.enumerationMs = performance.now() - enumerationStarted + for (const run of runs) { + // Missing logs were not replayed by the old enumeration; never leave their checkpoints trusted. + if (!persisted.has(run.runId)) { + await this.#enqueue(run.runId, async () => this.#invalidateCheckpoint(run.runId)) + continue + } + await this.#enqueue(run.runId, async () => { + summary.scanned++ + const fingerprint = await this.#fingerprint(run.runId, run.conversationId, summary) + if (!options.force && !this.#fatalFailure && canReuseRunEventCheckpoint(run, fingerprint) + && this.#queries.hasCheckpoint(run.runId, run.checkpoint)) { + this.#nextSequences.set(run.runId, run.checkpoint.lastSequence + 1) + this.#compactionCheckedRuns.set(run.runId, run.checkpoint) + summary.skipped++ + return + } + summary.events += await this.#replay(run.runId, run.conversationId, summary, fingerprint) + summary.rebuilt++ + }) + } + return Object.freeze({ ...summary, durationMs: performance.now() - started }) + } + catch (error) { + summary.failed++ + throw error + } + finally { + summary.durationMs = performance.now() - started + try { + this.#onRecovery?.(Object.freeze({ ...summary })) + } + catch {} + } } compactTerminalRun(runId: string): Promise { @@ -161,8 +226,16 @@ export class RunEventLog implements RunEventLogPort { async compactTerminalRuns(): Promise { this.#assertMutationAllowed() let removed = 0 - for (const runId of this.#queries.listCompactableTerminalRunIds()) - removed += await this.compactTerminalRun(runId) + for (const runId of this.#queries.listCompactableTerminalRunIds()) { + removed += await this.#enqueueMutation(runId, async () => { + const checked = this.#compactionCheckedRuns.get(runId) + if (checked && this.#queries.hasCheckpoint(runId, checked) + && await this.#fingerprint(runId) === checked.fileFingerprint) { + return 0 + } + return this.#compactTerminalRun(runId) + }) + } return removed } @@ -200,6 +273,42 @@ export class RunEventLog implements RunEventLogPort { throw this.#fatalFailure } + async #replay( + runId: string, + conversationId?: string, + summary?: RunEventRecoverySummary, + knownFingerprint?: string | null, + ): Promise { + const before = knownFingerprint === undefined ? await this.#fingerprint(runId, conversationId, summary) : knownFingerprint + const readStarted = performance.now() + const events = await this.#readAndRepair(runId, conversationId) + const compacted = createRunEventCompactionPlan(events).removed.length === 0 + if (summary) + summary.readMs += performance.now() - readStarted + const after = await this.#fingerprint(runId, conversationId, summary) + this.#nextSequences.set(runId, nextEventSequence(events)) + // Repairs or concurrent external changes are not certified by this read. + const checkpoint = after !== null && before === after && compacted + ? { fileFingerprint: after, lastSequence: nextEventSequence(events) - 1, projectionVersion: RUN_EVENT_PROJECTION_VERSION } + : undefined + const projectionStarted = performance.now() + const count = this.#rebuildProjection(runId, events, checkpoint) + if (summary) + summary.projectionMs += performance.now() - projectionStarted + return count + } + + async #fingerprint(runId: string, conversationId?: string, summary?: RunEventRecoverySummary): Promise { + const started = performance.now() + try { + return await this.#runStoreOperation(() => this.#store.fingerprint(runId, conversationId)) + } + finally { + if (summary) + summary.inspectionMs += performance.now() - started + } + } + async #compactTerminalRun(runId: string): Promise { if (!this.#isTerminalRun(runId)) return 0 @@ -241,9 +350,12 @@ export class RunEventLog implements RunEventLogPort { } } - #rebuildProjection(runId: string, events: readonly BuddyRunEvent[]): number { + #rebuildProjection(runId: string, events: readonly BuddyRunEvent[], checkpoint?: RunEventCheckpoint): number { try { - return this.#projector.rebuild(runId, events) + const count = this.#projector.rebuild(runId, events, checkpoint) + if (checkpoint && this.#queries.hasCheckpoint(runId, checkpoint)) + this.#compactionCheckedRuns.set(runId, checkpoint) + return count } catch (cause) { this.#fail(new RunEventProjectionError(runId, events, { cause })) @@ -261,8 +373,19 @@ export class RunEventLog implements RunEventLogPort { throw this.#fatalFailure } - async #readAndRepair(runId: string): Promise { - return this.#runStoreOperation(() => this.#store.readAndRepair(runId)) + #invalidateCheckpoint(runId: string): void { + this.#compactionCheckedRuns.delete(runId) + try { + this.#projector.invalidateCheckpoint(runId) + } + catch (cause) { + this.#fail(new RunEventCheckpointError(runId, { cause })) + } + } + + async #readAndRepair(runId: string, conversationId?: string): Promise { + this.#invalidateCheckpoint(runId) + return this.#runStoreOperation(() => this.#store.readAndRepair(runId, conversationId)) } async #runStoreOperation(operation: () => Promise): Promise { diff --git a/apps/buddy/service/src/events/RunEventPorts.ts b/apps/buddy/service/src/events/RunEventPorts.ts index 20e4e7cd..c980a0d5 100644 --- a/apps/buddy/service/src/events/RunEventPorts.ts +++ b/apps/buddy/service/src/events/RunEventPorts.ts @@ -5,6 +5,7 @@ import type { BuddyRunEvent, ListBuddyRunEventsOptions, } from './BuddyRunEvent' +import type { RunEventRecoverySummary } from './RunEventRecovery' export interface RunEventWriter { append: (input: AppendBuddyRunEventInput) => Promise @@ -28,6 +29,7 @@ export interface RunEventMaintenance { close: () => Promise compactTerminalRun: (runId: string) => Promise compactTerminalRuns: () => Promise + recoverAll: (options?: { force?: boolean }) => Promise> replay: (runId: string) => Promise replayAll: () => Promise } diff --git a/apps/buddy/service/src/events/RunEventProjector.ts b/apps/buddy/service/src/events/RunEventProjector.ts index b09185a9..2d7bbbcc 100644 --- a/apps/buddy/service/src/events/RunEventProjector.ts +++ b/apps/buddy/service/src/events/RunEventProjector.ts @@ -1,5 +1,6 @@ import type { DatabaseSync } from 'node:sqlite' import type { BuddyRunEvent } from './BuddyRunEvent' +import type { RunEventCheckpoint } from './RunEventRecovery' import { z } from 'zod' import { BUDDY_ATTACHMENT_COUNT_LIMIT } from '../../../shared/conversation/attachmentPolicy' import { MAX_BUDDY_MESSAGE_TEXT_LENGTH } from '../../../shared/conversation/buddyMessageContent' @@ -111,8 +112,13 @@ export class RunEventProjector { )) } - rebuild(runId: string, events: readonly BuddyRunEvent[]): number { + invalidateCheckpoint(runId: string): void { + this.#database.prepare('DELETE FROM run_event_checkpoints WHERE run_id = ?').run(runId) + } + + rebuild(runId: string, events: readonly BuddyRunEvent[], checkpoint?: RunEventCheckpoint): number { return withTransaction(this.#database, () => { + this.invalidateCheckpoint(runId) const messageIds = new Set() for (const event of events) { if (event.runId !== runId) @@ -139,6 +145,13 @@ export class RunEventProjector { if (!messageIds.has(message.id)) remove.run(message.id, runId) } + if (checkpoint) { + this.#database.prepare(` + INSERT INTO run_event_checkpoints (run_id, last_sequence, projection_version, file_fingerprint) + SELECT ?, ?, ?, ? FROM runs + WHERE id = ? AND status IN ('completed', 'failed', 'cancelled') + `).run(runId, checkpoint.lastSequence, checkpoint.projectionVersion, checkpoint.fileFingerprint, runId) + } return count }) } diff --git a/apps/buddy/service/src/events/RunEventQueries.ts b/apps/buddy/service/src/events/RunEventQueries.ts index 62c00522..6e8c1e3e 100644 --- a/apps/buddy/service/src/events/RunEventQueries.ts +++ b/apps/buddy/service/src/events/RunEventQueries.ts @@ -1,5 +1,6 @@ import type { DatabaseSync } from 'node:sqlite' import type { BuddyRunEvent, ListBuddyRunEventsOptions } from './BuddyRunEvent' +import type { RunEventCheckpoint, RunEventRecoveryRun } from './RunEventRecovery' import { buddyRunEventSchema } from './BuddyRunEvent' interface RunEventRow { @@ -83,6 +84,42 @@ export class RunEventQueries { return rows.map(row => row.id) } + listRecoveryRuns(): RunEventRecoveryRun[] { + const rows = this.#database.prepare(` + SELECT runs.id AS runId, runs.conversation_id AS conversationId, runs.status, + checkpoints.last_sequence AS lastSequence, + checkpoints.projection_version AS projectionVersion, + checkpoints.file_fingerprint AS fileFingerprint, + COALESCE((SELECT MAX(sequence) FROM run_events WHERE run_id = runs.id), 0) AS projectedLastSequence + FROM runs LEFT JOIN run_event_checkpoints AS checkpoints ON checkpoints.run_id = runs.id + ORDER BY CASE WHEN runs.status IN ('queued', 'running') THEN 0 ELSE 1 END, runs.started_at, runs.id + `).all() as unknown as Array<{ + conversationId: string + fileFingerprint: string | null + lastSequence: number | null + projectedLastSequence: number + projectionVersion: number | null + runId: string + status: string + }> + return rows.map(row => ({ + conversationId: row.conversationId, + runId: row.runId, + status: row.status, + projectedLastSequence: row.projectedLastSequence, + checkpoint: row.lastSequence === null || row.projectionVersion === null || row.fileFingerprint === null + ? null + : { lastSequence: row.lastSequence, projectionVersion: row.projectionVersion, fileFingerprint: row.fileFingerprint }, + })) + } + + hasCheckpoint(runId: string, checkpoint: RunEventCheckpoint): boolean { + return this.#database.prepare(` + SELECT 1 FROM run_event_checkpoints + WHERE run_id = ? AND file_fingerprint = ? AND projection_version = ? AND last_sequence = ? + `).get(runId, checkpoint.fileFingerprint, checkpoint.projectionVersion, checkpoint.lastSequence) !== undefined + } + listRunIds(): string[] { const rows = this.#database.prepare(` SELECT id FROM runs ORDER BY started_at, id diff --git a/apps/buddy/service/src/events/RunEventRecovery.ts b/apps/buddy/service/src/events/RunEventRecovery.ts new file mode 100644 index 00000000..a8d761c0 --- /dev/null +++ b/apps/buddy/service/src/events/RunEventRecovery.ts @@ -0,0 +1,49 @@ +// Bump when projection semantics or the compaction plan become incompatible. +export const RUN_EVENT_PROJECTION_VERSION = 1 + +export interface RunEventCheckpoint { + fileFingerprint: string + lastSequence: number + projectionVersion: number +} + +export interface RunEventRecoveryRun { + checkpoint: RunEventCheckpoint | null + conversationId: string + projectedLastSequence: number + runId: string + status: string +} + +export interface RunEventRecoverySummary { + durationMs: number + enumerationMs: number + inspectionMs: number + readMs: number + projectionMs: number + scanned: number + skipped: number + rebuilt: number + events: number + failed: number +} + +export function canReuseRunEventCheckpoint( + run: RunEventRecoveryRun, + fileFingerprint: string | null, +): run is RunEventRecoveryRun & { checkpoint: RunEventCheckpoint } { + const checkpoint = run.checkpoint + return isTerminalRecoveryRun(run.status) + && checkpoint !== null + && checkpoint.projectionVersion === RUN_EVENT_PROJECTION_VERSION + && Number.isSafeInteger(checkpoint.lastSequence) + && checkpoint.lastSequence >= 0 + && checkpoint.lastSequence < Number.MAX_SAFE_INTEGER + && checkpoint.lastSequence === run.projectedLastSequence + && fileFingerprint !== null + && checkpoint.fileFingerprint === fileFingerprint +} + +export function isTerminalRecoveryRun(status: string): boolean { + return status === 'completed' || status === 'failed' || status === 'cancelled' +} diff --git a/apps/buddy/service/src/events/RunEventStore.ts b/apps/buddy/service/src/events/RunEventStore.ts index 796501aa..f0f4d2aa 100644 --- a/apps/buddy/service/src/events/RunEventStore.ts +++ b/apps/buddy/service/src/events/RunEventStore.ts @@ -8,6 +8,7 @@ import { open, readdir, readFile, + stat, unlink, } from 'node:fs/promises' import { dirname, join } from 'node:path' @@ -94,8 +95,22 @@ export class RunEventStore { } } - async readAndRepair(runId: string): Promise { - const path = this.#eventPath(runId) + async fingerprint(runId: string, conversationId?: string): Promise { + try { + const metadata = await stat(this.#eventPath(runId, conversationId), { bigint: true }) + if (!metadata.isFile()) + throw new Error('Lexora Buddy event path is not a regular file') + return [metadata.dev, metadata.ino, metadata.size, metadata.mtimeNs, metadata.ctimeNs, metadata.birthtimeNs].join(':') + } + catch (cause) { + if (isFileNotFound(cause)) + return null + throw new RunEventStorageError(runEventFailureScope(runId), 'read', 'stat', 'not_applicable', { cause }) + } + } + + async readAndRepair(runId: string, conversationId?: string): Promise { + const path = this.#eventPath(runId, conversationId) const readScope = runEventFailureScope(runId) let content: string try { @@ -192,12 +207,12 @@ export class RunEventStore { } } - async listPersistedRunIds(runIds: readonly string[]): Promise { + async listPersistedRunIds(runIds: readonly string[], conversations?: ReadonlyMap): Promise { const entriesByDirectory = new Map>() const persistedRunIds: string[] = [] for (const runId of runIds) { const scope = runEventFailureScope(runId) - const directory = this.#eventDirectory(runId) + const directory = this.#eventDirectory(runId, conversations?.get(runId)) let entries = entriesByDirectory.get(directory) if (!entries) { let names: string[] @@ -243,12 +258,12 @@ export class RunEventStore { return persistedRunIds } - #eventPath(runId: string): string { - return join(this.#eventDirectory(runId), `${buddyRunIdSchema.parse(runId)}.jsonl`) + #eventPath(runId: string, conversationId?: string): string { + return join(this.#eventDirectory(runId, conversationId), `${buddyRunIdSchema.parse(runId)}.jsonl`) } - #eventDirectory(runId: string): string { - const conversationId = this.#resolveConversationId(runId) + #eventDirectory(runId: string, knownConversationId?: string): string { + const conversationId = knownConversationId ?? this.#resolveConversationId(runId) if (conversationId === null) throw new Error(`Lexora Buddy run was not found: ${runId}`) return join( diff --git a/apps/buddy/service/src/events/__tests__/RunEventLog.spec.ts b/apps/buddy/service/src/events/__tests__/RunEventLog.spec.ts index 4e974a09..42e839e9 100644 --- a/apps/buddy/service/src/events/__tests__/RunEventLog.spec.ts +++ b/apps/buddy/service/src/events/__tests__/RunEventLog.spec.ts @@ -202,7 +202,8 @@ describe('runEventLog', () => { commitState: 'unknown', operation: 'append', runId: 'run-1', - stage: 'open', + // Windows can open a directory and reject the write rather than the open. + stage: expect.stringMatching(/^(?:open|write)$/), }) expect(failures).toEqual([failure]) }) diff --git a/apps/buddy/service/src/events/__tests__/RunEventRecovery.spec.ts b/apps/buddy/service/src/events/__tests__/RunEventRecovery.spec.ts new file mode 100644 index 00000000..5cb3c11a --- /dev/null +++ b/apps/buddy/service/src/events/__tests__/RunEventRecovery.spec.ts @@ -0,0 +1,445 @@ +import type { DatabaseSync } from 'node:sqlite' +import type { AppendBuddyRunEventInput, BuddyRunEvent } from '../BuddyRunEvent' +import type { RunEventLogCallbacks } from '../RunEventLog' +import { appendFile, mkdir, mkdtemp, readFile, rename, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import process from 'node:process' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { openBuddyDatabase } from '../../storage/database' +import { createRunEventLog } from '../createRunEventLog' +import { RunEventProjector } from '../RunEventProjector' +import { RunEventQueries } from '../RunEventQueries' +import { RUN_EVENT_PROJECTION_VERSION } from '../RunEventRecovery' +import { RunEventStore } from '../RunEventStore' + +const NOW = '2026-10-04T00:00:00.000Z' +const databases = new Set() +const directories: string[] = [] + +afterEach(async () => { + vi.restoreAllMocks() + for (const database of databases) + database.close() + databases.clear() + await Promise.all(directories.splice(0).map(path => rm(path, { force: true, recursive: true }))) +}) + +describe('startup run event recovery', () => { + it.skipIf(process.env.LEXORA_RUN_RECOVERY_BENCHMARK !== '1')('measures ten warmed full/skip recovery pairs on isolated persistent data', async () => { + const f = await fixture(true) + const eventStore = store(f) + const runCount = 60 + const eventsPerRun = 80 + for (let index = 0; index < runCount; index++) { + const runId = index === 0 ? 'run-1' : `benchmark-${index}` + if (index > 0) { + f.database.prepare(`INSERT INTO runs (id, conversation_id, branch_id, triggering_message_id, provider, model, purpose, status, started_at) + SELECT ?, conversation_id, branch_id, triggering_message_id, provider, model, purpose, 'running', started_at FROM runs WHERE id = 'run-1'`).run(runId) + } + const events: BuddyRunEvent[] = Array.from({ length: eventsPerRun }, (_, position) => ({ + runId, + sequence: position + 1, + type: position === eventsPerRun - 1 ? 'run.completed' : 'audit.sample', + payload: position === eventsPerRun - 1 ? {} : { text: 'fixture'.repeat(128) }, + createdAt: NOW, + })) + await eventStore.append(events) + } + await f.log.recoverAll({ force: true }) + const full: number[] = [] + const skipped: number[] = [] + for (let sample = 0; sample < 10; sample++) { + await f.restart() + const rebuilt = await f.log.recoverAll({ force: true }) + expect(rebuilt.rebuilt).toBe(runCount) + full.push(rebuilt.durationMs) + await f.restart() + const reused = await f.log.recoverAll() + expect(reused).toMatchObject({ skipped: runCount, rebuilt: 0, events: 0 }) + skipped.push(reused.durationMs) + } + const quantiles = (values: number[]) => { + const sorted = values.toSorted((a, b) => a - b) + return { medianMs: (sorted[4]! + sorted[5]!) / 2, p90Ms: sorted[8]! } + } + const fullSummary = quantiles(full) + const skippedSummary = quantiles(skipped) + process.stdout.write(`${JSON.stringify({ + event: 'run.recovery.benchmark', + samples: 10, + warmed: true, + runs: runCount, + eventsPerRun, + full: fullSummary, + skipped: skippedSummary, + reductionPercent: (1 - skippedSummary.medianMs / fullSummary.medianMs) * 100, + })}\n`) + // Timing is diagnostic, not a flaky unit-test pass threshold. + const before = projection(f.database) + await f.log.recoverAll({ force: true }) + expect(projection(f.database)).toEqual(before) + }, 60_000) + + it('rebuilds old data once, then skips stable history across database reopen without reading or projecting', async () => { + const f = await fixture(true) + await seedHistory(f.log) + expect(await f.log.recoverAll()).toMatchObject({ scanned: 1, rebuilt: 1, skipped: 0, failed: 0 }) + const before = projection(f.database) + await f.reopen() + const read = vi.spyOn(RunEventStore.prototype, 'readAndRepair') + const rebuild = vi.spyOn(RunEventProjector.prototype, 'rebuild') + const lookup = vi.spyOn(RunEventQueries.prototype, 'findConversationId') + expect(await f.log.recoverAll()).toMatchObject({ scanned: 1, rebuilt: 0, skipped: 1, events: 0 }) + expect(read).not.toHaveBeenCalled() + expect(rebuild).not.toHaveBeenCalled() + expect(lookup).not.toHaveBeenCalled() + expect(projection(f.database)).toEqual(before) + expect(await f.log.append({ runId: 'run-1', type: 'audit.after_restart', payload: {} })).toMatchObject({ sequence: 6 }) + expect(read).not.toHaveBeenCalled() + expect(checkpoint(f.database)).toBeUndefined() + }) + + it('only rebuilds the changed run and leaves other stable history alone', async () => { + const f = await fixture() + await seedHistory(f.log) + f.database.exec(`INSERT INTO runs (id, conversation_id, branch_id, triggering_message_id, provider, model, purpose, status, started_at) + SELECT 'run-2', conversation_id, branch_id, triggering_message_id, provider, model, purpose, 'running', started_at FROM runs WHERE id = 'run-1'`) + await f.log.append({ runId: 'run-2', type: 'run.completed', payload: {} }) + await f.log.recoverAll() + await f.log.append({ runId: 'run-1', type: 'audit.new', payload: {} }) + const read = vi.spyOn(RunEventStore.prototype, 'readAndRepair') + expect(await f.log.recoverAll()).toMatchObject({ scanned: 2, skipped: 1, rebuilt: 1 }) + expect(read.mock.calls.map(call => call[0])).toEqual(['run-1']) + expect(await f.log.list('run-1')).toHaveLength(6) + }) + + it('forces full validation without changing messages, approvals, usage or event results', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + const before = projection(f.database) + const read = vi.spyOn(RunEventStore.prototype, 'readAndRepair') + expect(await f.log.recoverAll({ force: true })).toMatchObject({ skipped: 0, rebuilt: 1, events: 5 }) + expect(await f.log.replay('run-1')).toBe(5) + expect(await f.log.replayAll()).toBe(5) + expect(read).toHaveBeenCalledTimes(3) + expect(projection(f.database)).toEqual(before) + }) + + it('does not mutate the log when checkpoint invalidation fails', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + const before = await readFile(f.eventFile, 'utf8') + f.database.exec(`CREATE TRIGGER reject_invalidation BEFORE DELETE ON run_event_checkpoints + BEGIN SELECT RAISE(ABORT, 'invalidation failed'); END;`) + const append = vi.spyOn(RunEventStore.prototype, 'append') + await expect(f.log.append({ runId: 'run-1', type: 'audit.new', payload: {} })).rejects.toMatchObject({ + code: 'EVENT_PROJECTION_FAILED', + commitState: 'not_applicable', + }) + expect(append).not.toHaveBeenCalled() + expect(await readFile(f.eventFile, 'utf8')).toBe(before) + expect(checkpoint(f.database)).toBeDefined() + expect(f.log.state).toBe('failed') + }) + + it('rolls back projection and checkpoint together when checkpoint publication fails', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + f.database.exec(`UPDATE messages SET content_json = '{"text":"damaged"}' WHERE id = 'answer-1'; + CREATE TRIGGER reject_checkpoint BEFORE INSERT ON run_event_checkpoints + BEGIN SELECT RAISE(ABORT, 'checkpoint failed'); END;`) + const damaged = projection(f.database) + await expect(f.log.recoverAll()).rejects.toMatchObject({ code: 'EVENT_PROJECTION_FAILED' }) + expect(checkpoint(f.database)).toBeUndefined() + expect(projection(f.database)).toEqual(damaged) + f.database.exec('DROP TRIGGER reject_checkpoint') + await f.restart() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(f.database.prepare('SELECT content_json FROM messages WHERE id = \'answer-1\'').get()).toEqual({ content_json: '{"text":"answer"}' }) + }) + + it('recovers the log-durable / projection-uncommitted crash window', async () => { + const f = await fixture(true) + await seedHistory(f.log) + await f.log.recoverAll() + new RunEventProjector(f.database).invalidateCheckpoint('run-1') + await store(f).append([{ runId: 'run-1', sequence: 6, type: 'audit.crash', payload: {}, createdAt: NOW }]) + expect(await f.log.list('run-1')).toHaveLength(5) + await f.reopen() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0, events: 6 }) + expect(await f.log.list('run-1')).toHaveLength(6) + const before = projection(f.database) + await f.log.recoverAll({ force: true }) + expect(projection(f.database)).toEqual(before) + }) + + it('safely rebuilds after invalidation even if the process stopped before modifying the log', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + new RunEventProjector(f.database).invalidateCheckpoint('run-1') + await f.restart() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + }) + + it.each([ + 'DELETE FROM run_events WHERE run_id = \'run-1\' AND sequence = 2', + 'DELETE FROM messages WHERE id = \'answer-1\'', + 'UPDATE messages SET content_json = \'{"text":"damaged"}\' WHERE id = \'answer-1\'', + 'DELETE FROM approvals WHERE id = \'approval-1\'', + 'UPDATE approvals SET status = \'denied\' WHERE id = \'approval-1\'', + 'DELETE FROM usage_records WHERE id = \'usage-1\'', + 'UPDATE usage_records SET total_cost = 99 WHERE id = \'usage-1\'', + 'UPDATE runs SET completed_at = \'2026-10-01T00:00:00.000Z\' WHERE id = \'run-1\'', + ])('invalidates direct database projection writes and restores the supported projection: %s', async (sql) => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + const before = projection(f.database) + f.database.exec(sql) + expect(checkpoint(f.database)).toBeUndefined() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(projection(f.database)).toEqual(before) + }) + + it('does not reread stable terminal history for a retained noncompactable tool update', async () => { + const f = await fixture() + await f.log.appendBatch([ + { runId: 'run-1', type: 'tool.updated', payload: { toolCallId: 'tool-1' } }, + { runId: 'run-1', type: 'run.completed', payload: {} }, + ]) + await f.log.recoverAll() + await f.restart() + const read = vi.spyOn(RunEventStore.prototype, 'readAndRepair') + expect(await f.log.recoverAll()).toMatchObject({ skipped: 1, rebuilt: 0 }) + expect(await f.log.compactTerminalRuns()).toBe(0) + expect(read).not.toHaveBeenCalled() + // An explicit maintenance request remains a full inspection. + expect(await f.log.compactTerminalRun('run-1')).toBe(0) + expect(read).toHaveBeenCalledTimes(1) + }) + + it('invalidates compaction, then certifies retained sequences without renumbering', async () => { + const f = await fixture() + await f.log.appendBatch([ + { runId: 'run-1', type: 'run.started', payload: {} }, + { runId: 'run-1', type: 'run.progress', payload: {} }, + { runId: 'run-1', type: 'run.completed', payload: {} }, + ]) + await f.log.recoverAll() + expect(checkpoint(f.database)).toBeUndefined() + expect(await f.log.compactTerminalRuns()).toBe(1) + await f.restart() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, events: 2 }) + await f.restart() + expect(await f.log.recoverAll()).toMatchObject({ skipped: 1 }) + const read = vi.spyOn(RunEventStore.prototype, 'readAndRepair') + expect(await f.log.append({ runId: 'run-1', type: 'audit.next', payload: {} })).toMatchObject({ sequence: 4 }) + expect(read).not.toHaveBeenCalled() + }) + + it('detects replacement even when the final sequence and file size have not changed', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + const content = await readFile(f.eventFile, 'utf8') + const replaced = content.replace('"text":"answer"', '"text":"edited"') + expect(replaced.length).toBe(content.length) + const replacement = `${f.eventFile}.replacement` + await writeFile(replacement, replaced) + await rm(f.eventFile) + await rename(replacement, f.eventFile) + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(f.database.prepare('SELECT content_json FROM messages WHERE id = \'answer-1\'').get()).toEqual({ content_json: '{"text":"edited"}' }) + }) + + it('repairs a torn tail but does not certify the changed file until a subsequent clean read', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + await appendFile(f.eventFile, '{"runId":"run-1"') + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(checkpoint(f.database)).toBeUndefined() + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1 }) + expect(await f.log.recoverAll()).toMatchObject({ skipped: 1 }) + }) + + it.each([ + 'UPDATE run_event_checkpoints SET last_sequence = last_sequence + 100', + 'UPDATE run_event_checkpoints SET projection_version = projection_version + 1', + 'DELETE FROM run_event_checkpoints', + ])('rebuilds an incompatible or missing checkpoint: %s', async (sql) => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + f.database.exec(sql) + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(checkpoint(f.database)).toMatchObject({ last_sequence: 5, projection_version: RUN_EVENT_PROJECTION_VERSION }) + }) + + it('does not skip a running task even with an externally inserted checkpoint', async () => { + const f = await fixture() + await f.log.append({ runId: 'run-1', type: 'run.started', payload: {} }) + const fingerprint = await store(f).fingerprint('run-1') + f.database.prepare('INSERT INTO run_event_checkpoints VALUES (?, ?, ?, ?)').run('run-1', 1, RUN_EVENT_PROJECTION_VERSION, fingerprint) + expect(await f.log.recoverAll()).toMatchObject({ rebuilt: 1, skipped: 0 }) + expect(checkpoint(f.database)).toBeUndefined() + }) + + it('invalidates missing logs without erasing projections or inventing unknown runs', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + const before = projection(f.database) + await rm(f.eventFile) + await writeFile(join(f.conversationsDirectory, 'conversation-1', 'events', 'unknown-run.jsonl'), '{broken}\n') + expect(await f.log.recoverAll()).toMatchObject({ scanned: 0, skipped: 0, rebuilt: 0 }) + expect(checkpoint(f.database)).toBeUndefined() + expect(projection(f.database)).toEqual(before) + expect(f.database.prepare('SELECT id FROM runs').all()).toEqual([{ id: 'run-1' }]) + }) + + it('does not revive a soft-deleted conversation and cascades checkpoints on physical run deletion', async () => { + const f = await fixture() + await seedHistory(f.log) + await f.log.recoverAll() + f.database.prepare('UPDATE conversations SET deleted_at = ? WHERE id = ?').run(NOW, 'conversation-1') + expect(await f.log.recoverAll()).toMatchObject({ skipped: 1 }) + expect(f.database.prepare('SELECT deleted_at FROM conversations').get()).toEqual({ deleted_at: NOW }) + f.database.exec('DELETE FROM runs WHERE id = \'run-1\'') + expect(checkpoint(f.database)).toBeUndefined() + expect(await f.log.recoverAll()).toMatchObject({ scanned: 0 }) + }) + + it('reports failures and does not hide corrupt log errors behind a checkpoint', async () => { + const onRecovery = vi.fn() + const f = await fixture(false, { onRecovery }) + await seedHistory(f.log) + await f.log.recoverAll() + await appendFile(f.eventFile, '{broken}\n') + await expect(f.log.recoverAll()).rejects.toMatchObject({ code: 'EVENT_LOG_CORRUPTED' }) + expect(onRecovery).toHaveBeenLastCalledWith(expect.objectContaining({ failed: 1, skipped: 0 })) + expect(f.log.state).toBe('failed') + expect(checkpoint(f.database)).toBeUndefined() + }) + + it('ignores diagnostic observer exceptions and preserves failure on invalid file paths', async () => { + const f = await fixture(false, { onRecovery: () => { + throw new Error('observer failed') + } }) + expect(await f.log.recoverAll()).toMatchObject({ failed: 0 }) + await mkdir(f.eventFile, { recursive: true }) + await expect(f.log.recoverAll()).rejects.toMatchObject({ code: 'EVENT_STORAGE_FAILED', stage: 'stat' }) + }) +}) + +async function fixture(persistent = false, callbacks: RunEventLogCallbacks = {}) { + const root = await mkdtemp(join(tmpdir(), 'buddy-event-recovery-')) + directories.push(root) + const databasePath = persistent ? join(root, 'buddy.sqlite3') : ':memory:' + const database = openBuddyDatabase({ databasePath }) + databases.add(database) + database.exec(` + INSERT INTO conversations (id, created_at, updated_at) VALUES ('conversation-1', '${NOW}', '${NOW}'); + INSERT INTO conversation_branches (id, conversation_id, created_at) VALUES ('branch-1', 'conversation-1', '${NOW}'); + INSERT INTO runs (id, conversation_id, branch_id, triggering_message_id, provider, model, purpose, status, started_at) + VALUES ('run-1', 'conversation-1', 'branch-1', 'user-1', 'fixture', 'fixture', 'chat', 'running', '${NOW}'); + `) + const conversationsDirectory = join(root, 'conversations') + const create = (db: DatabaseSync) => createRunEventLog({ conversationsDirectory, database: db, ...callbacks }) + const f = { + root, + database, + conversationsDirectory, + eventFile: join(conversationsDirectory, 'conversation-1', 'events', 'run-1.jsonl'), + log: create(database), + async restart() { + await f.log.close() + f.log = create(f.database) + }, + async reopen() { + await f.log.close() + f.database.close() + databases.delete(f.database) + f.database = openBuddyDatabase({ databasePath }) + databases.add(f.database) + f.log = create(f.database) + }, + } + return f +} + +function store(f: Awaited>) { + return new RunEventStore({ conversationsDirectory: f.conversationsDirectory, resolveConversationId: () => 'conversation-1' }) +} + +function checkpoint(database: DatabaseSync) { + return database.prepare('SELECT * FROM run_event_checkpoints WHERE run_id = \'run-1\'').get() +} + +function projection(database: DatabaseSync) { + return { + runs: database.prepare('SELECT * FROM runs ORDER BY id').all(), + messages: database.prepare('SELECT rowid, * FROM messages ORDER BY id').all(), + approvals: database.prepare('SELECT * FROM approvals ORDER BY id').all(), + usage: database.prepare('SELECT * FROM usage_records ORDER BY id').all(), + events: database.prepare('SELECT * FROM run_events ORDER BY run_id, sequence').all(), + } +} + +async function seedHistory(log: ReturnType) { + const events: AppendBuddyRunEventInput[] = [ + { + runId: 'run-1', + type: 'message.completed', + createdAt: NOW, + payload: { messageId: 'answer-1', role: 'assistant', content: { text: 'answer' }, stopReason: 'completed' }, + }, + { + runId: 'run-1', + type: 'approval.requested', + createdAt: NOW, + payload: { + id: 'approval-1', + runId: 'run-1', + toolCallId: 'tool-1', + kind: 'read', + status: 'pending', + summary: 'Read fixture', + createdAt: NOW, + resolvedAt: null, + payload: { card: 'paths', access: 'read', grant: null, toolName: 'read', targets: [{ path: '/fixture.txt', zone: 'workspace' }] }, + }, + }, + { runId: 'run-1', type: 'approval.resolved', createdAt: NOW, payload: { id: 'approval-1', status: 'approved', resolvedAt: NOW } }, + { + runId: 'run-1', + type: 'usage.recorded', + createdAt: NOW, + payload: { + usageRecordId: 'usage-1', + sourceEntryId: 'entry-1', + provider: 'fixture', + model: 'fixture', + purpose: 'turn', + inputTokens: 10, + outputTokens: 5, + cacheReadTokens: 0, + cacheWriteTokens: 0, + reasoningTokens: null, + totalTokens: 15, + inputCost: 0.01, + outputCost: 0.02, + cacheReadCost: 0, + cacheWriteCost: 0, + totalCost: 0.03, + }, + }, + { runId: 'run-1', type: 'run.completed', createdAt: NOW, payload: {} }, + ] + await log.appendBatch(events) +} diff --git a/apps/buddy/service/src/events/createRunEventLog.ts b/apps/buddy/service/src/events/createRunEventLog.ts index ce05604a..665747ee 100644 --- a/apps/buddy/service/src/events/createRunEventLog.ts +++ b/apps/buddy/service/src/events/createRunEventLog.ts @@ -15,6 +15,7 @@ export function createRunEventLog(options: CreateRunEventLogOptions): RunEventLo return new RunEventLog({ onObserverError: options.onObserverError, onFatalFailure: options.onFatalFailure, + onRecovery: options.onRecovery, projector: new RunEventProjector(options.database), queries, store: new RunEventStore({ diff --git a/apps/buddy/service/src/index.ts b/apps/buddy/service/src/index.ts index cf226665..a7256246 100644 --- a/apps/buddy/service/src/index.ts +++ b/apps/buddy/service/src/index.ts @@ -129,6 +129,14 @@ async function runBuddyService(): Promise { conversationsDirectory: join(buddyHome, 'conversations'), database: openedDatabase, onObserverError: error => record({ event: 'run.observer_failed', level: 'warn', errorCode: readDiagnosticErrorCode(error) }), + onRecovery: (summary) => { + record({ event: 'run.recovery.enumerated', level: 'info', durationMs: summary.enumerationMs, count: summary.scanned }) + record({ event: 'run.recovery.inspected', level: 'info', durationMs: summary.inspectionMs, count: summary.scanned }) + record({ event: 'run.recovery.read', level: 'info', durationMs: summary.readMs, count: summary.events }) + record({ event: 'run.recovery.projected', level: 'info', durationMs: summary.projectionMs, count: summary.rebuilt }) + record({ event: 'run.recovery.skipped', level: 'info', count: summary.skipped }) + record({ event: 'run.recovery.finished', level: summary.failed ? 'error' : 'info', durationMs: summary.durationMs, count: summary.failed }) + }, onFatalFailure: (error) => { record({ event: 'run.storage_failed', level: 'error', runId: error.runId, errorCode: error.code }) notifyFailure(readBuddyServiceFailureCode(error)) @@ -157,7 +165,9 @@ async function runBuddyService(): Promise { }) return log }, ['runtime.database']) - await host.step('runtime.event_replay', () => eventLog!.replayAll()) + await host.step('runtime.event_replay', () => eventLog!.recoverAll({ + force: process.env.LEXORA_BUDDY_EVENT_RECOVERY === 'full', + })) await host.start('runtime.services', async ({ defer }) => { serviceHandle = await startBuddyService({ buddyHome, diff --git a/apps/buddy/service/src/storage/__tests__/schema.spec.ts b/apps/buddy/service/src/storage/__tests__/schema.spec.ts index 52044d57..022e8f21 100644 --- a/apps/buddy/service/src/storage/__tests__/schema.spec.ts +++ b/apps/buddy/service/src/storage/__tests__/schema.spec.ts @@ -80,6 +80,41 @@ function seedRun( } describe('buddy schema', { timeout: MIGRATION_TEST_TIMEOUT }, () => { + it('upgrades v22 with empty checkpoints without changing historical data', () => { + const root = mkdtempSync(join(tmpdir(), 'buddy-checkpoint-migration-')) + directories.push(root) + const databasePath = join(root, 'buddy.sqlite3') + const legacy = openMigrationFixtureDatabase(databasePath) + for (const migration of BUDDY_SCHEMA_MIGRATIONS.filter(migration => migration.version <= 22)) { + legacy.exec(migration.sql) + legacy.exec(`PRAGMA user_version = ${migration.version}`) + } + seedRun(legacy) + const before = { + runs: legacy.prepare('SELECT * FROM runs').all(), + messages: legacy.prepare('SELECT * FROM messages').all(), + conversations: legacy.prepare('SELECT * FROM conversations').all(), + } + legacy.close() + const migrated = openBuddyDatabase({ databasePath }) + databases.push(migrated) + expect(migrated.prepare('SELECT * FROM run_event_checkpoints').all()).toEqual([]) + expect(migrated.prepare('SELECT * FROM runs').all()).toEqual(before.runs) + expect(migrated.prepare('SELECT * FROM messages').all()).toEqual(before.messages) + expect(migrated.prepare('SELECT * FROM conversations').all()).toEqual(before.conversations) + expect(migrated.prepare('PRAGMA foreign_key_check').all()).toEqual([]) + expect(migrated.prepare('PRAGMA user_version').get()).toEqual({ user_version: BUDDY_SCHEMA_VERSION }) + }) + + it('rejects a current database missing checkpoint invalidation triggers', () => { + const root = mkdtempSync(join(tmpdir(), 'buddy-checkpoint-schema-')) + directories.push(root) + const databasePath = join(root, 'buddy.sqlite3') + const database = openBuddyDatabase({ databasePath }) + database.exec('DROP TRIGGER invalidate_run_checkpoint_messages_delete') + database.close() + expect(() => openBuddyDatabase({ databasePath })).toThrow(/incomplete checkpoint invalidation schema/i) + }) it('preserves v21 usage and attributes independent actions without modifying completed runs', () => { const directory = mkdtempSync(join(tmpdir(), 'buddy-action-migration-')) directories.push(directory) diff --git a/apps/buddy/service/src/storage/__tests__/usageAnalyticsRepository.spec.ts b/apps/buddy/service/src/storage/__tests__/usageAnalyticsRepository.spec.ts index ba732864..160de26b 100644 --- a/apps/buddy/service/src/storage/__tests__/usageAnalyticsRepository.spec.ts +++ b/apps/buddy/service/src/storage/__tests__/usageAnalyticsRepository.spec.ts @@ -5,6 +5,7 @@ import { join } from 'node:path' import { afterEach, describe, expect, it } from 'vitest' import { usageAnalyticsSchema, usagePeriodSchema, usageTopTasksRequestSchema, usageTrendRequestSchema, usageTrendSchema } from '../../../../shared/usage/usageAnalyticsApi' import { openBuddyDatabase } from '../database' +import { BUDDY_RUN_EVENT_CHECKPOINT_TRIGGER_NAMES } from '../migrations/v23RunEventCheckpoints' import { createUsageAnalyticsRepository } from '../usageAnalyticsRepository' const databases: DatabaseSync[] = [] @@ -119,6 +120,8 @@ describe('usage analytics', () => { f.usage('retained') const before = f.repository.analytics(period) f.database.exec(` + ${BUDDY_RUN_EVENT_CHECKPOINT_TRIGGER_NAMES.map(name => `DROP TRIGGER ${name};`).join('\n ')} + DROP TABLE run_event_checkpoints; DROP TABLE extension_invocations; DROP TABLE connector_tool_catalogs; DROP TABLE skill_file_cleanup; diff --git a/apps/buddy/service/src/storage/database.ts b/apps/buddy/service/src/storage/database.ts index d2210afd..ae4866a0 100644 --- a/apps/buddy/service/src/storage/database.ts +++ b/apps/buddy/service/src/storage/database.ts @@ -4,6 +4,7 @@ import { mkdirSync } from 'node:fs' import { dirname, isAbsolute, join } from 'node:path' import { DatabaseSync as NodeDatabaseSync } from 'node:sqlite' import { BUDDY_V15_CAPABILITY_OVERRIDES_SCHEMA_SQL, BUDDY_V15_CATALOG_MODEL_ID_SCHEMA_SQL, BUDDY_V15_CATALOG_SELECTION_SCHEMA_SQL, BUDDY_V15_PROVIDER_INSTANCES_SCHEMA_SQL, BUDDY_V15_REQUEST_HEADERS_SCHEMA_SQL } from './migrations/v15ModelServices' +import { BUDDY_RUN_EVENT_CHECKPOINT_TRIGGER_NAMES } from './migrations/v23RunEventCheckpoints' import { BUDDY_SCHEMA_MIGRATIONS, @@ -17,6 +18,7 @@ export interface OpenBuddyDatabaseOptions { } const BUDDY_CURRENT_SCHEMA_COLUMNS = { + run_event_checkpoints: ['run_id', 'last_sequence', 'projection_version', 'file_fingerprint'], extension_invocations: ['id', 'extension_id', 'action_id', 'conversation_id', 'trigger', 'status', 'branch_id', 'source_message_id', 'extension_name', 'action_title', 'result_message'], usage_records: ['run_id', 'invocation_id'], conversations: ['title_source', 'title_revision'], @@ -150,6 +152,9 @@ function assertCurrentSchema(database: DatabaseSync, excludedTables: readonly st if (requiredColumns.some(column => !columns.has(column) && !(table === 'provider_model_states' && excludedModelColumns.includes(column)))) throw new BuddyDatabaseVersionError('incomplete schema version') } + const triggers = new Set((database.prepare('SELECT name FROM sqlite_master WHERE type = \'trigger\'').all() as Array<{ name: string }>).map(row => row.name)) + if (BUDDY_RUN_EVENT_CHECKPOINT_TRIGGER_NAMES.some(name => !triggers.has(name))) + throw new BuddyDatabaseVersionError('incomplete checkpoint invalidation schema') } function applyMigration(database: DatabaseSync, migration: BuddySchemaMigration): void { diff --git a/apps/buddy/service/src/storage/migrations/v23RunEventCheckpoints.ts b/apps/buddy/service/src/storage/migrations/v23RunEventCheckpoints.ts new file mode 100644 index 00000000..f074b2a0 --- /dev/null +++ b/apps/buddy/service/src/storage/migrations/v23RunEventCheckpoints.ts @@ -0,0 +1,38 @@ +const projectionTables = ['run_events', 'messages', 'approvals', 'usage_records'] as const + +// Database-level invalidation covers existing repository writes as well as the projector. +// An absent checkpoint is untrusted; only a successful full rebuild publishes one. +export const BUDDY_RUN_EVENT_CHECKPOINT_TRIGGER_NAMES = [ + ...projectionTables.flatMap(table => ['insert', 'update', 'delete'].map(operation => `invalidate_run_checkpoint_${table}_${operation}`)), + 'invalidate_run_checkpoint_runs_update', +] + +export const BUDDY_V23_RUN_EVENT_CHECKPOINTS_SCHEMA_SQL = ` +CREATE TABLE run_event_checkpoints ( + run_id TEXT PRIMARY KEY REFERENCES runs(id) ON DELETE CASCADE, + last_sequence INTEGER NOT NULL CHECK (last_sequence >= 0), + projection_version INTEGER NOT NULL CHECK (projection_version > 0), + file_fingerprint TEXT NOT NULL +); + +${projectionTables.map(table => ` +CREATE TRIGGER invalidate_run_checkpoint_${table}_insert AFTER INSERT ON ${table} +BEGIN + DELETE FROM run_event_checkpoints WHERE run_id = NEW.run_id; +END; +CREATE TRIGGER invalidate_run_checkpoint_${table}_update AFTER UPDATE ON ${table} +BEGIN + DELETE FROM run_event_checkpoints WHERE run_id IN (OLD.run_id, NEW.run_id); +END; +CREATE TRIGGER invalidate_run_checkpoint_${table}_delete AFTER DELETE ON ${table} +BEGIN + DELETE FROM run_event_checkpoints WHERE run_id = OLD.run_id; +END; +`).join('\n')} + +CREATE TRIGGER invalidate_run_checkpoint_runs_update +AFTER UPDATE OF status, started_at, completed_at, error_code, conversation_id, branch_id ON runs +BEGIN + DELETE FROM run_event_checkpoints WHERE run_id IN (OLD.id, NEW.id); +END; +` diff --git a/apps/buddy/service/src/storage/schema.ts b/apps/buddy/service/src/storage/schema.ts index d348b05c..efba64f6 100644 --- a/apps/buddy/service/src/storage/schema.ts +++ b/apps/buddy/service/src/storage/schema.ts @@ -23,6 +23,7 @@ import { BUDDY_V20_TASK_DRAFTS_SCHEMA_SQL } from './migrations/v20TaskDrafts' import { BUDDY_V21_TITLE_SCHEMA_SQL } from './migrations/v21Title' import { BUDDY_V22_EXTENSION_SCHEMA_SQL } from './migrations/v22Extension' +import { BUDDY_V23_RUN_EVENT_CHECKPOINTS_SCHEMA_SQL } from './migrations/v23RunEventCheckpoints' export interface BuddySchemaMigration { foreignKeys?: 'off' @@ -30,7 +31,7 @@ export interface BuddySchemaMigration { version: number } -export const BUDDY_SCHEMA_VERSION = 22 as const +export const BUDDY_SCHEMA_VERSION = 23 as const export const BUDDY_SCHEMA_MIGRATIONS: readonly BuddySchemaMigration[] = [ { sql: BUDDY_V1_INITIAL_SCHEMA_SQL, version: 1 }, @@ -55,4 +56,5 @@ export const BUDDY_SCHEMA_MIGRATIONS: readonly BuddySchemaMigration[] = [ { foreignKeys: 'off', sql: BUDDY_V20_TASK_DRAFTS_SCHEMA_SQL, version: 20 }, { sql: BUDDY_V21_TITLE_SCHEMA_SQL, version: 21 }, { sql: BUDDY_V22_EXTENSION_SCHEMA_SQL, version: 22 }, + { sql: BUDDY_V23_RUN_EVENT_CHECKPOINTS_SCHEMA_SQL, version: 23 }, ]