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..5e1dc7d1f
--- /dev/null
+++ b/packages/sdk/src/journal-offset.ts
@@ -0,0 +1,60 @@
+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) {
+ // 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) {
+ 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..884230660
--- /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') = 'step_done'
+ 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..d74ebbb13
--- /dev/null
+++ b/packages/sdk/tests/cli-replay.test.ts
@@ -0,0 +1,306 @@
+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('--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");
+ 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]');
+ });
+});