Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 92 additions & 0 deletions .playwright/scripts/__tests__/runEventRecovery.e2e.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
import fs from 'node:fs/promises'
import path from 'node:path'
import { DatabaseSync } from 'node:sqlite'
import { expect, test } from '../fixtures/electron.mjs'

test('startup recovers history once and only rebuilds changed runs after restart', async ({ buddy }, testInfo) => {
test.setTimeout(120_000)
const instance = await buddy.createInstance('event-recovery')
await instance.launch()
await instance.stop()
const runCount = 24
const messagesPerRun = 20
const now = '2026-10-08T00:00:00.000Z'
const eventDirectory = path.join(instance.home, 'buddy/conversations/recovery-conversation/events')
await fs.mkdir(eventDirectory, { recursive: true, mode: 0o700 })
const databasePath = path.join(instance.home, 'buddy/buddy.sqlite3')
const database = new DatabaseSync(databasePath)
try {
database.exec('PRAGMA foreign_keys = ON')
database.prepare('INSERT INTO conversations (id, created_at, updated_at) VALUES (?, ?, ?)')
.run('recovery-conversation', now, now)
database.prepare('INSERT INTO conversation_branches (id, conversation_id, created_at) VALUES (?, ?, ?)')
.run('recovery-branch', 'recovery-conversation', now)
const insertRun = database.prepare(`INSERT INTO runs
(id, conversation_id, branch_id, triggering_message_id, provider, model, purpose, status, started_at)
VALUES (?, 'recovery-conversation', 'recovery-branch', 'recovery-input', 'fixture', 'fixture', 'chat', 'running', ?)`)
for (let index = 0; index < runCount; index++) {
const runId = `recovery-${index}`
insertRun.run(runId, now)
const events = [
{ type: 'run.started', payload: {} },
...Array.from({ length: messagesPerRun }, (_, message) => ({
type: 'message.completed',
payload: { messageId: `${runId}-answer-${message}`, role: 'assistant', content: { text: 'Recovery fixture' }, stopReason: 'completed' },
})),
{ type: 'run.completed', payload: {} },
].map((event, position) => ({ ...event, runId, sequence: position + 1, createdAt: now }))
await fs.writeFile(path.join(eventDirectory, `${runId}.jsonl`), `${events.map(event => JSON.stringify(event)).join('\n')}\n`, { mode: 0o600 })
}
}
finally {
database.close()
}

const measurements = []
const launch = async () => {
const { page, diagnostics } = await instance.launch()
expect(diagnostics.console.filter(item => item.type === 'pageerror')).toEqual([])
expect((await page.evaluate(() => window.lexoraDesktop.app.startup.getState())).status).toBe('ready')
await instance.stop()
const records = (await fs.readFile(path.join(instance.home, '.runtime/state/logs/application.jsonl'), 'utf8'))
.trim()
.split('\n')
.map(line => JSON.parse(line))
const ready = records.findLast(record => record.event === 'app.ready')
const current = records.filter(record => record.launchId === ready.launchId)
const measurement = {
startupMs: ready.durationMs,
recoveryMs: current.find(record => record.event === 'run.recovery.finished').durationMs,
rebuilt: current.find(record => record.event === 'run.recovery.projected').count,
skipped: current.find(record => record.event === 'run.recovery.skipped').count,
}
measurements.push(measurement)
return measurement
}
const messages = () => {
const database = new DatabaseSync(databasePath, { readOnly: true })
try {
return database.prepare('SELECT rowid, id, content_json FROM messages ORDER BY id').all()
}
finally {
database.close()
}
}

expect(await launch()).toMatchObject({ rebuilt: runCount, skipped: 0 })
const recovered = messages()
expect(recovered).toHaveLength(runCount * messagesPerRun)
expect(await launch()).toMatchObject({ rebuilt: 0, skipped: runCount })
expect(messages()).toEqual(recovered)

await fs.appendFile(path.join(eventDirectory, 'recovery-0.jsonl'), `${JSON.stringify({
runId: 'recovery-0',
sequence: messagesPerRun + 3,
createdAt: now,
type: 'message.completed',
payload: { messageId: 'recovery-late-answer', role: 'assistant', content: { text: 'Durable after checkpoint' }, stopReason: 'completed' },
})}\n`)
expect(await launch()).toMatchObject({ rebuilt: 1, skipped: runCount - 1 })
expect(messages()).toEqual([...recovered, expect.objectContaining({ id: 'recovery-late-answer', content_json: '{"text":"Durable after checkpoint"}' })])
await testInfo.attach('startup-recovery', { body: JSON.stringify(measurements, null, 2), contentType: 'application/json' })
})
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
4 changes: 4 additions & 0 deletions apps/buddy/electron/main/runtime/buddyServiceEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
},
Expand Down
12 changes: 11 additions & 1 deletion apps/buddy/service/src/events/RunEventFailure.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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,
Expand Down
149 changes: 136 additions & 13 deletions apps/buddy/service/src/events/RunEventLog.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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<RunEventRecoverySummary>) => 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'
Expand All @@ -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 {
Expand All @@ -61,9 +68,11 @@ export class RunEventLogClosedError extends Error {

export class RunEventLog implements RunEventLogPort {
readonly #nextSequences = new Map<string, number>()
readonly #compactionCheckedRuns = new Map<string, RunEventCheckpoint>()
readonly #committed: Emitter<EventSnapshot<BuddyRunEvent>>
readonly onDidCommit: RunEventLogPort['onDidCommit']
readonly #onFatalFailure?: (error: RunEventLogFatalError) => void
readonly #onRecovery?: RunEventLogCallbacks['onRecovery']
readonly #projector: RunEventProjectorPort
readonly #queries: RunEventQueryPort
readonly #store: RunEventStorePort
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -147,11 +158,65 @@ export class RunEventLog implements RunEventLogPort {
}

replay(runId: string): Promise<number> {
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<Readonly<RunEventRecoverySummary>> {
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<number> {
Expand All @@ -161,8 +226,16 @@ export class RunEventLog implements RunEventLogPort {
async compactTerminalRuns(): Promise<number> {
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
}

Expand Down Expand Up @@ -200,6 +273,42 @@ export class RunEventLog implements RunEventLogPort {
throw this.#fatalFailure
}

async #replay(
runId: string,
conversationId?: string,
summary?: RunEventRecoverySummary,
knownFingerprint?: string | null,
): Promise<number> {
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<string | null> {
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<number> {
if (!this.#isTerminalRun(runId))
return 0
Expand Down Expand Up @@ -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 }))
Expand All @@ -261,8 +373,19 @@ export class RunEventLog implements RunEventLogPort {
throw this.#fatalFailure
}

async #readAndRepair(runId: string): Promise<BuddyRunEvent[]> {
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<BuddyRunEvent[]> {
this.#invalidateCheckpoint(runId)
return this.#runStoreOperation(() => this.#store.readAndRepair(runId, conversationId))
}

async #runStoreOperation<T>(operation: () => Promise<T>): Promise<T> {
Expand Down
Loading