From 98104757ce27a948e747254cb02b2d128efbc656 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 26 Sep 2026 04:37:11 -0700 Subject: [PATCH] perf(crabbox): share the final source mirror inventory (#158828) * perf(crabbox): share the final source mirror inventory Collect mirror stamps during the fresh staging manifest walk, retaining the separate pre-Git integrity scan and both source freeze checks. Remove one complete payload traversal without changing persisted formats or cache custody. * fix(crabbox): require recorded mirror allocations Enforce the staging owner contract before opening the mirror database, including cold rebuilds. Dispose invalid allocations and report the invariant instead of accepting an empty final inventory. --- scripts/crabbox-source-capsule.mts | 32 ++++++-- scripts/crabbox-source-mirror.mts | 57 ++++++++------ scripts/crabbox-staging.mts | 16 +++- test/scripts/crabbox-source-capsule.test.ts | 82 ++++++++++++++++++++- 4 files changed, 154 insertions(+), 33 deletions(-) diff --git a/scripts/crabbox-source-capsule.mts b/scripts/crabbox-source-capsule.mts index aa8ecadde60a..78267f9b988f 100644 --- a/scripts/crabbox-source-capsule.mts +++ b/scripts/crabbox-source-capsule.mts @@ -23,7 +23,12 @@ import { copyFileDescriptorSync } from "@openclaw/fs-safe/advanced"; import { sha256FileSync } from "@openclaw/fs-safe/durability"; import { FsSafeError } from "@openclaw/fs-safe/errors"; import { z } from "zod"; -import { mirrorStatStamp, openSourceMirror, type MirrorFile } from "./crabbox-source-mirror.mts"; +import { + mirrorStatStamp, + openSourceMirror, + recordMirrorEntry, + type MirrorFile, +} from "./crabbox-source-mirror.mts"; import { captureSourceWitness } from "./crabbox-staging-witness.mts"; import { createMirrorStaging, createStaging, type StagingHandle } from "./crabbox-staging.mts"; @@ -289,8 +294,15 @@ export function prepareCrabboxSourceCapsule(options: { } mkdirSync(options.syncRoot, { recursive: true }); const witness = captureSourceWitness(repoRoot, sourceSha); - let mirror = - options.reuseMirror && witness ? createMirrorStaging(options.syncRoot, repoRoot) : undefined; + function allocateMirror() { + const allocated = createMirrorStaging(options.syncRoot, repoRoot); + if (allocated && !allocated.staging.recorded) { + allocated.discard(); + throw new Error("source mirror requires recorded staging; source was not uploaded"); + } + return allocated; + } + let mirror = options.reuseMirror && witness ? allocateMirror() : undefined; let cache: ReturnType | undefined; try { if (mirror) { @@ -311,7 +323,7 @@ export function prepareCrabboxSourceCapsule(options: { } console.error("[crabbox] source mirror failed verification; rebuilding a cold capsule"); mirror.discard(); - mirror = createMirrorStaging(options.syncRoot, repoRoot); + mirror = allocateMirror(); if (mirror) { cache = openSourceMirror( mirror.staging.root, @@ -1124,6 +1136,7 @@ export function prepareCrabboxSourceCapsule(options: { rmSync(linkBlobs, { recursive: true, force: true }); rmSync(join(temporary, "sparse-blobs"), { force: true }); rmSync(shallow, { force: true }); + const mirrorInventory = cache ? new Map() : undefined; if (staging.recorded) { checkPreparation(sourceEnv); checkPreparation(nativeGitEnv); @@ -1137,9 +1150,16 @@ export function prepareCrabboxSourceCapsule(options: { deleted, }, witness, + mirrorInventory + ? (path, stat) => { + if (path.startsWith("source/")) { + recordMirrorEntry(mirrorInventory, path.slice("source/".length), stat); + } + } + : undefined, ); } - if (cache) { + if (cache && mirrorInventory) { const next = new Map(); for (const path of paths) { const entry = frozen.get(path)!; @@ -1156,7 +1176,7 @@ export function prepareCrabboxSourceCapsule(options: { blob: entry.blob!, }); } - cache.save(next, trackedRecords); + cache.save(next, trackedRecords, mirrorInventory); console.error( `[crabbox] source mirror ${warm ? "warm" : "cold"}: copied ${copiedFiles} files, reused ${reusedFiles} files; preparation ${Date.now() - startedAt}ms`, ); diff --git a/scripts/crabbox-source-mirror.mts b/scripts/crabbox-source-mirror.mts index 3ff84d96ff5d..978c13eb616c 100644 --- a/scripts/crabbox-source-mirror.mts +++ b/scripts/crabbox-source-mirror.mts @@ -37,37 +37,49 @@ export function mirrorStatStamp(stat: Stats) { ].join(":"); } +function mirrorArtifact(path: string) { + return [".crabbox/runs", ".crabbox/captures"].some( + (root) => path === root || path.startsWith(root + "/"), + ); +} + +export function recordMirrorEntry(entries: Map, path: string, stat: Stats) { + if (mirrorArtifact(path)) { + return; + } + // Native sync may refresh its index; selection/candidate indexes stay sealed. + if (path === ".git/index") { + if (!stat.isFile() || stat.nlink !== 1) { + throw new Error("source mirror index is not a private regular file"); + } + entries.set(path, "mutable-index"); + } else if (stat.isDirectory()) { + if (path !== ".crabbox") { + entries.set(path, `dir:${stat.dev}:${stat.ino}:${stat.mode}`); + } + } else if (stat.isFile() || stat.isSymbolicLink()) { + entries.set(path, mirrorStatStamp(stat)); + } else { + throw new Error("source mirror contains an unsupported file kind"); + } +} + function payloadInventory(directory: string) { const entries = new Map(); let objectBytes = 0; function walk(parent: string) { for (const name of readdirSync(join(directory, parent))) { const path = parent ? `${parent}/${name}` : name; - // Native sync may refresh its index; the sealed selection/candidate indexes - // are separate files. Run outputs belong to staging's preservation owner. - if (path === ".crabbox/runs" || path === ".crabbox/captures") { + // Run outputs belong to staging's preservation owner. + if (mirrorArtifact(path)) { continue; } const stat = lstatSync(join(directory, path)); - if (path === ".git/index") { - if (!stat.isFile() || stat.nlink !== 1) { - throw new Error("source mirror index is not a private regular file"); - } - entries.set(path, "mutable-index"); - continue; - } + recordMirrorEntry(entries, path, stat); if (stat.isDirectory()) { - if (path !== ".crabbox") { - entries.set(path, `dir:${stat.dev}:${stat.ino}:${stat.mode}`); - } walk(path); - } else if (stat.isFile() || stat.isSymbolicLink()) { - entries.set(path, mirrorStatStamp(stat)); - if (path.startsWith(".git/objects/")) { - objectBytes += stat.size; - } - } else { - throw new Error("source mirror contains an unsupported file kind"); + } else if (path.startsWith(".git/objects/")) { + objectBytes += stat.size; } } } @@ -153,8 +165,9 @@ export function openSourceMirror( files, tracked: metadata?.tracked, close, - save(next: Map, tracked: string) { - const inventory = payloadInventory(directory).entries; + save(next: Map, tracked: string, inventory: Map) { + // Staging's final payload walk supplies these fresh observations. No payload + // writes may follow it; both seals describe the same frozen filesystem. for (const [path, file] of next) { if (inventory.get(path) !== file.stamp) { throw new Error( diff --git a/scripts/crabbox-staging.mts b/scripts/crabbox-staging.mts index b5f140acab36..3f09280230ad 100644 --- a/scripts/crabbox-staging.mts +++ b/scripts/crabbox-staging.mts @@ -17,6 +17,7 @@ import { rmSync, rmdirSync, writeFileSync, + type Stats, } from "node:fs"; import { basename, dirname, isAbsolute, join } from "node:path"; import { performance } from "node:perf_hooks"; @@ -345,6 +346,7 @@ function inventory( payload: string, known = new Map(), preservedArtifacts = false, + observeEntry?: (path: string, stat: Stats) => void, ): Entry[] { const entries: Entry[] = []; const walk = (directory: string, parent: string) => { @@ -358,6 +360,7 @@ function inventory( } const absolute = join(payload, path); const stat = lstatSync(absolute); + observeEntry?.(path, stat); if (stat.isDirectory()) { if (!preservedArtifacts || path !== "source/.crabbox") { entries.push({ path, kind: "directory" }); @@ -415,7 +418,11 @@ export type StagingHandle = { recorded: boolean; root: string; payload: string; - prepared: (source: FrozenSource, witness?: SourceWitness) => void; + prepared: ( + source: FrozenSource, + witness?: SourceWitness, + observeEntry?: (path: string, stat: Stats) => void, + ) => void; admitted: (claims?: ClaimNamespace, leases?: string[]) => void; settled: (leases?: string[]) => void; preserved: (artifacts: CrabboxArtifactEvidence) => void; @@ -498,7 +505,7 @@ function stagingHandle(root: string, initialReceipt: Receipt, recorded: boolean) recorded, root, payload, - prepared(source, witness) { + prepared(source, witness, observeEntry) { if (!recorded) { return; } @@ -512,7 +519,10 @@ function stagingHandle(root: string, initialReceipt: Receipt, recorded: boolean) }, ]), ); - const manifest: Manifest = { source, entries: inventory(payload, known) }; + const manifest: Manifest = { + source, + entries: inventory(payload, known, false, observeEntry), + }; const bytes = JSON.stringify(manifest) + "\n"; if (Buffer.byteLength(bytes) > manifestLimit) { throw new Error("staging manifest exceeds the recovery metadata limit"); diff --git a/test/scripts/crabbox-source-capsule.test.ts b/test/scripts/crabbox-source-capsule.test.ts index cf47006be6b1..f0c26bb4af3c 100644 --- a/test/scripts/crabbox-source-capsule.test.ts +++ b/test/scripts/crabbox-source-capsule.test.ts @@ -7,6 +7,7 @@ import { mkdirSync, readFileSync, readlinkSync, + renameSync, rmSync, symlinkSync, writeFileSync, @@ -17,9 +18,20 @@ import { prepareCrabboxSourceCapsule, type CrabboxSourceCapsule, } from "../../scripts/crabbox-source-capsule.mts"; +import { createMirrorStaging } from "../../scripts/crabbox-staging.mts"; import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js"; import { createNestedGitEnv } from "../helpers/temp-repo.js"; +vi.mock("node:fs", async (importOriginal) => { + const fs = await importOriginal(); + return { ...fs, lstatSync: vi.fn(fs.lstatSync) }; +}); + +vi.mock("../../scripts/crabbox-staging.mts", async (importOriginal) => { + const staging = await importOriginal(); + return { ...staging, createMirrorStaging: vi.fn(staging.createMirrorStaging) }; +}); + const temporary = useAutoCleanupTempDirTracker(afterEach); afterEach(() => { vi.restoreAllMocks(); @@ -126,9 +138,47 @@ function fileIdentity(path: string) { } describe.skipIf(process.platform === "win32")("persistent Crabbox source capsules", () => { + it.each(["initial allocation", "cold rebuild"])( + "rejects an unrecorded mirror before freezing during %s", + async (allocation) => { + const f = fixture(); + if (allocation === "cold rebuild") { + const first = f.prepare(); + first.cleanup(); + writeFileSync(join(first.directory, "stable.txt"), "corrupted cached bytes\n"); + } + const actual = await vi.importActual( + "../../scripts/crabbox-staging.mts", + ); + let unrecordedRoot: string | undefined; + const rejectRecording = (...args: Parameters) => { + const mirror = actual.createMirrorStaging(...args); + if (!mirror) { + throw new Error("fixture requires a mirror allocation"); + } + mirror.staging.recorded = false; + unrecordedRoot = mirror.staging.root; + return mirror; + }; + const allocate = vi.mocked(createMirrorStaging); + if (allocation === "cold rebuild") { + allocate.mockImplementationOnce(actual.createMirrorStaging); + } + allocate.mockImplementationOnce(rejectRecording); + const selected = join(f.root, "selected"); + expect(() => + f.prepare(true, `require("node:fs").writeFileSync(${JSON.stringify(selected)}, "ran");`), + ).toThrow("source mirror requires recorded staging"); + expect(unrecordedRoot).toBeDefined(); + expect(existsSync(unrecordedRoot!)).toBe(false); + expect(existsSync(selected)).toBe(false); + }, + ); + it("keeps unchanged files and a warm index while updating eligibility and raw source bytes", () => { - const f = fixture(); + const f = fixture({ "rename.txt": "renamed source\n" }); writeFileSync(join(f.repository, "future.txt"), "initial untracked bytes\n"); + writeFileSync(join(f.repository, "promoted.ignored"), "initially ignored bytes\n"); writeFileSync(join(f.repository, "staged.ignored"), "staged ignored bytes\r\n"); writeFileSync(join(f.repository, "credentials.secret"), "synthetic secret\n"); f.git(f.repository, "add", "--force", "staged.ignored"); @@ -141,11 +191,16 @@ describe.skipIf(process.platform === "win32")("persistent Crabbox source capsule first.cleanup(); const arrivingSecret = join(f.repository, "arriving.secret"); + vi.mocked(lstatSync).mockClear(); const unchanged = f.prepare( true, `require("node:fs").writeFileSync(${JSON.stringify(arrivingSecret)},"synthetic secret");`, ); try { + // One pre-Git integrity check and one shared final seal per retained file. + expect( + vi.mocked(lstatSync).mock.calls.filter(([path]) => path === join(directory, "stable.txt")), + ).toHaveLength(2); expect(unchanged.directory).toBe(directory); for (const [path, identity] of before) { expect(fileIdentity(join(directory, path)), path).toEqual(identity); @@ -161,6 +216,8 @@ describe.skipIf(process.platform === "win32")("persistent Crabbox source capsule writeFileSync(join(f.repository, "change.txt"), "updated raw bytes\r\n"); writeFileSync(join(f.repository, ".gitignore"), "*.secret\n*.ignored\nfuture.txt\n"); writeFileSync(join(f.repository, "new.txt"), "new untracked bytes\n"); + renameSync(join(f.repository, "rename.txt"), join(f.repository, "renamed.txt")); + f.git(f.repository, "add", "--force", "promoted.ignored"); rmSync(join(f.repository, "deleted.txt")); const warm = f.prepare(); try { @@ -176,10 +233,12 @@ describe.skipIf(process.platform === "win32")("persistent Crabbox source capsule ".gitignore", "change.txt", "new.txt", + "promoted.ignored", + "renamed.txt", "stable.txt", "staged.ignored", ]); - for (const path of ["credentials.secret", "deleted.txt", "future.txt"]) { + for (const path of ["credentials.secret", "deleted.txt", "future.txt", "rename.txt"]) { expect(existsSync(join(directory, path)), path).toBe(false); } f.expectColdEquivalent(warm); @@ -291,6 +350,9 @@ describe.skipIf(process.platform === "win32")("persistent Crabbox source capsule it.each([ "missing file", "corrupt file", + "extra file", + "corrupt Git config", + "corrupt Git object", "corrupt metadata", "missing metadata", "hardlink metadata", @@ -308,6 +370,22 @@ describe.skipIf(process.platform === "win32")("persistent Crabbox source capsule rmSync(join(first.directory, "stable.txt")); } else if (damage === "corrupt file") { writeFileSync(join(first.directory, "stable.txt"), "corrupted cached bytes\n"); + } else if (damage === "extra file") { + writeFileSync(join(first.directory, "unexpected.txt"), "outside-owner bytes\n"); + } else if (damage === "corrupt Git config") { + writeFileSync(join(first.directory, ".git", "config"), "invalid Git configuration\n"); + } else if (damage === "corrupt Git object") { + const object = join( + first.directory, + ".git", + "objects", + first.carrier.slice(0, 2), + first.carrier.slice(2), + ); + const mode = lstatSync(object).mode & 0o777; + chmodSync(object, 0o600); + writeFileSync(object, "corrupted private Git object\n"); + chmodSync(object, mode); } else if (damage === "missing index") { rmSync(join(first.directory, ".git", "mirror-candidate-index")); } else if (damage === "symlink index" || damage === "hardlink index") {