From b405cfa6ab03735e0f9693635005f962db81ecd8 Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 11 Sep 2026 10:20:23 +0200 Subject: [PATCH 1/2] feat(sdk): add read-only flows replay journal walker (#309) Session-Id: 01a08f80-0bc7-7e03-91e1-fb1f112a0616 --- packages/sdk/src/cli.ts | 5 + packages/sdk/src/cli/replay.ts | 72 +++++++ packages/sdk/src/journal-client.ts | 1 + packages/sdk/src/journal-offset.ts | 57 ++++++ packages/sdk/src/journal-reader.ts | 165 +++++++++++++++ packages/sdk/tests/cli-replay.test.ts | 276 ++++++++++++++++++++++++++ 6 files changed, 576 insertions(+) create mode 100644 packages/sdk/src/cli/replay.ts create mode 100644 packages/sdk/src/journal-offset.ts create mode 100644 packages/sdk/src/journal-reader.ts create mode 100644 packages/sdk/tests/cli-replay.test.ts diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index 4458b515a..786bfd0db 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -16,6 +16,7 @@ import { type RunReport, } from './cli/run.js'; import { runDirectFlow } from './cli/direct-run.js'; +import { parseReplayArgs, replayJournal, type ReplayArgs } from './cli/replay.js'; import { runCloudCli } from './cli/cloud-run.js'; import { isAuthoredFlowPath } from './direct-input.js'; import { runHnMonitor } from './cli/hn-monitor.js'; @@ -35,6 +36,7 @@ export interface CliIo { type CliExitCode = 0 | 1 | 2 | 3; type ParsedArgs = + | ReplayArgs | { command: 'cloud-run'; value: string; json: boolean; wait: boolean } | { command: 'check'; json: boolean; value: string } | { command: 'run'; localAgent: boolean; dataDir: string; input: string | undefined; json: boolean; spawn: boolean; noObserverLink: boolean; value: string } @@ -54,6 +56,7 @@ const USAGE = [ 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir ] [--local-agent] --input ', 'flows tick start --schedule-id --interval-ms [--epoch-ms ] [--max-catch-up ] [--poll-interval-ms ] [--data-dir ] ', 'flows resume [--json] [--no-spawn] [--no-observer-link] [--data-dir ] ', + 'flows replay [--json] [--data-dir ] [--at ]', 'flows observer [--data-dir ]', 'flows hn-monitor start [--data-dir ] [--poll-interval-ms ] ', ].join('\n'); @@ -90,6 +93,7 @@ export async function runCli( } if (parsed.command === 'cloud-run') return runCloudCli(parsed, io); + if (parsed.command === 'replay') return replayJournal(parsed, io); if (parsed.command === 'check') { // Deliberately daemon-free (kernel/DAEMON-LIFECYCLE.md ยง4). `checkFlow` is @@ -356,6 +360,7 @@ function emitWait( function parseArgs(args: readonly string[]): ParsedArgs | undefined { const command = args[0]; + if (command === 'replay') return parseReplayArgs(args.slice(1)); if (command === 'hn-monitor') return parseHnMonitorArgs(args.slice(1)); if (command === 'tick') return parseTickArgs(args.slice(1)); if (command === 'observer') return parseObserverArgs(args.slice(1)); diff --git a/packages/sdk/src/cli/replay.ts b/packages/sdk/src/cli/replay.ts new file mode 100644 index 000000000..80c9a6128 --- /dev/null +++ b/packages/sdk/src/cli/replay.ts @@ -0,0 +1,72 @@ +import type { CliIo } from '../cli.js'; +import { canonicalize } from '../canonical.js'; +import { JournalReadError, walkJournal } from '../journal-client.js'; + +export interface ReplayArgs { + command: 'replay'; + value: string; + json: boolean; + dataDir: string; + at?: string; +} + +export function parseReplayArgs(args: readonly string[]): ReplayArgs | undefined { + let json = false; + let dataDir: string | undefined; + let at: string | undefined; + const positionals: string[] = []; + for (let index = 0; index < args.length; index += 1) { + const argument = args[index]!; + if (argument === '--json') { + if (json) return undefined; + json = true; + } else if (argument === '--data-dir' || argument === '--at') { + const value = args[++index]; + if (value === undefined || value.length === 0 || value.startsWith('-')) return undefined; + if (argument === '--data-dir') { + if (dataDir !== undefined) return undefined; + dataDir = value; + } else { + if (at !== undefined) return undefined; + at = value; + } + } else if (argument.startsWith('-')) { + return undefined; + } else { + positionals.push(argument); + } + } + if (positionals.length !== 1) return undefined; + return { command: 'replay', value: positionals[0]!, json, dataDir: dataDir ?? '.relayflowd', at }; +} + +export async function replayJournal(args: ReplayArgs, io: CliIo): Promise<0 | 1 | 2> { + let emitted = false; + try { + for await (const event of walkJournal(args.value, args.dataDir, { at: args.at })) { + const payload = event.payload !== null && typeof event.payload === 'object' && !Array.isArray(event.payload) + ? event.payload as Record : {}; + io.stdout(args.json ? canonicalize({ + step_id: event.step_id, + kind: event.entry_type, + event, + verification: payload['verification'] ?? null, + spend: payload['budget'] ?? payload['budget_total'] ?? null, + }) + : `${event.seq} ${event.at_ms} ${event.entry_type}` + + (event.step_id === null ? '' : ` step=${JSON.stringify(event.step_id)}`) + + (event.attempt === null ? '' : ` attempt=${event.attempt}`) + + ` ${canonicalize(event.payload)}`); + emitted = true; + } + return 0; + } catch (error) { + const diagnostic = { + severity: 'refusal', + kind: error instanceof JournalReadError ? error.code : 'journal_read_failed', + message: error instanceof Error ? error.message : String(error), + }; + io.stderr(`${emitted ? 'FAILED' : 'REFUSED'} [${diagnostic.kind}] ${diagnostic.message}`); + return emitted ? 1 : 2; + } +} diff --git a/packages/sdk/src/journal-client.ts b/packages/sdk/src/journal-client.ts index ff6aeb184..02013b14f 100644 --- a/packages/sdk/src/journal-client.ts +++ b/packages/sdk/src/journal-client.ts @@ -10,6 +10,7 @@ // loopback double in tests. import { EventEmitter } from 'node:events'; +export { walkJournal, JournalReadError, type JournalEvent, type JournalReadFailure } from './journal-reader.js'; import { randomUUID } from 'node:crypto'; import { createConnection, type Socket } from 'node:net'; import type { VerbContract, EventSubmitParams } from './protocol.js'; diff --git a/packages/sdk/src/journal-offset.ts b/packages/sdk/src/journal-offset.ts new file mode 100644 index 000000000..515c4f970 --- /dev/null +++ b/packages/sdk/src/journal-offset.ts @@ -0,0 +1,57 @@ +import { readFileSync } from 'node:fs'; + +/** Locate a table row's cell in a SQLite snapshot for read-error diagnostics. */ +export function journalRecordOffset(path: string, rootPage: number, seq: number): string { + const main = readFileSync(path); + const encodedSize = main.readUInt16BE(16); + const pageSize = encodedSize === 1 ? 65536 : encodedSize; + const pages = new Map(); + let wal: Buffer | undefined; + try { wal = readFileSync(`${path}-wal`); } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + if (wal !== undefined) { + const frameSize = pageSize + 24; + let lastCommit = 0; + for (let frame = 32; frame + frameSize <= wal.length; frame += frameSize) { + if (wal.readUInt32BE(frame + 4) !== 0) lastCommit = frame; + } + for (let frame = 32; frame <= lastCommit; frame += frameSize) { + pages.set(wal.readUInt32BE(frame), { bytes: wal, offset: frame + 24, file: 'WAL' }); + } + } + const visited = new Set(); + let page = rootPage; + while (!visited.has(page)) { + visited.add(page); + const source = pages.get(page) ?? { bytes: main, offset: (page - 1) * pageSize, file: 'journal' }; + const header = source.offset + (page === 1 ? 100 : 0); + const kind = source.bytes[header]; + if (kind !== 5 && kind !== 13) break; + const count = source.bytes.readUInt16BE(header + 3); + let next = kind === 5 ? source.bytes.readUInt32BE(header + 8) : 0; + for (let index = 0; index < count; index += 1) { + const cell = source.offset + source.bytes.readUInt16BE(header + (kind === 5 ? 12 : 8) + 2 * index); + const rowIdOffset = kind === 5 ? cell + 4 : varint(source.bytes, cell).next; + const rowId = varint(source.bytes, rowIdOffset).value; + if (kind === 13 && rowId === BigInt(seq)) return `${source.file} byte offset ${cell}`; + if (kind === 5 && BigInt(seq) <= rowId) { + next = source.bytes.readUInt32BE(cell); + break; + } + } + if (next === 0) break; + page = next; + } + throw new Error('Cannot locate journal record cell.'); +} + +function varint(bytes: Buffer, offset: number): { value: bigint; next: number } { + let value = 0n; + for (let index = 0; index < 9; index += 1) { + const byte = bytes.readUInt8(offset++); + value = (value << (index === 8 ? 8n : 7n)) | BigInt(index === 8 ? byte : byte & 0x7f); + if (byte < 0x80 || index === 8) return { value, next: offset }; + } + throw new Error('Invalid SQLite varint.'); +} diff --git a/packages/sdk/src/journal-reader.ts b/packages/sdk/src/journal-reader.ts new file mode 100644 index 000000000..39a702de7 --- /dev/null +++ b/packages/sdk/src/journal-reader.ts @@ -0,0 +1,165 @@ +import { constants } from 'node:fs'; +import { copyFile, mkdtemp, open, rm, stat } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { createRequire } from 'node:module'; +import { join, resolve } from 'node:path'; +import { journalRecordOffset } from './journal-offset.js'; + +// Journal version 1's closed vocabulary (relayflowd-core/src/entry.rs). +const ENTRY_TYPES = new Set([ + 'run.spawned', 'run.cancel.requested', 'event.received', 'subscription.registered', + 'subscription.matched', 'subscription.stale', 'step.routed', 'step.attempt.started', + 'step.completed', 'wait.event', 'wait.human', 'sleep.until', 'wait.completed', + 'stream.appended', 'memory.injected', 'channel.appended', 'channel.delivered', + 'channel.acknowledged', 'effect.recorded', 'effect.confirmed', 'epoch.summary', + 'segment.closed', 'run.completed', +]); + +/** The persisted envelope from relayflowd-core/src/entry.rs. */ +export interface JournalEvent { + seq: number; + segment_id: number; + entry_type: string; + run_id: string; + step_id: string | null; + attempt: number | null; + at_ms: number; + payload: unknown; +} + +export type JournalReadFailure = + | 'invalid_run_id' | 'run_not_found' | 'step_not_found' | 'journal_read_failed'; + +export class JournalReadError extends Error { + constructor(readonly code: JournalReadFailure, message: string) { + super(message); + this.name = 'JournalReadError'; + } +} + +async function fingerprint(path: string): Promise { + try { + const info = await stat(path, { bigint: true }); + if (!info.isFile()) throw new Error(`Journal path is not a regular file: ${path}`); + return [info.dev, info.ino, info.size, info.mtimeNs, info.ctimeNs].join(':'); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return undefined; + throw error; + } +} + +/** + * Walk a stable on-disk journal in sequence order, across every segment. + * `at` includes the step's terminal completion, or its last journaled event + * if it has not terminated. Retries are included. + * + * SQLite read-only connections can still create/change WAL shared-memory files. + * Copy the database and its WAL into private scratch space before opening it, + * and reject concurrent source changes instead of returning a torn snapshot. + * Only scratch files are written; the run, registry and daemon are untouched. + */ +export async function* walkJournal( + runId: string, + dataDir: string, + options: { at?: string } = {}, +): AsyncIterable { + // Run ids are path components, never paths. Permit imported run names as + // well as the ULIDs produced by the kernel. + if (!/^[A-Za-z0-9][A-Za-z0-9_-]*$/.test(runId)) { + throw new JournalReadError('invalid_run_id', 'Run id must contain only letters, digits, underscores or hyphens, starting with a letter or digit.'); + } + const source = join(resolve(dataDir), 'runs', `${runId}.sqlite3`); + let scratch: string | undefined; + let database: import('node:sqlite').DatabaseSync | undefined; + try { + const before = await Promise.all([fingerprint(source), fingerprint(`${source}-wal`)]); + if (before[0] === undefined) { + throw new JournalReadError('run_not_found', `Run "${runId}" does not exist in "${dataDir}".`); + } + scratch = await mkdtemp(join(tmpdir(), 'flows-replay-')); + const snapshot = join(scratch, 'journal.sqlite3'); + await copyFile(source, snapshot, constants.COPYFILE_EXCL); + if (before[1] !== undefined) await copyFile(`${source}-wal`, `${snapshot}-wal`, constants.COPYFILE_EXCL); + const after = await Promise.all([fingerprint(source), fingerprint(`${source}-wal`)]); + if (before.some((value, index) => value !== after[index])) { + throw new Error('Journal changed while taking the replay snapshot; retry replay.'); + } + const file = await open(snapshot, 'r'); + try { + const header = Buffer.alloc(16); + const { bytesRead } = await file.read(header, 0, 16, 0); + const expected = Buffer.from('SQLite format 3\0'); + for (let offset = 0; offset < expected.length; offset += 1) { + if (offset >= bytesRead || header[offset] !== expected[offset]) { + throw new Error(`Invalid SQLite header at journal byte offset ${offset}.`); + } + } + } finally { await file.close(); } + + // Lazy loading keeps all existing CLI verbs independent of node:sqlite. + const { DatabaseSync } = createRequire(import.meta.url)('node:sqlite') as typeof import('node:sqlite'); + database = new DatabaseSync(snapshot, { readOnly: true }); + const metadata = database.prepare("SELECT value FROM meta WHERE key = 'run_id'").get(); + if (metadata?.['value'] !== runId) throw new Error('Journal run id does not match the requested run.'); + const version = database.prepare("SELECT value FROM meta WHERE key = 'journal_version'").get(); + if (version?.['value'] !== '1') throw new Error('Unsupported journal version.'); + const integrity = database.prepare('PRAGMA quick_check').all(); + if (integrity.length !== 1 || integrity[0]?.['quick_check'] !== 'ok') throw new Error('Journal integrity check failed.'); + + let through: number | undefined; + if (options.at !== undefined) { + const row = database.prepare(` + SELECT COALESCE(MAX(CASE WHEN entry_type = 'step.completed' + AND json_valid(payload) AND json_extract(payload, '$.disposition') IN ('step_done', 'park') + THEN seq END), MAX(seq)) AS seq FROM entries WHERE step_id = ? + `).get(options.at); + if (row?.['seq'] === null || row === undefined) { + throw new JournalReadError('step_not_found', `Step "${options.at}" is not journaled in run "${runId}".`); + } + through = integer(row['seq'], 'seq'); + } + const rows = database.prepare( + 'SELECT seq, segment_id, entry_type, step_id, attempt, at_ms, payload FROM entries' + + (through === undefined ? '' : ' WHERE seq <= ?') + ' ORDER BY seq', + ).iterate(...(through === undefined ? [] : [through])); + for (const row of rows) { + try { + if (typeof row['entry_type'] !== 'string' || !ENTRY_TYPES.has(row['entry_type']) + || typeof row['payload'] !== 'string' + || (row['step_id'] !== null && typeof row['step_id'] !== 'string')) { + throw new Error('Invalid journal entry or unknown record type.'); + } + yield { + seq: integer(row['seq'], 'seq'), + segment_id: integer(row['segment_id'], 'segment_id'), + entry_type: row['entry_type'], + run_id: runId, + step_id: row['step_id'], + attempt: row['attempt'] === null ? null : integer(row['attempt'], 'attempt'), + at_ms: integer(row['at_ms'], 'at_ms'), + payload: JSON.parse(row['payload']), + }; + } catch (error) { + const root = database.prepare("SELECT rootpage FROM sqlite_schema WHERE name = 'entries'").get(); + const location = journalRecordOffset(snapshot, integer(root?.['rootpage'], 'rootpage'), integer(row['seq'], 'seq')); + throw new Error(`Entry seq ${row['seq']} at ${location}: ${error instanceof Error ? error.message : String(error)}`); + } + } + } catch (error) { + if (error instanceof JournalReadError) throw error; + throw new JournalReadError('journal_read_failed', `Cannot read journal for run "${runId}": ${error instanceof Error ? error.message : String(error)}`); + } finally { + try { + database?.close(); + } finally { + if (scratch !== undefined) await rm(scratch, { recursive: true, force: true }); + } + } +} + +function integer(value: unknown, field: string): number { + if (typeof value !== 'number' || !Number.isSafeInteger(value)) { + throw new Error(`Invalid journal ${field}: expected a safe integer.`); + } + return value; +} diff --git a/packages/sdk/tests/cli-replay.test.ts b/packages/sdk/tests/cli-replay.test.ts new file mode 100644 index 000000000..8faffb477 --- /dev/null +++ b/packages/sdk/tests/cli-replay.test.ts @@ -0,0 +1,276 @@ +import { spawnSync } from 'node:child_process'; +import { createHash } from 'node:crypto'; +import * as fsPromises from 'node:fs/promises'; +import { mkdirSync, mkdtempSync, readFileSync, readdirSync, rmSync, statSync, writeFileSync } from 'node:fs'; +import { Socket } from 'node:net'; +import { createRequire } from 'node:module'; +import { tmpdir } from 'node:os'; +import { dirname, join, resolve } from 'node:path'; +import type { DatabaseSync as Database } from 'node:sqlite'; +import { fileURLToPath } from 'node:url'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { runCli } from '../src/cli.js'; +import { walkJournal, type JournalEvent } from '../src/journal-client.js'; + +vi.mock('node:fs/promises', async (importOriginal) => ({ + ...await importOriginal(), +})); + +const ROOT = resolve(dirname(fileURLToPath(import.meta.url)), '../../..'); +const { DatabaseSync } = createRequire(import.meta.url)('node:sqlite') as typeof import('node:sqlite'); +const CLI = join(ROOT, 'packages/sdk/dist/cli.js'); +// Captured from a real completed kernel run; replay must not execute its +// recorded agent command or its 35-second deterministic command. +const EVENTS = readFileSync(join(ROOT, 'docs/evidence/journal-close-0909/completed.journal.jsonl'), 'utf8') + .trim().split('\n').map((line) => JSON.parse(line) as JournalEvent); +const RUN_ID = EVENTS[0]!.run_id; +const directories: string[] = []; + +afterEach(() => { + vi.restoreAllMocks(); + for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true }); +}); + +function temporaryDirectory(): string { + const directory = mkdtempSync(join(tmpdir(), 'cli-replay-')); + directories.push(directory); + return directory; +} + +function fixture(events = EVENTS, wal = false): { dataDir: string; path: string; writer: Database } { + const dataDir = temporaryDirectory(); + mkdirSync(join(dataDir, 'runs')); + const path = join(dataDir, 'runs', `${RUN_ID}.sqlite3`); + const writer = new DatabaseSync(path); + if (wal) writer.exec('PRAGMA journal_mode = WAL'); + // Persist the real envelope using the schema in relayflowd-journal/src/lib.rs. + writer.exec(` + CREATE TABLE meta (key TEXT PRIMARY KEY, value TEXT NOT NULL) WITHOUT ROWID; + INSERT INTO meta VALUES ('run_id', '${RUN_ID}'), ('journal_version', '1'); + CREATE TABLE segments (segment_id INTEGER PRIMARY KEY, journal_version INTEGER NOT NULL, opened_seq INTEGER NOT NULL); + INSERT INTO segments VALUES (1, 1, 1); + CREATE TABLE entries ( + seq INTEGER PRIMARY KEY, segment_id INTEGER NOT NULL REFERENCES segments(segment_id), + entry_type TEXT NOT NULL, step_id TEXT, attempt INTEGER, at_ms INTEGER NOT NULL, payload TEXT NOT NULL + ); + `); + const insert = writer.prepare('INSERT INTO entries VALUES (?, ?, ?, ?, ?, ?, ?)'); + for (const event of events) { + insert.run(event.seq, event.segment_id, event.entry_type, event.step_id, event.attempt, event.at_ms, JSON.stringify(event.payload)); + } + return { dataDir, path, writer }; +} + +async function replay(dataDir: string, extra: string[] = [], runId = RUN_ID) { + const stdout: string[] = []; + const stderr: string[] = []; + const code = await runCli(['replay', '--data-dir', dataDir, runId, ...extra], { + stdout: (line) => stdout.push(line), stderr: (line) => stderr.push(line), + }); + return { code, stdout, stderr }; +} + +function diskState(directory: string): unknown { + return readdirSync(directory).sort().map((name) => { + const path = join(directory, name); + const info = statSync(path); + return [name, info.mtimeMs, info.isDirectory() ? diskState(path) : createHash('sha256').update(readFileSync(path)).digest('hex')]; + }); +} + +describe('flows replay', () => { + it('full walk emits every journal event in order through the terminal event without a daemon or writes', async () => { + const { dataDir, writer } = fixture(EVENTS, true); + writer.close(); + const before = diskState(dataDir); + const connect = vi.spyOn(Socket.prototype, 'connect').mockImplementation(() => { throw new Error('replay opened a socket'); }); + const json = await replay(dataDir, ['--json']); + expect(json.code).toBe(0); + expect(json.stderr).toEqual([]); + expect(json.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS); + expect(JSON.parse(json.stdout.at(-1)!).kind).toBe('run.completed'); + expect(JSON.parse(json.stdout[0]!)).toMatchObject({ step_id: null, kind: 'run.spawned', verification: null, spend: null }); + const completedPayload = EVENTS[3]!.payload as Record; + expect(JSON.parse(json.stdout[3]!)).toEqual({ + step_id: 'implement', kind: 'step.completed', event: EVENTS[3], + verification: completedPayload['verification'], spend: completedPayload['budget'], + }); + const text = await replay(dataDir); + expect(text.code).toBe(0); + expect(text.stderr).toEqual([]); + expect(text.stdout).toHaveLength(EVENTS.length); + EVENTS.forEach((event, index) => expect(text.stdout[index]).toContain(`${event.seq} ${event.at_ms} ${event.entry_type}`)); + expect(text.stdout.every((line) => !line.includes('\n'))).toBe(true); + expect(connect).not.toHaveBeenCalled(); + expect(diskState(dataDir)).toEqual(before); + }); + + it('--at emits the full-walk prefix inclusive of the requested step', async () => { + const { dataDir, writer } = fixture(); + writer.close(); + for (const flags of [[], ['--json']]) { + const full = await replay(dataDir, flags); + const prefix = await replay(dataDir, [...flags, '--at', 'verify']); + expect(prefix.code).toBe(0); + expect(prefix.stderr).toEqual([]); + expect(prefix.stdout).toEqual(full.stdout.slice(0, 7)); + expect(prefix.stdout.at(-1)).toContain('step.completed'); + } + }); + + it('--json is byte-identical across two CLI invocations (diff)', () => { + const { dataDir, writer } = fixture(); + writer.close(); + const outputs = [0, 1].map((index) => { + const result = spawnSync(process.execPath, [CLI, 'replay', '--json', '--data-dir', dataDir, RUN_ID], { encoding: 'utf8' }); + expect(result.status, result.stderr).toBe(0); + expect(result.stdout.trim().split('\n').map((line) => JSON.parse(line).event)).toEqual(EVENTS); + const path = join(dataDir, `replay-${index}.jsonl`); + writeFileSync(path, result.stdout); + return path; + }); + const diff = spawnSync('diff', ['-u', ...outputs], { encoding: 'utf8' }); + expect(diff.status, diff.stderr).toBe(0); + expect(diff.stdout).toBe(''); + }); + + it('refuses run_not_found without creating a data directory', async () => { + const dataDir = join(temporaryDirectory(), 'absent'); + const output = await replay(dataDir, ['--json']); + expect(output.code).toBe(2); + expect(output.stdout).toEqual([]); + expect(output.stderr[0]).toContain('REFUSED [run_not_found]'); + expect(readdirSync(dirname(dataDir))).toEqual([]); + }); + + it('refuses step_not_found before emitting journal events', async () => { + const { dataDir, writer } = fixture(); + writer.close(); + const output = await replay(dataDir, ['--at', 'absent']); + expect(output.code).toBe(2); + expect(output.stdout).toEqual([]); + expect(output.stderr[0]).toContain('REFUSED [step_not_found]'); + }); + + it.each(['../escape', '/absolute', 'a/b', 'a\\b', '.', '..', 'bad\nname', ''])('refuses invalid_run_id: %j', async (runId) => { + const output = await replay(temporaryDirectory(), [], runId); + expect(output.code).toBe(2); + expect(output.stdout).toEqual([]); + expect(output.stderr[0]).toContain('REFUSED [invalid_run_id]'); + }); + + it('refuses journal_read_failed for a corrupted journal', async () => { + const { dataDir, path, writer } = fixture(); + writer.close(); + const bytes = readFileSync(path); + bytes[0] = 0xff; + writeFileSync(path, bytes); + const output = await replay(dataDir, ['--json']); + expect(output.code).toBe(2); + expect(output.stdout).toEqual([]); + expect(output.stderr[0]).toContain('REFUSED [journal_read_failed]'); + expect(output.stderr[0]).toContain('byte offset 0'); + expect(readFileSync(path)).toEqual(bytes); + }); + + it('reads committed WAL events without changing the live database or sidecars', async () => { + const { dataDir, writer } = fixture(EVENTS, true); + try { + const before = diskState(dataDir); + const output = await replay(dataDir, ['--json']); + expect(output.code).toBe(0); + expect(output.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS); + expect(diskState(dataDir)).toEqual(before); + } finally { writer.close(); } + }); + + it('walks all segments and includes all retries when stopping at a step', async () => { + const { dataDir, writer } = fixture(); + writer.exec("INSERT INTO segments VALUES (2, 1, 8); UPDATE entries SET segment_id = 2 WHERE seq >= 8; UPDATE entries SET step_id = 'verify', attempt = 2 WHERE seq BETWEEN 8 AND 10; UPDATE entries SET payload = json_set(payload, '$.disposition', 'retry') WHERE seq = 7"); + writer.close(); + const output = await replay(dataDir, ['--json', '--at', 'verify']); + expect(output.code).toBe(0); + const events = output.stdout.map((line) => JSON.parse(line).event); + expect(events).toHaveLength(10); + expect(events.at(-1)).toMatchObject({ seq: 10, segment_id: 2, step_id: 'verify', attempt: 2 }); + }); + + it('stops at the terminal completion even if a later channel record references the step', async () => { + const { dataDir, writer } = fixture(); + writer.exec("UPDATE entries SET step_id = 'verify', entry_type = 'channel.acknowledged' WHERE seq = 10"); + writer.close(); + const output = await replay(dataDir, ['--json', '--at', 'verify']); + expect(output.code).toBe(0); + expect(output.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS.slice(0, 7)); + }); + + it.each([false, true])('returns exit 1 and a byte offset after a mid-journal parse error (WAL=%s)', async (wal) => { + const { dataDir, path, writer } = fixture(EVENTS, wal); + writer.exec("UPDATE entries SET payload = '{' WHERE seq = 7"); + if (!wal) writer.close(); + try { + const output = await replay(dataDir, ['--json']); + expect(output.code).toBe(1); + expect(output.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS.slice(0, 6)); + expect(output.stderr[0]).toContain('FAILED [journal_read_failed]'); + expect(output.stderr[0]).toContain('Entry seq 7'); + const location = output.stderr[0]!.match(/(journal|WAL) byte offset (\d+)/); + expect(location).not.toBeNull(); + const bytes = readFileSync(location![1] === 'WAL' ? `${path}-wal` : path); + expect(Number(location![2])).toBeLessThan(bytes.length); + // The cell contains this record's kind and invalid JSON payload. + expect(bytes.subarray(Number(location![2]), Number(location![2]) + 100).toString()).toContain('step.completed'); + } finally { if (wal) writer.close(); } + }); + + it('refuses a journal that changes while the snapshot is copied', async () => { + const { dataDir, path, writer } = fixture(); + writer.close(); + const originalCopy = fsPromises.copyFile; + vi.spyOn(fsPromises, 'copyFile').mockImplementation(async (source, destination, mode) => { + await originalCopy(source, destination, mode); + writeFileSync(path, Buffer.concat([readFileSync(path), Buffer.from('changed')])); + }); + const output = await replay(dataDir); + expect(output.code).toBe(2); + expect(output.stdout).toEqual([]); + expect(output.stderr[0]).toContain('Journal changed while taking the replay snapshot'); + }); + + it('can stop at an unfinished step and release the walker early', async () => { + const { dataDir, writer } = fixture(EVENTS.slice(0, 6)); + writer.close(); + const output = await replay(dataDir, ['--json', '--at', 'verify']); + expect(output.code).toBe(0); + expect(output.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS.slice(0, 6)); + for await (const event of walkJournal(RUN_ID, dataDir)) { + expect(event).toEqual(EVENTS[0]); + break; + } + }); + + it.each([ + "UPDATE meta SET value = 'wrong-run' WHERE key = 'run_id'", + "UPDATE meta SET value = '999' WHERE key = 'journal_version'", + "UPDATE entries SET payload = '{' WHERE seq = 1", + "UPDATE entries SET at_ms = 9223372036854775807 WHERE seq = 1", + "UPDATE entries SET entry_type = 'unknown.record' WHERE seq = 1", + ])('fails closed on invalid journal contents: %s', async (sql) => { + const { dataDir, writer } = fixture(); + writer.exec(sql); + writer.close(); + const output = await replay(dataDir); + expect(output.code).toBe(2); + expect(output.stderr[0]).toContain('REFUSED [journal_read_failed]'); + }); + + it.each([ + [], [RUN_ID, '--at'], [RUN_ID, '--data-dir'], [RUN_ID, '--no-spawn'], + [RUN_ID, '--json', '--json'], [RUN_ID, '--at', 'x', '--at', 'y'], + [RUN_ID, '--data-dir', 'x', '--data-dir', 'y'], [RUN_ID, 'extra'], + ].map((args) => ({ args })))('refuses invalid invocation $args', async ({ args }) => { + const stderr: string[] = []; + expect(await runCli(['replay', ...args], { stdout: () => {}, stderr: (line) => stderr.push(line) })).toBe(2); + expect(stderr.join('\n')).toContain('REFUSED [invalid_invocation]'); + }); +}); From e2463c3b0401c9a54c3dc54aa95996e23502683d Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 11 Sep 2026 10:34:11 +0200 Subject: [PATCH 2/2] fix(sdk): preserve parked replay waits and current WAL offsets Session-Id: 01a08f80-0bc7-7e03-91e1-fb1f112a0616 --- packages/sdk/src/journal-offset.ts | 3 +++ packages/sdk/src/journal-reader.ts | 2 +- packages/sdk/tests/cli-replay.test.ts | 30 +++++++++++++++++++++++++++ 3 files changed, 34 insertions(+), 1 deletion(-) diff --git a/packages/sdk/src/journal-offset.ts b/packages/sdk/src/journal-offset.ts index 515c4f970..5e1dc7d1f 100644 --- a/packages/sdk/src/journal-offset.ts +++ b/packages/sdk/src/journal-offset.ts @@ -14,6 +14,9 @@ export function journalRecordOffset(path: string, rootPage: number, seq: number) const frameSize = pageSize + 24; let lastCommit = 0; for (let frame = 32; frame + frameSize <= wal.length; frame += frameSize) { + // RESTART checkpoints reuse the WAL without truncating its old tail. + // Frames from that prior generation are not part of SQLite's view. + if (!wal.subarray(frame + 8, frame + 16).equals(wal.subarray(16, 24))) break; if (wal.readUInt32BE(frame + 4) !== 0) lastCommit = frame; } for (let frame = 32; frame <= lastCommit; frame += frameSize) { diff --git a/packages/sdk/src/journal-reader.ts b/packages/sdk/src/journal-reader.ts index 39a702de7..884230660 100644 --- a/packages/sdk/src/journal-reader.ts +++ b/packages/sdk/src/journal-reader.ts @@ -110,7 +110,7 @@ export async function* walkJournal( if (options.at !== undefined) { const row = database.prepare(` SELECT COALESCE(MAX(CASE WHEN entry_type = 'step.completed' - AND json_valid(payload) AND json_extract(payload, '$.disposition') IN ('step_done', 'park') + AND json_valid(payload) AND json_extract(payload, '$.disposition') = 'step_done' THEN seq END), MAX(seq)) AS seq FROM entries WHERE step_id = ? `).get(options.at); if (row?.['seq'] === null || row === undefined) { diff --git a/packages/sdk/tests/cli-replay.test.ts b/packages/sdk/tests/cli-replay.test.ts index 8faffb477..d74ebbb13 100644 --- a/packages/sdk/tests/cli-replay.test.ts +++ b/packages/sdk/tests/cli-replay.test.ts @@ -204,6 +204,36 @@ describe('flows replay', () => { expect(output.stdout.map((line) => JSON.parse(line).event)).toEqual(EVENTS.slice(0, 7)); }); + it('--at includes the human-wait record after a parked completion', async () => { + const { dataDir, writer } = fixture(EVENTS.slice(0, 8)); + writer.exec("UPDATE entries SET payload = json_set(payload, '$.disposition', 'park', '$.completionReason', 'needs_human') WHERE seq = 7"); + writer.prepare("UPDATE entries SET step_id = 'verify', entry_type = 'wait.human', payload = ? WHERE seq = 8") + .run(JSON.stringify({ prompt: 'Review the workspace changes', diff_ref: 'diff-1' })); + writer.close(); + const output = await replay(dataDir, ['--json', '--at', 'verify']); + expect(output.code).toBe(0); + expect(output.stdout).toHaveLength(8); + expect(JSON.parse(output.stdout.at(-1)!)).toMatchObject({ + kind: 'wait.human', step_id: 'verify', + event: { payload: { prompt: 'Review the workspace changes', diff_ref: 'diff-1' } }, + }); + }); + + it('locates corruption in the current WAL generation after a checkpoint', async () => { + const { dataDir, path, writer } = fixture(EVENTS, true); + try { + writer.exec('PRAGMA wal_checkpoint(RESTART)'); + writer.exec("UPDATE entries SET payload = 'invalid-after-checkpoint' WHERE seq = 7"); + const output = await replay(dataDir, ['--json']); + expect(output.code).toBe(1); + expect(output.stdout).toHaveLength(6); + const location = output.stderr[0]!.match(/WAL byte offset (\d+)/); + expect(location, output.stderr[0]).not.toBeNull(); + const bytes = readFileSync(`${path}-wal`); + expect(bytes.subarray(Number(location![1]), Number(location![1]) + 100).toString()).toContain('invalid-after-checkpoint'); + } finally { writer.close(); } + }); + it.each([false, true])('returns exit 1 and a byte offset after a mid-journal parse error (WAL=%s)', async (wal) => { const { dataDir, path, writer } = fixture(EVENTS, wal); writer.exec("UPDATE entries SET payload = '{' WHERE seq = 7");