diff --git a/packages/cli/src/commands/sync.ts b/packages/cli/src/commands/sync.ts index 8ec5429..82f0f5f 100644 --- a/packages/cli/src/commands/sync.ts +++ b/packages/cli/src/commands/sync.ts @@ -12,7 +12,7 @@ import { writeFile, } from 'node:fs/promises'; import { dirname, join, relative, sep } from 'node:path'; -import { ensureRunsDir, resolveRunSelector, runDirFor } from '@orcareplay/core'; +import { ensureRunsDir, resolveRunSelector, runDirFor, sha256File } from '@orcareplay/core'; import type { ParsedArgs } from '../args.js'; import type { Output } from '../out.js'; import { @@ -953,6 +953,117 @@ export async function swapStagedRun( * * `orca pull [--gateway URL] [--force]` */ + +/** + * Spec §6: "Readers SHOULD verify it and MUST report, not repair, a mismatch." + * + * Nothing did. `verifyIntegrity` existed and `replay` called it; `pull` wrote a run to disk and + * said `pull.done`. A truncated download, a proxy that re-encoded the stream, an archive that lost + * a byte — every one of them landed silently, and the first sign would be a replay that diverged + * for no reason anyone could name. + * + * Checked on the STAGED copy, before the swap: a corrupt transfer must not be what replaces a local + * run that `--force` merely asked to update. + */ +async function verifyStaged(staging: string, runId: string): Promise { + const manifestPath = join(staging, 'manifest.json'); + const raw = await readFile(manifestPath, 'utf8').catch(() => undefined); + if (raw === undefined) return undefined; + let expected: unknown; + try { + expected = (JSON.parse(raw) as { integrity?: { events_sha256?: unknown } }).integrity + ?.events_sha256; + } catch { + throw new Error(`${runId}: the manifest that arrived is not JSON`); + } + // An unsealed run — one whose recorder crashed — has no root to check. Spec §6 describes the + // root as written "at the moment the run ended", so its absence is a state, not a failure. + if (typeof expected !== 'string') return undefined; + const actual = await sha256File(join(staging, 'events.jsonl')); + if (actual !== expected) { + throw new Error( + `${runId}: events.jsonl does not match the integrity root in its own manifest ` + + `(manifest ${expected.slice(0, 16)}…, file ${actual.slice(0, 16)}…). Nothing was ` + + 'replaced. Pull it again; if it persists the copy on the gateway is damaged.', + ); + } + return actual; +} + +/** Every event, in order, with keys sorted deeply — so two serialisations of one run compare equal. */ +async function canonicalEvents(dir: string): Promise { + const raw = await readFile(join(dir, 'events.jsonl'), 'utf8').catch(() => undefined); + if (raw === undefined) return undefined; + const sortDeep = (v: unknown): unknown => { + if (Array.isArray(v)) return v.map(sortDeep); + if (v === null || typeof v !== 'object') return v; + const out: Record = {}; + for (const k of Object.keys(v as Record).sort()) + out[k] = sortDeep((v as Record)[k]); + return out; + }; + try { + return raw + .split('\n') + .filter((l) => l.trim().length > 0) + .map((l) => JSON.stringify(sortDeep(JSON.parse(l)))) + .join('\n'); + } catch { + return undefined; + } +} + +/** + * Whether the copy that arrived is the copy that was pushed. + * + * A gateway is free to store the trace however it likes, and this one re-serialises it: Go's + * `encoding/json` sorts the keys of a map and escapes `<`, `>` and `&`, so the bytes come back + * different and the root in the manifest is the root of THOSE bytes. Both copies are internally + * consistent, and neither is wrong — but the digest no longer tells you the two are the same run, + * which is the one question a digest is for. + * + * So say which it is. Identical content under a different serialisation is worth a line; content + * that actually differs is worth a loud one. + */ +async function reportAgainstLocal( + out: Output, + runId: string, + localDir: string, + stagedDir: string, + stagedRoot: string | undefined, +): Promise { + const localRaw = await readFile(join(localDir, 'manifest.json'), 'utf8').catch(() => undefined); + if (localRaw === undefined || stagedRoot === undefined) return; + let localRoot: unknown; + try { + localRoot = (JSON.parse(localRaw) as { integrity?: { events_sha256?: unknown } }).integrity + ?.events_sha256; + } catch { + return; + } + if (typeof localRoot !== 'string' || localRoot === stagedRoot) return; + + const [a, b] = await Promise.all([canonicalEvents(localDir), canonicalEvents(stagedDir)]); + const fields = { + run: runId, + local: `${localRoot.slice(0, 16)}…`, + pulled: `${stagedRoot.slice(0, 16)}…`, + }; + if (a !== undefined && b !== undefined && a === b) { + out.warn('pull.reserialized', { + ...fields, + why: 'same events, different bytes — the gateway stores its own serialisation and recomputes the root over it', + effect: 'the two copies can no longer be compared by digest, though nothing was lost', + }); + return; + } + out.warn('pull.differs', { + ...fields, + why: 'the events themselves differ, not just how they were written', + next: `orca show ${runId} # the copy that just replaced the local one`, + }); +} + export async function pullCommand( args: ParsedArgs, out: Output, @@ -1080,6 +1191,10 @@ export async function pullCommand( try { await stageRunEntries(entries, runId, staging, stillHeld); + // Before the swap, never after: what fails verification must not be what replaced the run. + const stagedRoot = await verifyStaged(staging, runId); + if (existing) await reportAgainstLocal(out, runId, dest, staging, stagedRoot); + movedAside = await swapStagedRun({ dest, staging, retired, existing: !!existing }, stillHeld); } catch (err) { await revertStagedSwap(stillHeld, { dest, staging, retired, existing: movedAside }); diff --git a/packages/cli/test/sync.test.ts b/packages/cli/test/sync.test.ts index ee380fe..f087b5b 100644 --- a/packages/cli/test/sync.test.ts +++ b/packages/cli/test/sync.test.ts @@ -1,3 +1,4 @@ +import { createHash } from 'node:crypto'; import { createServer, type Server } from 'node:http'; import { mkdtemp, mkdir, readFile, readdir, rm, stat, utimes, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; @@ -192,6 +193,102 @@ describe('push and pull', () => { expect(Array.from(await readFile(join(dir, 'blobs', 'ab', 'abcdef')))).toEqual([9, 8, 7]); }); + /** + * SPEC §6: "Readers SHOULD verify it and MUST report, not repair, a mismatch." + * + * `verifyIntegrity` existed and `replay` called it. `pull` wrote a run to disk and said + * `pull.done` — so a download that arrived altered landed silently, and the first sign of it + * would be a replay that diverged for a reason nobody could name. + * + * The zip's own CRC catches a flipped bit. It cannot catch an archive that is internally + * consistent and whose events simply are not the ones the manifest attests to, which is the + * shape a storage bug on the far side takes. + */ + it('refuses a pulled run whose events do not match its own integrity root', async () => { + const events = '{"seq":0,"type":"run.start"}\n'; + const zip = await writeArchive([ + { + name: `${runId}/manifest.json`, + bytes: new TextEncoder().encode( + `{"run_id":"${runId}","integrity":{"events_sha256":"${'0'.repeat(64)}"}}\n`, + ), + }, + { name: `${runId}/events.jsonl`, bytes: new TextEncoder().encode(events) }, + ]); + reply = { status: 200, body: Buffer.from(zip), headers: { 'content-type': 'application/zip' } }; + + await expect(pullCommand(parseArgs(['pull', runId]), out, workspace, env())).rejects.toThrow( + /does not match the integrity root/, + ); + // Verified on the staged copy, before the swap: what fails must not be what replaced a run. + await expect( + readFile(join(workspace, '.orca', 'runs', runId, 'events.jsonl')), + ).rejects.toThrow(); + }); + + /** + * The same run, pushed and pulled back, does not come back as the same bytes: a gateway is free + * to store the trace however it likes, and this one re-serialises it — Go's `encoding/json` + * sorts a map's keys and escapes `<`, `>` and `&` — then recomputes the root over its own bytes. + * + * Both copies are internally consistent and neither is wrong. What is lost is the one thing a + * digest is for: it no longer tells you the two are the same run. So say which it is, rather + * than leaving someone to discover two roots and not know whether to worry. + */ + it('says when the copy that arrived is the same events under a different serialisation', async () => { + await seedRun(); + const local = join(workspace, '.orca', 'runs', runId); + // `seedRun` writes no integrity root, and without one on both sides there is nothing to + // compare — which is the case this test exists to cover, so give the local copy a real one. + // Written in the order a writer writes them, which is NOT alphabetical — `seq,ts,type` already + // is, so sorting it would be a no-op and the test would pass by testing nothing. + const localEvents = '{"type":"run.start","seq":0,"actor":"orca"}\n'; + await writeFile(join(local, 'events.jsonl'), localEvents); + await writeFile( + join(local, 'manifest.json'), + `${JSON.stringify({ + run_id: runId, + schema_version: '0.1.0', + integrity: { + events_sha256: createHash('sha256').update(localEvents).digest('hex'), + blob_count: 0, + }, + })}\n`, + ); + // What the gateway stores: the same events, keys sorted, root recomputed over those bytes. + const resorted = localEvents + .split('\n') + .filter((l) => l.trim().length > 0) + .map((l) => { + const o = JSON.parse(l) as Record; + const sorted: Record = {}; + for (const k of Object.keys(o).sort()) sorted[k] = o[k]; + return JSON.stringify(sorted); + }) + .join('\n'); + const bytes = new TextEncoder().encode(`${resorted}\n`); + const root = createHash('sha256').update(bytes).digest('hex'); + const zip = await writeArchive([ + { + name: `${runId}/manifest.json`, + bytes: new TextEncoder().encode( + `{"run_id":"${runId}","integrity":{"events_sha256":"${root}"}}\n`, + ), + }, + { name: `${runId}/events.jsonl`, bytes }, + ]); + reply = { status: 200, body: Buffer.from(zip), headers: { 'content-type': 'application/zip' } }; + + await pullCommand(parseArgs(['pull', runId, '--force']), out, workspace, env()); + + const said = logs.find((e) => e.event === 'pull.reserialized'); + expect( + said, + 'a re-serialised copy should be reported, not passed off as identical', + ).toBeDefined(); + expect(logs.some((e) => e.event === 'pull.differs')).toBe(false); + }); + /** * A truncated recording is one the gateway could not store whole. Pulling it is allowed — half a * trace still debugs — but saying so is not optional: every later command reads the result as if