diff --git a/src/graph/deeplake-pull.ts b/src/graph/deeplake-pull.ts index b627fd971..64c8ae644 100644 --- a/src/graph/deeplake-pull.ts +++ b/src/graph/deeplake-pull.ts @@ -35,7 +35,7 @@ import { execFileSync } from "node:child_process"; import { createHash } from "node:crypto"; -import { existsSync, mkdirSync, renameSync, writeFileSync } from "node:fs"; +import { existsSync, mkdirSync, readFileSync, renameSync, writeFileSync } from "node:fs"; import { dirname, join } from "node:path"; /** @@ -187,8 +187,9 @@ export async function pullSnapshot( } const cloudTs = parseTs(row.ts); - // Compare with local. readLastBuild returns null on missing/corrupt - // files; in that case we ALWAYS pull (no comparison possible). + // Compare with local. The snapshot is the payload authority; the sidecar + // only contributes when it reconciles with that payload. A missing, + // corrupt, stale, or inconsistent local state falls through to pull. // // Codex P1 fix: gate the comparison on local.commit_sha === head. // `.last-build.json` records the last build for ANY commit in the @@ -205,7 +206,11 @@ export async function pullSnapshot( // sibling's. Without this, after pull worktree-A would overwrite // worktree-B's metadata (or vice versa). const worktreeId = workTreeIdFor(cwd); - const local = readLastBuild(baseDir, worktreeId); + const snapshotsDir = join(baseDir, "snapshots"); + const snapshotPath = join(snapshotsDir, `${head}.json`); + const sidecar = readLastBuild(baseDir, worktreeId); + const localSnapshot = readLocalSnapshotState(snapshotPath, head, repoKey); + const local = reconcileLocalSnapshotState(sidecar, localSnapshot); if (local !== null && local.commit_sha === head) { // CodeRabbit P1: empty cloud sha (legacy rows without the column // populated) is NOT proof local is current — it's "we don't know". @@ -222,14 +227,24 @@ export async function pullSnapshot( cloudTs, }; } + } else if ( + // Preserve the existing equal-hash fast path when the snapshot file is + // absent: a matching sidecar is enough to prove there is no cloud write + // to perform. It is deliberately NOT used for local-newer decisions; + // those require a valid commit-addressed snapshot on disk. + !existsSync(snapshotPath) && + sidecar !== null && + sidecar.commit_sha === head && + cloudSha256 !== "" && + sidecar.snapshot_sha256 === cloudSha256 + ) { + return { kind: "up-to-date", commitSha: head, snapshotSha256: cloudSha256 }; } // Write payload + sidecars. The payload IS the canonical bytes // (canonicalJSON(snapshot)) — same function as writeSnapshot uses // locally — so the file we write here is byte-identical to what a // local build would have produced. - const snapshotsDir = join(baseDir, "snapshots"); - const snapshotPath = join(snapshotsDir, `${head}.json`); const worktreeRoot = join(baseDir, "worktrees", worktreeId); try { writeFileAtomic(snapshotPath, cloudPayload); @@ -305,6 +320,107 @@ function parseTs(raw: unknown): number { return 0; } +interface LocalSnapshotState { + ts: number; + commit_sha: string; + snapshot_sha256: string; + node_count: number; + edge_count: number; +} + +interface LocalSidecarState { + ts: number; + commit_sha: string | null; + snapshot_sha256: string; +} + +/** + * Read freshness from the commit-addressed snapshot, not from a sidecar in + * isolation. The sidecar is only consistency metadata: its hash must match + * the recomputed stable-field hash before its timestamp is considered. + */ +function readLocalSnapshotState( + snapshotPath: string, + head: string, + repoKey: string, +): LocalSnapshotState | null { + if (!existsSync(snapshotPath)) return null; + try { + const parsed = JSON.parse(readFileSync(snapshotPath, "utf8")) as unknown; + if (!isCompatibleLocalSnapshot(parsed, head, repoKey)) return null; + const snapshot = parsed as GraphSnapshot; + const ts = parseTs(snapshot.observation.ts); + // A future observation is not freshness evidence. Local graph files are + // not authenticated, so an impossible clock value must fail closed and + // let the cloud repair the snapshot rather than protecting it. + if (!Number.isFinite(ts) || ts <= 0 || ts > Date.now()) return null; + return { + ts, + commit_sha: head, + snapshot_sha256: computeSnapshotSha256(snapshot), + node_count: snapshot.nodes.length, + edge_count: snapshot.links.length, + }; + } catch { + return null; + } +} + +/** + * A sidecar may be missing, stale, or left over from a different snapshot. + * Reconcile it with the actual payload before using it for freshness. A + * matching sidecar hash is consistency evidence, not authentication. + */ +function reconcileLocalSnapshotState( + sidecar: LocalSidecarState | null, + snapshot: LocalSnapshotState | null, +): LocalSnapshotState | null { + if (snapshot === null) return null; + if ( + sidecar === null || + sidecar.commit_sha !== snapshot.commit_sha || + sidecar.snapshot_sha256 !== snapshot.snapshot_sha256 || + sidecar.ts > Date.now() + ) { + return snapshot; + } + return { ...snapshot, ts: Math.max(snapshot.ts, sidecar.ts) }; +} + +/** Validate local identity and required metadata before a no-write freshness decision. */ +function isCompatibleLocalSnapshot(raw: unknown, head: string, repoKey: string): raw is GraphSnapshot { + if (raw === null || typeof raw !== "object") return false; + const snapshot = raw as Record; + if (snapshot.directed !== true || snapshot.multigraph !== true) return false; + if (!Array.isArray(snapshot.nodes) || !Array.isArray(snapshot.links)) return false; + + const graph = snapshot.graph; + if (graph === null || typeof graph !== "object") return false; + const metadata = graph as Record; + if (metadata.schema_version !== 1) return false; + if (metadata.generator !== "hivemind-graph") return false; + if (metadata.commit_sha !== head) return false; + if (metadata.repo_key !== repoKey) return false; + + // Observation is required even though the stable-field hash excludes it. + // Its fields are shape checks, not identity keys: repo_project and + // worktree_path legitimately differ across checkouts of one remote. + const observation = snapshot.observation; + if (observation === null || typeof observation !== "object") return false; + const o = observation as Record; + if (typeof o.ts !== "string") return false; + if (o.branch !== null && typeof o.branch !== "string") return false; + if (typeof o.worktree_path !== "string") return false; + if (typeof o.repo_project !== "string") return false; + if (typeof o.generator_version !== "string") return false; + if (!finiteNonNegative(o.source_files_extracted) || !finiteNonNegative(o.source_files_skipped)) return false; + return true; +} + +function finiteNonNegative(value: unknown): value is number { + return typeof value === "number" && Number.isFinite(value) && value >= 0; +} + function numOrUndefined(raw: unknown): number | undefined { if (typeof raw === "number" && Number.isFinite(raw) && raw >= 0) return raw; if (typeof raw === "string") { diff --git a/tests/shared/graph/deeplake-pull.test.ts b/tests/shared/graph/deeplake-pull.test.ts index ccd9f87fd..08e384c04 100644 --- a/tests/shared/graph/deeplake-pull.test.ts +++ b/tests/shared/graph/deeplake-pull.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; -import { mkdtempSync, rmSync, readFileSync, existsSync, mkdirSync } from "node:fs"; +import { mkdtempSync, rmSync, readFileSync, existsSync, mkdirSync, writeFileSync } from "node:fs"; import { createHash } from "node:crypto"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -246,10 +246,19 @@ describe("pullSnapshot — outcome resolution", () => { it("local ts > cloud ts → local-newer (NO overwrite)", async () => { mkdirSync(baseDir, { recursive: true }); + const localSnapshot: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + observation: { ...FIXTURE_SNAPSHOT.observation, ts: "2026-06-03T00:00:00.000Z" }, + nodes: [{ ...FIXTURE_SNAPSHOT.nodes[0]!, label: "local" }], + }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(localSnapshot)); writeLastBuild(baseDir, { - ts: 2_000_000_000_000, // year 2033 in epoch ms + ts: Date.parse("2026-06-04T00:00:00.000Z"), // valid sidecar, newer than cloud commit_sha: "head1234abcd", - snapshot_sha256: "different-local-sha", + snapshot_sha256: computeSnapshotSha256(localSnapshot), node_count: 1, edge_count: 0, }); @@ -273,11 +282,11 @@ describe("pullSnapshot — outcome resolution", () => { expect(result.kind).toBe("local-newer"); if (result.kind === "local-newer") { expect(result.commitSha).toBe("head1234abcd"); - expect(result.localTs).toBe(2_000_000_000_000); + expect(result.localTs).toBe(Date.parse("2026-06-04T00:00:00.000Z")); expect(result.cloudTs).toBeLessThan(result.localTs); } - // No file written - expect(existsSync(join(baseDir, "snapshots", "head1234abcd.json"))).toBe(false); + // No snapshot content was overwritten. + expect(readFileSync(snapshotPath, "utf8")).toBe(canonicalSnapshot(localSnapshot)); }); it("local missing → pulls (creates snapshot file + sidecars + history entry)", async () => { @@ -328,6 +337,206 @@ describe("pullSnapshot — outcome resolution", () => { expect(parsed.commit_sha).toBe("head1234abcd"); }); + it("preserves a valid newer same-HEAD snapshot when its sidecar is missing", async () => { + const localSnapshot: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + observation: { ...FIXTURE_SNAPSHOT.observation, ts: "2026-06-03T00:00:00.000Z" }, + nodes: [{ ...FIXTURE_SNAPSHOT.nodes[0]!, label: "local" }], + }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(localSnapshot)); + + const { api } = makeMockApi({ + selectReturns: [{ + snapshot_jsonb: CLOUD_PAYLOAD, + snapshot_sha256: CLOUD_PAYLOAD_SHA, + ts: "2026-06-02T00:00:00.000Z", + node_count: 1, edge_count: 0, + worktree_id: "remote-wt", + }], + }); + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, + readHead: () => "head1234abcd", + makeApi: () => api, + }); + + expect(result.kind).toBe("local-newer"); + expect(readFileSync(snapshotPath, "utf8")).toBe(canonicalSnapshot(localSnapshot)); + }); + + it("repairs a same-HEAD snapshot with foreign embedded identity instead of trusting its future observation", async () => { + const foreignSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { + ...FIXTURE_SNAPSHOT.graph, + schema_version: 999, + generator: "other-generator", + repo_key: "WRONG-REPO", + }, + observation: { + ...FIXTURE_SNAPSHOT.observation, + ts: "2099-06-03T00:00:00.000Z", + }, + } as unknown as GraphSnapshot; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(foreignSnapshot)); + + const { api } = makeMockApi({ + selectReturns: [{ + snapshot_jsonb: CLOUD_PAYLOAD, + snapshot_sha256: CLOUD_PAYLOAD_SHA, + ts: "2026-06-02T00:00:00.000Z", + node_count: 1, edge_count: 0, + worktree_id: "remote-wt", + }], + }); + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, + readHead: () => "head1234abcd", + makeApi: () => api, + }); + + expect(result.kind).toBe("pulled"); + expect(readFileSync(snapshotPath, "utf8")).toBe(CLOUD_PAYLOAD); + }); + + it("does not use a future observation as freshness evidence even with current stable identity", async () => { + const futureSnapshot: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + observation: { ...FIXTURE_SNAPSHOT.observation, ts: "2099-06-03T00:00:00.000Z" }, + nodes: [{ ...FIXTURE_SNAPSHOT.nodes[0]!, label: "future-local" }], + }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(futureSnapshot)); + + const { api } = makeMockApi({ + selectReturns: [{ + snapshot_jsonb: CLOUD_PAYLOAD, + snapshot_sha256: CLOUD_PAYLOAD_SHA, + ts: "2026-06-02T00:00:00.000Z", + node_count: 1, edge_count: 0, + worktree_id: "remote-wt", + }], + }); + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, + readHead: () => "head1234abcd", + makeApi: () => api, + }); + + expect(result.kind).toBe("pulled"); + expect(readFileSync(snapshotPath, "utf8")).toBe(CLOUD_PAYLOAD); + }); + + it("does not let a matching future sidecar block a newer cloud snapshot", async () => { + const cloud: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + }; + const local: GraphSnapshot = { + ...cloud, + observation: { ...cloud.observation, ts: "2026-06-01T00:00:00.000Z" }, + nodes: [{ ...cloud.nodes[0]!, label: "older-local" }], + }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(local)); + writeLastBuild(baseDir, { + ts: Date.now() + 60_000, + commit_sha: "head1234abcd", + snapshot_sha256: computeSnapshotSha256(local), + node_count: 1, edge_count: 0, + }, worktreeIdFromCwd(tmpCwd)); + const payload = canonicalSnapshot(cloud); + const { api } = makeMockApi({ selectReturns: [{ + snapshot_jsonb: payload, snapshot_sha256: computeSnapshotSha256(cloud), + ts: cloud.observation.ts, node_count: 1, edge_count: 0, + }] }); + + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, readHead: () => "head1234abcd", makeApi: () => api, + }); + + expect(result.kind).toBe("pulled"); + expect(readFileSync(snapshotPath, "utf8")).toBe(payload); + }); + + it.each([ + ["missing", undefined], + ["null", null], + ["incomplete", {}], + ])("repairs a %s observation despite matching snapshot and sidecar hashes", async (_label, observation) => { + const cloud: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + }; + const malformed = { ...cloud, observation }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, JSON.stringify(malformed)); + writeLastBuild(baseDir, { + ts: Date.parse(cloud.observation.ts), + commit_sha: "head1234abcd", + snapshot_sha256: computeSnapshotSha256(cloud), + node_count: 1, edge_count: 0, + }, worktreeIdFromCwd(tmpCwd)); + const payload = canonicalSnapshot(cloud); + const { api } = makeMockApi({ selectReturns: [{ + snapshot_jsonb: payload, snapshot_sha256: computeSnapshotSha256(cloud), + ts: cloud.observation.ts, node_count: 1, edge_count: 0, + }] }); + + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, readHead: () => "head1234abcd", makeApi: () => api, + }); + + expect(result.kind).toBe("pulled"); + expect(readFileSync(snapshotPath, "utf8")).toBe(payload); + }); + + it("reconciles a stale inconsistent sidecar before comparing a newer local snapshot", async () => { + const localSnapshot: GraphSnapshot = { + ...FIXTURE_SNAPSHOT, + graph: { ...FIXTURE_SNAPSHOT.graph, repo_key: deriveProjectKey(tmpCwd).key }, + observation: { ...FIXTURE_SNAPSHOT.observation, ts: "2026-06-03T00:00:00.000Z" }, + nodes: [{ ...FIXTURE_SNAPSHOT.nodes[0]!, label: "local" }], + }; + const snapshotPath = join(baseDir, "snapshots", "head1234abcd.json"); + mkdirSync(join(baseDir, "snapshots"), { recursive: true }); + writeFileSync(snapshotPath, canonicalSnapshot(localSnapshot)); + writeLastBuild(baseDir, { + ts: Date.parse("2026-06-01T00:00:00.000Z"), + commit_sha: "head1234abcd", + snapshot_sha256: CLOUD_PAYLOAD_SHA, + node_count: 1, + edge_count: 0, + }, worktreeIdFromCwd(tmpCwd)); + + const { api } = makeMockApi({ + selectReturns: [{ + snapshot_jsonb: CLOUD_PAYLOAD, + snapshot_sha256: CLOUD_PAYLOAD_SHA, + ts: "2026-06-02T00:00:00.000Z", + node_count: 1, edge_count: 0, + worktree_id: "remote-wt", + }], + }); + const result = await pullSnapshot(tmpCwd, { + loadConfig: makeConfig, + readHead: () => "head1234abcd", + makeApi: () => api, + }); + + expect(result.kind).toBe("local-newer"); + expect(readFileSync(snapshotPath, "utf8")).toBe(canonicalSnapshot(localSnapshot)); + }); + it("codex P1 regression: local sidecar points at a DIFFERENT commit → ignored, pull proceeds", async () => { // Scenario from codex review: // 1. User built commit B locally (ts=2_000_000_000_000)