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
117 changes: 116 additions & 1 deletion packages/cli/src/commands/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -953,6 +953,117 @@ export async function swapStagedRun(
*
* `orca pull <run> [--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<string | undefined> {
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<string | undefined> {
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<string, unknown> = {};
for (const k of Object.keys(v as Record<string, unknown>).sort())
out[k] = sortDeep((v as Record<string, unknown>)[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<void> {
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,
Expand Down Expand Up @@ -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 });
Expand Down
97 changes: 97 additions & 0 deletions packages/cli/test/sync.test.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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<string, unknown>;
const sorted: Record<string, unknown> = {};
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
Expand Down
Loading