From b9e5d6ecdb05ab7244c90976e07c950f8d09b15b Mon Sep 17 00:00:00 2001 From: Dotta Date: Thu, 1 Oct 2026 22:03:21 -0500 Subject: [PATCH] perf(runner): keep verified snapshot copies progressing Co-Authored-By: Paperclip --- .../native-distribution-integrity.test.ts | 89 ++++++++++++++++++- .../acpx/native-distribution-integrity.ts | 48 +++++----- 2 files changed, 113 insertions(+), 24 deletions(-) diff --git a/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.test.ts b/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.test.ts index 3b69adcd2b..b130bb1185 100644 --- a/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.test.ts +++ b/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.test.ts @@ -1,6 +1,8 @@ import { createHash } from "node:crypto"; import { once } from "node:events"; -import { chmod, copyFile, mkdir, mkdtemp, open, readFile, readdir, rm, stat, symlink, writeFile, type FileHandle } from "node:fs/promises"; +import { constants } from "node:fs"; +import { execFileSync } from "node:child_process"; +import { chmod, copyFile, link, mkdir, mkdtemp, open, readFile, readdir, realpath, rename, rm, stat, symlink, writeFile, type FileHandle } from "node:fs/promises"; import { createServer, type Server } from "node:net"; import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; @@ -11,7 +13,7 @@ import { createNativeAcpxDistributionSnapshot, readNativeAcpxDistributionEntries vi.mock("node:fs/promises", async (importOriginal) => { const original = await importOriginal(); - return { ...original, rm: vi.fn(original.rm), mkdir: vi.fn(original.mkdir), chmod: vi.fn(original.chmod) }; + return { ...original, open: vi.fn(original.open), rm: vi.fn(original.rm), mkdir: vi.fn(original.mkdir), chmod: vi.fn(original.chmod) }; }); const roots: string[] = []; @@ -107,6 +109,66 @@ describe("native ACPX execution closure", () => { await symlink("closure.json", join(declaration.distributionRoot, "runtime")); await expect(installation.openCommand()).rejects.toThrow("symbolic link"); }); + it.each(["fifo", "directory", "hardlink", "symlink", "size", "mode"] as const)( + "rejects a %s swapped in immediately before open without reading it or leaking its descriptor", async kind => { + const declaration = await fixture(); + const entries = await readNativeAcpxDistributionEntries(declaration); + const original = await vi.importActual("node:fs/promises"); + const source = await realpath(declaration.distributionRoot); + const path = join(source, "runtime"); + const originalPath = join(source, "original"); + const prototype = await filePrototype(path); + const read = vi.spyOn(prototype, "read"); + let opened: FileHandle | undefined; let closed = false; + vi.mocked(open).mockImplementation(async (selected: any, flags: any, ...rest: any[]): Promise => { + if (String(selected) !== path) return original.open(selected, flags, rest[0]); + expect(flags & constants.O_NOFOLLOW).not.toBe(0); + expect(flags & constants.O_NONBLOCK).not.toBe(0); + await rename(path, originalPath); + if (kind === "fifo") execFileSync("mkfifo", [path], { timeout: 2_000 }); + else if (kind === "directory") await mkdir(path); + else if (kind === "hardlink") await link(originalPath, path); + else if (kind === "symlink") await symlink(originalPath, path); + else { + await copyFile(originalPath, path); + if (kind === "size") await writeFile(path, "short"); + else await chmod(path, 0o600); + } + opened = await original.open(selected, flags, rest[0]); + const close = opened.close.bind(opened); + opened.close = async () => { await close(); closed = true; }; + return opened; + }); + let unexpected: Awaited> | undefined; + try { + const error = await createNativeAcpxDistributionSnapshot(declaration, entries).then( + value => { unexpected = value; return undefined; }, error => error); + expect(error).toBeInstanceOf(Error); + expect(read).not.toHaveBeenCalled(); + if (kind === "symlink") expect(opened).toBeUndefined(); + else { expect(opened).toBeDefined(); expect(closed).toBe(true); } + } finally { + vi.mocked(open).mockImplementation(original.open); + if (unexpected) { await unexpected.commandDirectory.close(); await unexpected.snapshot.close(); } + } + }); + it("rejects a same-byte pathname replacement after reading the held descriptor", async () => { + const declaration = await fixture(); const entries = await readNativeAcpxDistributionEntries(declaration); + const path = join(declaration.distributionRoot, "runtime"); + const prototype = await filePrototype(path); const originalRead = prototype.read; + let replaced = false; + vi.spyOn(prototype, "read").mockImplementation(async function (this: FileHandle, ...args: any[]): Promise { + const result = await originalRead.apply(this, args as never); + if (!replaced) { + replaced = true; + await rename(path, path + ".original"); + await copyFile(path + ".original", path); + } + return result; + }); + await expect(createNativeAcpxDistributionSnapshot(declaration, entries)).rejects.toThrow("changed while read"); + expect(replaced).toBe(true); + }); it("launches only fixed arguments from a frozen snapshot after installed files change", async () => { const declaration = await fixture(); const lease = await (await verifyNativeAcpxInstallation(declaration)).openCommand(); expect(() => lease.spawn(["--untrusted-override"])).toThrow("fixed profile arguments"); @@ -231,6 +293,27 @@ describe("native ACPX execution closure", () => { for (const entry of entries) expect(hash(await readFile(join(created.snapshot.roots[0]!, entry.path)))).toBe(entry.sha256); } finally { await created.commandDirectory.close(); await created.snapshot.close(); } }); + it("uses released capacity while an earlier copy is still pending", async () => { + const { declaration, entries } = await manyFileFixture(Array.from({ length: 40 }, (_, index) => 100 + index)); + const prototype = await filePrototype(join(declaration.distributionRoot, "runtime")); + const originalRead = prototype.read; const hold = gate(); let laterReached = false; + let active = 0; let peak = 0; + vi.spyOn(prototype, "read").mockImplementation(async function (this: FileHandle, ...args: any[]): Promise { + active++; peak = Math.max(peak, active); + try { + if (args[0].length === 100) await hold.promise; + if (args[0].length === 139) laterReached = true; + return await originalRead.apply(this, args as never); + } finally { active--; } + }); + const creating = createNativeAcpxDistributionSnapshot(declaration, entries); + let created: Awaited | undefined; + try { + try { await vi.waitFor(() => expect(laterReached).toBe(true)); expect(active).toBeGreaterThan(0); expect(peak).toBeLessThanOrEqual(32); } + finally { hold.release(); created = await creating; } + expect(Object.keys(created.snapshot.digests).slice(0, entries.length)).toEqual(entries.map(entry => join(created!.snapshot.roots[0]!, entry.path))); + } finally { if (created) { await created.commandDirectory.close(); await created.snapshot.close(); } } + }); it("bounds simultaneous buffers and gives an oversized admitted file exclusive capacity", async () => { const mib = 1024 * 1024; const { declaration, entries } = await manyFileFixture([17 * mib, 17 * mib, 33 * mib]); @@ -266,7 +349,7 @@ describe("native ACPX execution closure", () => { this.close = async () => { await originalClose(); invalidClosed.release(); }; await entered.promise; } - if (size === 92) { entered.release(); await hold.promise; } + if (size !== 91) { entered.release(); await hold.promise; } return originalRead.apply(this, args as never); }); const creating = createNativeAcpxDistributionSnapshot(declaration, entries); diff --git a/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.ts b/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.ts index 287bcca8c7..11d18eba9e 100644 --- a/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.ts +++ b/packages/paperclip-runner/src/drivers/acpx/native-distribution-integrity.ts @@ -127,11 +127,13 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist const copyEntry = async (entry: NativeAcpxDistributionEntry): Promise => { const path = join(source, ...entry.path.split("/")); if (await realpath(path) !== path) throw new Error("Native ACPX closure contains a symbolic link"); - const before = await lstat(path, { bigint: true }); - if (!before.isFile() || before.nlink !== 1n || before.size !== BigInt(entry.size) || Boolean(before.mode & 0o111n) !== entry.executable) throw new Error("Native ACPX closure file identity is invalid"); - const file = await open(path, constants.O_RDONLY | constants.O_NOFOLLOW); + // Bind metadata to the descriptor we will read, without a redundant + // pathname stat. NONBLOCK prevents a raced-in FIFO from blocking open; + // only a bounded regular single-link file may reach the read below. + const file = await open(path, constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK); try { - if (!same(before, await file.stat({ bigint: true }))) throw new Error("Native ACPX closure file changed before read"); + const before = await file.stat({ bigint: true }); + if (!before.isFile() || before.nlink !== 1n || before.size !== BigInt(entry.size) || Boolean(before.mode & 0o111n) !== entry.executable) throw new Error("Native ACPX closure file identity is invalid"); const bytes = Buffer.alloc(entry.size); let offset = 0; while (offset < bytes.length) { @@ -145,25 +147,29 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist await writeFile(target, bytes, { mode: entry.executable ? 0o500 : 0o400, flag: "wx" }); } finally { await file.close(); } }; - for (let start = 0; start < entries.length;) { - let end = start; - let bytes = 0; - while (end < entries.length && end - start < NATIVE_COPY_CONCURRENCY) { - const size = entries[end]!.size; - if (end > start && bytes + size > NATIVE_COPY_BUFFER_BYTES) break; - bytes += size; - end++; + // Keep bounded capacity occupied when one file is slower than its peers. + // Every task catches its rejection before releasing capacity; once failed, + // no new copy is admitted and all owned descriptors drain before cleanup. + const active = new Set>(); + let activeBytes = 0; + let failed = false; + let failure: unknown; + for (const entry of entries) { + while (!failed && (active.size >= NATIVE_COPY_CONCURRENCY + || (active.size > 0 && activeBytes + entry.size > NATIVE_COPY_BUFFER_BYTES))) { + await Promise.race(active); } - const batch = entries.slice(start, end); - // Never clean the private root while another worker can still write or - // close a descriptor. Stop scheduling new batches after any rejection. - const copied = await Promise.allSettled(batch.map(copyEntry)); - const failure = copied.find((result) => result.status === "rejected"); - if (failure?.status === "rejected") throw failure.reason; - // Worker completion order must not change the module guard or manifest. - for (const entry of batch) digests[join(packageRoot, ...entry.path.split("/"))] = entry.sha256; - start = end; + if (failed) break; + activeBytes += entry.size; + const copying = copyEntry(entry).catch(error => { + if (!failed) { failed = true; failure = error; } + }).finally(() => { activeBytes -= entry.size; active.delete(copying); }); + active.add(copying); } + await Promise.all(active); + if (failed) throw failure; + // Completion order must not change the module guard or manifest. + for (const entry of entries) digests[join(packageRoot, ...entry.path.split("/"))] = entry.sha256; if (!same(rootBefore, await heldRoot.stat({ bigint: true })) || !same(rootBefore, await lstat(source, { bigint: true }))) throw new Error("Native ACPX distribution root changed during snapshot"); const executable = join(packageRoot, ...input.executable.split("/")); const args = [...(input.entrypoint === undefined ? [] : ["--require", join(packageRoot, GUARD), join(packageRoot, ...input.entrypoint.split("/"))]), ...input.fixedArguments];