mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 21:05:21 +02:00
perf(runner): bound cold runtime directory work
Batch Pi inventory metadata checks while preserving depth-first ordering and revalidating directories immediately before descent. Deduplicate native snapshot parent creation and bound sealing work without changing the complete byte verification, private snapshot, or runtime deadlines. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
5cd6d3d523
commit
586af4d7f3
4 files changed
+180
-23
No files matched your search
@@ -11,7 +11,7 @@ import { createNativeAcpxDistributionSnapshot, readNativeAcpxDistributionEntries
|
||||
|
||||
vi.mock("node:fs/promises", async (importOriginal) => {
|
||||
const original = await importOriginal<typeof import("node:fs/promises")>();
|
||||
return { ...original, rm: vi.fn(original.rm) };
|
||||
return { ...original, rm: vi.fn(original.rm), mkdir: vi.fn(original.mkdir), chmod: vi.fn(original.chmod) };
|
||||
});
|
||||
|
||||
const roots: string[] = [];
|
||||
@@ -57,6 +57,22 @@ async function manyFileFixture(sizes: number[]) {
|
||||
await writeFile(declaration.manifestPath, JSON.stringify({ entries }));
|
||||
return { declaration: { ...declaration, expectedClosureSha256: hash(JSON.stringify(entries)) }, entries };
|
||||
}
|
||||
async function directoryFixture() {
|
||||
const declaration = await fixture();
|
||||
const entries = await readNativeAcpxDistributionEntries(declaration);
|
||||
for (let index = 0; index < 40; index++) {
|
||||
const parent = `dir-${String(index).padStart(2, "0")}/nested`;
|
||||
await mkdir(join(declaration.distributionRoot, parent), { recursive: true });
|
||||
for (const name of ["a", "b"]) {
|
||||
const path = `${parent}/${name}`; const bytes = Buffer.from(path);
|
||||
await writeFile(join(declaration.distributionRoot, path), bytes, { mode: 0o600 });
|
||||
entries.push({ path, sha256: hash(bytes), size: bytes.length, executable: false });
|
||||
}
|
||||
}
|
||||
entries.sort((a, b) => a.path < b.path ? -1 : a.path > b.path ? 1 : 0);
|
||||
await writeFile(declaration.manifestPath, JSON.stringify({ entries }));
|
||||
return { declaration: { ...declaration, expectedClosureSha256: hash(JSON.stringify(entries)) }, entries };
|
||||
}
|
||||
async function filePrototype(path: string): Promise<FileHandle> {
|
||||
const probe = await open(path, "r"); const prototype = Object.getPrototypeOf(probe) as FileHandle;
|
||||
await probe.close(); return prototype;
|
||||
@@ -143,6 +159,56 @@ describe("native ACPX execution closure", () => {
|
||||
const denied = await output((await (await verifyNativeAcpxInstallation(evil)).openCommand()).spawn());
|
||||
expect(denied.code).not.toBe(0); expect(denied.error).toContain("escaped its closed distribution");
|
||||
}, 30_000);
|
||||
it("creates each private parent once and seals every directory before returning", async () => {
|
||||
const { declaration, entries } = await directoryFixture();
|
||||
const creatingStart = vi.mocked(mkdir).mock.calls.length;
|
||||
const created = await createNativeAcpxDistributionSnapshot(declaration, entries);
|
||||
try {
|
||||
const root = created.snapshot.roots[0]!;
|
||||
const calls = vi.mocked(mkdir).mock.calls.slice(creatingStart).map(([path]) => String(path)).filter(path => path.startsWith(root));
|
||||
expect(calls).toHaveLength(81); expect(new Set(calls).size).toBe(81);
|
||||
for (const path of calls) expect((await stat(path)).mode & 0o777).toBe(0o500);
|
||||
for (const entry of entries) {
|
||||
const path = join(root, entry.path);
|
||||
expect(hash(await readFile(path))).toBe(entry.sha256);
|
||||
expect((await stat(path)).mode & 0o777).toBe(entry.executable ? 0o500 : 0o400);
|
||||
}
|
||||
} finally { await created.commandDirectory.close(); await created.snapshot.close(); }
|
||||
});
|
||||
it.each(["create", "seal"] as const)("bounds %s batches and drains failures before cleanup or further scheduling", async stage => {
|
||||
const { declaration, entries } = await directoryFixture();
|
||||
const original = await vi.importActual<typeof import("node:fs/promises")>("node:fs/promises");
|
||||
const hold = gate(); let started = 0; let active = 0; let peak = 0; let settled = false;
|
||||
const removalStart = vi.mocked(rm).mock.calls.length;
|
||||
const operation = async (path: any, mode: any) => {
|
||||
const selected = stage === "create"
|
||||
? /paperclip-acpx-native-.*\/distribution\/dir-\d+$/.test(String(path))
|
||||
: mode === 0o500;
|
||||
if (!selected) return stage === "create" ? original.mkdir(path, mode) : original.chmod(path, mode);
|
||||
const index = started++; active++; peak = Math.max(peak, active);
|
||||
try {
|
||||
if (index === 0) throw new Error(`fixture ${stage} failure`);
|
||||
await hold.promise;
|
||||
return stage === "create" ? await original.mkdir(path, mode) : await original.chmod(path, mode);
|
||||
} finally { active--; }
|
||||
};
|
||||
if (stage === "create") vi.mocked(mkdir).mockImplementation(operation as typeof mkdir);
|
||||
else vi.mocked(chmod).mockImplementation(operation as typeof chmod);
|
||||
const creating = createNativeAcpxDistributionSnapshot(declaration, entries);
|
||||
void creating.then(() => { settled = true; }, () => { settled = true; });
|
||||
try {
|
||||
await vi.waitFor(() => expect(started).toBe(32));
|
||||
expect(peak).toBeGreaterThan(1); expect(peak).toBeLessThanOrEqual(32); expect(settled).toBe(false);
|
||||
expect(vi.mocked(rm).mock.calls.slice(removalStart).filter(([path]) => String(path).includes("paperclip-acpx-native-"))).toHaveLength(0);
|
||||
} finally { hold.release(); }
|
||||
try {
|
||||
await expect(creating).rejects.toThrow(`fixture ${stage} failure`);
|
||||
expect(active).toBe(0); expect(started).toBe(32);
|
||||
const removed = vi.mocked(rm).mock.calls.slice(removalStart).map(([path]) => String(path)).filter(path => path.includes("paperclip-acpx-native-"));
|
||||
expect(removed).toHaveLength(1);
|
||||
await expect(stat(removed[0]!)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
} finally { vi.mocked(mkdir).mockImplementation(original.mkdir); vi.mocked(chmod).mockImplementation(original.chmod); }
|
||||
});
|
||||
it("copies concurrently within its descriptor bound and retains canonical manifest order", async () => {
|
||||
const { declaration, entries } = await manyFileFixture(Array.from({ length: 40 }, (_, index) => 100 + index));
|
||||
const prototype = await filePrototype(join(declaration.distributionRoot, "runtime"));
|
||||
|
||||
@@ -2,7 +2,7 @@ import { createHash } from "node:crypto";
|
||||
import { constants } from "node:fs";
|
||||
import { chmod, lstat, mkdir, mkdtemp, open, realpath, rm, writeFile, type FileHandle } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { dirname, isAbsolute, join, relative, resolve } from "node:path";
|
||||
import { isAbsolute, join, relative, resolve } from "node:path";
|
||||
import type { AcpxPrivateSnapshot } from "./private-snapshot.js";
|
||||
|
||||
export interface NativeAcpxDistributionEntry { path: string; sha256: string; size: number; executable: boolean }
|
||||
@@ -31,6 +31,7 @@ const MAX_NATIVE_MANIFEST_BYTES = 4 * 1024 * 1024;
|
||||
// Bound both file descriptors and buffers. A single larger admitted file runs
|
||||
// alone and remains subject to MAX_NATIVE_FILE_BYTES.
|
||||
const NATIVE_COPY_CONCURRENCY = 32;
|
||||
const NATIVE_DIRECTORY_CONCURRENCY = 32;
|
||||
const NATIVE_COPY_BUFFER_BYTES = 32 * 1024 * 1024;
|
||||
const BOOTSTRAP = ".paperclip-native-entry.cjs";
|
||||
const GUARD = ".paperclip-native-module-guard.cjs";
|
||||
@@ -102,6 +103,27 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist
|
||||
if (!same(rootBefore, await heldRoot.stat({ bigint: true }))) throw new Error("Native ACPX distribution root changed before snapshot");
|
||||
await mkdir(packageRoot, { mode: 0o700 });
|
||||
if (input.isolatedCacheEnvironmentName) await mkdir(cacheRoot, { mode: 0o700 });
|
||||
// Build each private parent once. Level ordering prevents child creation
|
||||
// from racing its parent; batches bound work and drain before any cleanup.
|
||||
const parentLevels = new Map<number, Set<string>>();
|
||||
for (const entry of entries) {
|
||||
const parts = entry.path.split("/");
|
||||
for (let depth = 1; depth < parts.length; depth++) {
|
||||
const path = join(packageRoot, ...parts.slice(0, depth));
|
||||
const level = parentLevels.get(depth) ?? new Set<string>();
|
||||
level.add(path); parentLevels.set(depth, level); directories.add(path);
|
||||
}
|
||||
}
|
||||
const directoryBatch = async (paths: string[], operation: (path: string) => Promise<unknown>): Promise<void> => {
|
||||
for (let start = 0; start < paths.length; start += NATIVE_DIRECTORY_CONCURRENCY) {
|
||||
const settled = await Promise.allSettled(paths.slice(start, start + NATIVE_DIRECTORY_CONCURRENCY).map(operation));
|
||||
const failed = settled.find(result => result.status === "rejected");
|
||||
if (failed?.status === "rejected") throw failed.reason;
|
||||
}
|
||||
};
|
||||
for (const depth of [...parentLevels.keys()].sort((a, b) => a - b)) {
|
||||
await directoryBatch([...parentLevels.get(depth)!].sort(), path => mkdir(path, { mode: 0o700 }));
|
||||
}
|
||||
const copyEntry = async (entry: NativeAcpxDistributionEntry): Promise<void> => {
|
||||
const path = join(source, ...entry.path.split("/"));
|
||||
if (await realpath(path) !== path) throw new Error("Native ACPX closure contains a symbolic link");
|
||||
@@ -120,9 +142,6 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist
|
||||
if (!same(before, await file.stat({ bigint: true })) || !same(before, await lstat(path, { bigint: true })) || await realpath(path) !== path) throw new Error("Native ACPX closure file changed while read");
|
||||
if (sha256(bytes) !== entry.sha256) throw new Error(`Native ACPX closure file digest mismatch: ${entry.path}`);
|
||||
const target = join(packageRoot, ...entry.path.split("/"));
|
||||
await mkdir(dirname(target), { recursive: true, mode: 0o700 });
|
||||
let directory = dirname(target);
|
||||
while (directory !== privateRoot) { directories.add(directory); directory = dirname(directory); }
|
||||
await writeFile(target, bytes, { mode: entry.executable ? 0o500 : 0o400, flag: "wx" });
|
||||
} finally { await file.close(); }
|
||||
};
|
||||
@@ -172,7 +191,7 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist
|
||||
const manifest = Buffer.from(JSON.stringify({ roots: [packageRoot], executable: null, digests }));
|
||||
const manifestPath = join(privateRoot, "manifest.json");
|
||||
await writeFile(manifestPath, manifest, { mode: 0o400, flag: "wx" });
|
||||
for (const directory of directories) await chmod(directory, 0o500);
|
||||
await directoryBatch([...directories], path => chmod(path, 0o500));
|
||||
commandDirectory = await open(packageRoot, constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_DIRECTORY);
|
||||
return { commandDirectory, bootstrap, snapshot: { roots: [packageRoot], executable: null, digests, handoff: { path: manifestPath, digest: sha256(manifest) }, close } };
|
||||
} catch (error) {
|
||||
|
||||
@@ -1,10 +1,15 @@
|
||||
import { link, mkdir, mkdtemp, open, realpath, rm, symlink, writeFile, type FileHandle } from "node:fs/promises";
|
||||
import { link, lstat, mkdir, mkdtemp, open, readdir, realpath, rename, rm, symlink, writeFile, type FileHandle } from "node:fs/promises";
|
||||
import { createHash } from "node:crypto";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { inventoryPiRuntimeFiles, PI_RUNTIME_MANIFEST_SCHEMA, verifyPiRuntimeManifest, type PiRuntimeManifest } from "./pi-verified-runtime.js";
|
||||
|
||||
vi.mock("node:fs/promises", async (importOriginal) => {
|
||||
const original = await importOriginal<typeof import("node:fs/promises")>();
|
||||
return { ...original, lstat: vi.fn(original.lstat), readdir: vi.fn(original.readdir) };
|
||||
});
|
||||
|
||||
const temporary: string[] = [];
|
||||
afterEach(async () => { vi.restoreAllMocks(); for (const root of temporary.splice(0)) await rm(root, { recursive: true, force: true }); });
|
||||
async function fixture() {
|
||||
@@ -72,6 +77,56 @@ describe("Pi complete runtime manifest", () => {
|
||||
expect(active).toBe(0); expect(started).toBe(32);
|
||||
});
|
||||
|
||||
it("bounds discovery and drains a failed batch before descent or later scheduling", async () => {
|
||||
const root = await mkdtemp(join(tmpdir(), "paperclip-pi-discovery-")); temporary.push(root);
|
||||
for (let index = 0; index < 40; index++) await writeFile(join(root, `file-${String(index).padStart(2, "0")}`), "data");
|
||||
const original = await vi.importActual<typeof import("node:fs/promises")>("node:fs/promises");
|
||||
let release!: () => void; const hold = new Promise<void>(resolve => { release = resolve; });
|
||||
let active = 0; let peak = 0; let started = 0; let settled = false;
|
||||
vi.mocked(lstat).mockImplementation((async (path: any, options: any) => {
|
||||
if (String(path) === root) return original.lstat(path, options);
|
||||
const index = started++; active++; peak = Math.max(peak, active);
|
||||
try {
|
||||
if (index === 0) throw new Error("fixture discovery failure");
|
||||
await hold; return await original.lstat(path, options);
|
||||
} finally { active--; }
|
||||
}) as typeof lstat);
|
||||
const inventory = inventoryPiRuntimeFiles(root);
|
||||
void inventory.then(() => { settled = true; }, () => { settled = true; });
|
||||
try {
|
||||
await vi.waitFor(() => expect(started).toBe(32));
|
||||
expect(peak).toBeLessThanOrEqual(32); expect(peak).toBeGreaterThan(1); expect(settled).toBe(false);
|
||||
} finally { release(); }
|
||||
try {
|
||||
await expect(inventory).rejects.toThrow("fixture discovery failure");
|
||||
expect(active).toBe(0); expect(started).toBe(32);
|
||||
} finally { vi.mocked(lstat).mockImplementation(original.lstat); }
|
||||
});
|
||||
|
||||
it("rechecks queued directories after earlier subtrees and never descends into a replacement symlink", async () => {
|
||||
const root = await realpath(await mkdtemp(join(tmpdir(), "paperclip-pi-discovery-race-"))); temporary.push(root);
|
||||
const outside = await realpath(await mkdtemp(join(tmpdir(), "paperclip-pi-outside-"))); temporary.push(outside);
|
||||
await mkdir(join(root, "a")); await mkdir(join(root, "b"));
|
||||
await writeFile(join(root, "a/trigger"), "inside"); await writeFile(join(outside, "private"), "must not read");
|
||||
const original = await vi.importActual<typeof import("node:fs/promises")>("node:fs/promises");
|
||||
let replaced = false;
|
||||
const readsStart = vi.mocked(readdir).mock.calls.length;
|
||||
vi.mocked(lstat).mockImplementation((async (path: any, options: any) => {
|
||||
if (String(path) === join(root, "a/trigger") && !replaced) {
|
||||
replaced = true;
|
||||
await rename(join(root, "b"), join(root, "b-retired"));
|
||||
await symlink(outside, join(root, "b"));
|
||||
}
|
||||
return original.lstat(path, options);
|
||||
}) as typeof lstat);
|
||||
try {
|
||||
await expect(inventoryPiRuntimeFiles(root)).rejects.toThrow("directory changed before traversal");
|
||||
expect(replaced).toBe(true);
|
||||
const visited = vi.mocked(readdir).mock.calls.slice(readsStart).map(([path]) => String(path));
|
||||
expect(visited).not.toContain(join(root, "b")); expect(visited).not.toContain(outside);
|
||||
} finally { vi.mocked(lstat).mockImplementation(original.lstat); }
|
||||
});
|
||||
|
||||
it("refuses unrecorded resources and duplicate manifest entries", async () => {
|
||||
const { root, manifest } = await fixture();
|
||||
await writeFile(join(root, "extra.js"), "untrusted");
|
||||
|
||||
@@ -29,6 +29,7 @@ function safeRelative(path: string): boolean {
|
||||
}
|
||||
// Hash streams are bounded; at most 32 descriptors/stream buffers are live.
|
||||
const PI_INVENTORY_HASH_CONCURRENCY = 32;
|
||||
const PI_INVENTORY_DISCOVERY_CONCURRENCY = 32;
|
||||
|
||||
function hash(bytes: Uint8Array | string): string { return `sha256:${createHash("sha256").update(bytes).digest("hex")}`; }
|
||||
|
||||
@@ -44,22 +45,38 @@ export async function inventoryPiRuntimeFiles(root: string): Promise<PiRuntimeFi
|
||||
const regular: Array<{ path: string; entry: PiRuntimeFile }> = [];
|
||||
const visit = async (directory: string): Promise<void> => {
|
||||
const entries = (await readdir(directory)).sort();
|
||||
for (const name of entries) {
|
||||
const path = join(directory, name);
|
||||
const rel = relative(physicalRoot, path).split(sep).join("/");
|
||||
if (!safeRelative(rel)) throw new Error("Pi runtime contains an invalid filename");
|
||||
const stat = await lstat(path);
|
||||
if (stat.isDirectory()) { await visit(path); continue; }
|
||||
if (files.length >= 100_000) throw new Error("Pi runtime file inventory exceeds its bound");
|
||||
if (stat.isSymbolicLink()) {
|
||||
const target = await readlink(path);
|
||||
if (isAbsolute(target) || !contained(physicalRoot, resolve(dirname(path), target)) || !contained(physicalRoot, await realpath(path))) throw new Error("Pi runtime link escapes its pack");
|
||||
files.push({ path: rel, kind: "symlink", target, sha256: hash(target) });
|
||||
} else if (stat.isFile()) {
|
||||
if (stat.nlink !== 1) throw new Error("Pi runtime file has another writable name");
|
||||
const entry: PiRuntimeFile = { path: rel, kind: "file", sha256: "" };
|
||||
files.push(entry); regular.push({ path, entry });
|
||||
} else throw new Error("Pi runtime contains a non-file resource");
|
||||
for (let start = 0; start < entries.length; start += PI_INVENTORY_DISCOVERY_CONCURRENCY) {
|
||||
const discovered = await Promise.allSettled(entries.slice(start, start + PI_INVENTORY_DISCOVERY_CONCURRENCY).map(async name => {
|
||||
const path = join(directory, name);
|
||||
const rel = relative(physicalRoot, path).split(sep).join("/");
|
||||
if (!safeRelative(rel)) throw new Error("Pi runtime contains an invalid filename");
|
||||
return { path, rel, stat: await lstat(path) };
|
||||
}));
|
||||
// Drain every admitted operation before failure or descent. Consuming the
|
||||
// sorted batch serially preserves the original depth-first manifest order.
|
||||
const failed = discovered.find(result => result.status === "rejected");
|
||||
if (failed?.status === "rejected") throw failed.reason;
|
||||
for (const result of discovered) {
|
||||
if (result.status !== "fulfilled") continue;
|
||||
const { path, rel, stat } = result.value;
|
||||
if (stat.isDirectory()) {
|
||||
// A previous sibling may have required a complete subtree walk since
|
||||
// this batch was inspected. Never descend using that stale identity.
|
||||
const current = await lstat(path);
|
||||
if (!current.isDirectory() || current.isSymbolicLink() || current.dev !== stat.dev || current.ino !== stat.ino || await realpath(path) !== path) throw new Error("Pi runtime directory changed before traversal");
|
||||
await visit(path); continue;
|
||||
}
|
||||
if (files.length >= 100_000) throw new Error("Pi runtime file inventory exceeds its bound");
|
||||
if (stat.isSymbolicLink()) {
|
||||
const target = await readlink(path);
|
||||
if (isAbsolute(target) || !contained(physicalRoot, resolve(dirname(path), target)) || !contained(physicalRoot, await realpath(path))) throw new Error("Pi runtime link escapes its pack");
|
||||
files.push({ path: rel, kind: "symlink", target, sha256: hash(target) });
|
||||
} else if (stat.isFile()) {
|
||||
if (stat.nlink !== 1) throw new Error("Pi runtime file has another writable name");
|
||||
const entry: PiRuntimeFile = { path: rel, kind: "file", sha256: "" };
|
||||
files.push(entry); regular.push({ path, entry });
|
||||
} else throw new Error("Pi runtime contains a non-file resource");
|
||||
}
|
||||
}
|
||||
};
|
||||
await visit(physicalRoot);
|
||||
|
||||
Reference in new issue
Block a user