mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 21:05:21 +02:00
perf(runner): keep verified snapshot copies progressing
Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
00462b048e
commit
b9e5d6ecdb
2 files changed
+113
-24
No files matched your search
@@ -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<typeof import("node:fs/promises")>();
|
||||
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<typeof import("node:fs/promises")>("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<any> => {
|
||||
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<ReturnType<typeof createNativeAcpxDistributionSnapshot>> | 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<any> {
|
||||
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<any> {
|
||||
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<typeof creating> | 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);
|
||||
|
||||
@@ -127,11 +127,13 @@ export async function createNativeAcpxDistributionSnapshot(input: NativeAcpxDist
|
||||
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");
|
||||
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<Promise<void>>();
|
||||
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];
|
||||
|
||||
Reference in new issue
Block a user