mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-09 06:15:21 +02:00
test(runner): bind remote native evidence to exact run leases
Arm bounded filesystem and process observers before publishing native actions, retain sealed file and retirement receipts before ephemeral lease teardown, and provide an independent attached-command fixture. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
d1bb5729f9
commit
9da7a7be62
2 files changed
+653
No files matched your search
@@ -0,0 +1,292 @@
|
||||
import { createHash } from "node:crypto";
|
||||
import { EventEmitter } from "node:events";
|
||||
import { mkdtemp, rename, rm, writeFile } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { Script } from "node:vm";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { bindRemoteNativeFixture, createRemoteTargetWatch, isRemoteRunRoot, parseRemoteProcStat, type RemoteNativeFixtureOptions, type RemoteNativeSnapshot } from "./remote-native-fixtures.js";
|
||||
|
||||
const hash = (s: string) => `sha256:${createHash("sha256").update(s).digest("hex")}`;
|
||||
const bootId = "12345678-1234-1234-1234-123456789abc";
|
||||
const authority = { companyId: "company", environmentId: "environment", runId: "run", leaseId: "lease", sandboxId: "sandbox", image: `registry/image@sha256:${"a".repeat(64)}` };
|
||||
const binding = { ...authority, remoteCwd: "/workspace" };
|
||||
const root = { pid: 21, ppid: 1, startTicks: "100", bootId };
|
||||
function snapshot(): RemoteNativeSnapshot {
|
||||
return { binding, observedAtMs: 1000, observedMonotonicNs: "10000", receivedAtMs: 1000, complete: true, workspace: { "existing.txt": hash("original") },
|
||||
targets: { "result.txt": { absent: true, sha256: null, parent: { dev: "1", ino: "2" }, mutationCount: 0, complete: true } },
|
||||
watcher: { complete: true, targetMutationCount: 0, workspaceMutationCount: 0 },
|
||||
processes: { captured: true, root, journal: [root], live: [21] }, setup: { path: "action.txt", sha256: null, published: false }, attached: null };
|
||||
}
|
||||
function harness() {
|
||||
let lease: Record<string, unknown> = { id: "lease", companyId: "company", environmentId: "environment", heartbeatRunId: "run", provider: "daytona", providerLeaseId: "sandbox", status: "active", releasedAt: null,
|
||||
metadata: { sandboxId: "sandbox", image: authority.image, reuseLease: false, remoteCwd: "/workspace", workspaceSentinel: { path: "/workspace/.paperclip-runtime/reusable-sandbox-lease.json", token: "fixture-sentinel-token", result: "written", runId: "run", providerLeaseId: "sandbox" } } };
|
||||
const labels: Record<string, string> = { "paperclip-provider": "daytona", "paperclip-company-id": "company", "paperclip-environment-id": "environment", "paperclip-run-id": "run", "paperclip-reuse-lease": "false" };
|
||||
const calls: Array<{ command: string; request: Record<string, any>; timeout: number | undefined }> = [];
|
||||
let resolveTerminal!: (v: unknown) => void, rejectTerminal!: (e: unknown) => void;
|
||||
const terminal = new Promise((resolve, reject) => { resolveTerminal = resolve; rejectTerminal = reject; });
|
||||
const current = snapshot();
|
||||
let override: ((request: Record<string, any>) => unknown) | undefined;
|
||||
const executeCommand = vi.fn(async (command: string, _cwd?: string, _env?: Record<string, string>, timeout?: number) => {
|
||||
const encoded = command.match(/ '([A-Za-z0-9+/=]+)'$/u)?.[1];
|
||||
if (!encoded) throw new Error("invalid command");
|
||||
const request = JSON.parse(Buffer.from(encoded, "base64").toString()); calls.push({ command, request, timeout });
|
||||
if (request.op === "wait") return { exitCode: 0, result: JSON.stringify({ ok: true, result: await terminal }) };
|
||||
if (override) { const result = override(request); if (result !== undefined) return result as { exitCode: number; result: string }; }
|
||||
let result: unknown = structuredClone(current);
|
||||
if (request.op === "publish") { current.setup = { path: request.path, published: true, sha256: hash(request.text) }; result = current.setup; }
|
||||
if (request.op === "close") result = { closed: true };
|
||||
if (request.op === "arm") result = { armed: true, sealed: false };
|
||||
if (request.op === "attached") result = { clientScript: `${request.root}/client.cjs`, clientSocket: `${request.root}/attached.sock` };
|
||||
return { exitCode: 0, result: JSON.stringify({ ok: true, result }) };
|
||||
});
|
||||
const get = vi.fn(async () => ({ id: "sandbox", labels, process: { executeCommand } }));
|
||||
const apiGet = vi.fn(async (path: string) => path.includes("/environments/") ? [structuredClone(lease)] : structuredClone(lease));
|
||||
const options: RemoteNativeFixtureOptions = { api: { get: apiGet as RemoteNativeFixtureOptions["api"]["get"] }, daytona: { get }, sdkVersion: "0.203.0", authority, nodeSha256: hash("node"), runnerdSha256: hash("runnerd"), targets: ["result.txt"], actionFile: "action.txt" };
|
||||
return { options, current, labels, calls, executeCommand, apiGet, get, resolveTerminal, rejectTerminal,
|
||||
setLease(value: Record<string, unknown>) { lease = value; }, lease: () => lease, override(fn: typeof override) { override = fn; } };
|
||||
}
|
||||
|
||||
describe("remote native lease admission", () => {
|
||||
it.each(["companyId", "environmentId", "heartbeatRunId", "providerLeaseId", "provider", "status"])("rejects wrong %s before executing any remote command", async key => {
|
||||
const h = harness(); h.setLease({ ...h.lease(), [key]: "foreign" });
|
||||
await expect(bindRemoteNativeFixture(h.options)).rejects.toThrow("lease_scope"); expect(h.executeCommand).not.toHaveBeenCalled();
|
||||
});
|
||||
it("requires exact SDK, immutable image and executable hashes", async () => {
|
||||
for (const input of [{ sdkVersion: "0.204.0" }, { nodeSha256: "unknown" }, { authority: { ...authority, image: "image:latest" } }]) {
|
||||
const h = harness(); await expect(bindRemoteNativeFixture({ ...h.options, ...input } as RemoteNativeFixtureOptions)).rejects.toThrow(); expect(h.get).not.toHaveBeenCalled();
|
||||
}
|
||||
});
|
||||
it("rejects reuse, foreign metadata and sentinel binding", async () => {
|
||||
for (const patch of [{ reuseLease: true }, { sandboxId: "other" }, { image: "image:latest" }, { remoteCwd: "/workspace/../foreign" }, { workspaceSentinel: { token: "fixture-sentinel-token" } }]) {
|
||||
const h = harness(); h.setLease({ ...h.lease(), metadata: { ...(h.lease().metadata as object), ...patch } });
|
||||
await expect(bindRemoteNativeFixture(h.options)).rejects.toThrow(); expect(h.executeCommand).not.toHaveBeenCalled();
|
||||
}
|
||||
});
|
||||
it("rejects wrong ownership labels and never discovers a sandbox by name", async () => {
|
||||
const h = harness(); h.labels["paperclip-run-id"] = "other";
|
||||
await expect(bindRemoteNativeFixture(h.options)).rejects.toThrow("sandbox_labels"); expect(h.get).toHaveBeenCalledExactlyOnceWith("sandbox"); expect(h.executeCommand).not.toHaveBeenCalled();
|
||||
});
|
||||
it("arms before publish, binds long receipt before teardown and preserves exact final bytes", async () => {
|
||||
const h = harness(), f = await bindRemoteNativeFixture(h.options);
|
||||
expect(h.calls.map(c => c.request.op)).toEqual(["install", "wait", "arm"]); expect(f.baseline.processes.live).toEqual([21]);
|
||||
expect(h.calls[1]!.timeout).toBe(295);
|
||||
await f.publishAction("action.txt", "write only result.txt");
|
||||
await expect(f.publishAction("action.txt", "retry")).rejects.toThrow("publish_bound");
|
||||
const bytes = "\nUnicode 🪴 literal \\n\n";
|
||||
const final = { ...structuredClone(h.current), processes: { ...h.current.processes, live: [] }, files: { "result.txt": Buffer.from(bytes).toString("base64") } };
|
||||
final.targets["result.txt"] = { ...final.targets["result.txt"]!, absent: false, sha256: hash(bytes), mutationCount: 1 };
|
||||
h.resolveTerminal(final); h.apiGet.mockRejectedValue(new Error("lease already deleted"));
|
||||
expect((await f.finish()).processes.live).toEqual([]); expect((await f.readFile("result.txt")).toString()).toBe(bytes);
|
||||
await f.close(); expect(h.calls.map(c => c.request.op)).toEqual(["install", "wait", "arm", "publish"]);
|
||||
});
|
||||
it("fails closed when lease deletion beats the terminal receipt", async () => {
|
||||
const h = harness(), f = await bindRemoteNativeFixture(h.options); await f.publishAction("action.txt", "task");
|
||||
h.rejectTerminal(new Error("channel closed before receipt")); await expect(f.finish()).rejects.toThrow("remote_command_failed_or_deadline");
|
||||
});
|
||||
it("rejects lease rotation before subsequent commands", async () => {
|
||||
const h = harness(), f = await bindRemoteNativeFixture(h.options); const count = h.calls.length;
|
||||
h.setLease({ ...h.lease(), heartbeatRunId: "next-run" });
|
||||
await expect(f.snapshot("pending")).rejects.toThrow("lease_scope"); expect(h.calls).toHaveLength(count);
|
||||
});
|
||||
it.each(["live", "incomplete", "missing-root", "bad-file"])("rejects %s terminal proof without weakening assertions", async type => {
|
||||
const h = harness(), f = await bindRemoteNativeFixture(h.options); await f.publishAction("action.txt", "task");
|
||||
const end: any = { ...structuredClone(h.current), processes: { ...h.current.processes, live: [] }, files: {} };
|
||||
if (type === "live") end.processes.live = [21];
|
||||
if (type === "incomplete") end.watcher.complete = false;
|
||||
if (type === "missing-root") end.processes = { captured: false, root: null, journal: [], live: [] };
|
||||
if (type === "bad-file") end.targets["result.txt"] = { ...end.targets["result.txt"], absent: false, sha256: hash("missing") };
|
||||
h.resolveTerminal(end); await expect(f.finish()).rejects.toThrow();
|
||||
});
|
||||
it("keeps a fixture-owned cross-root sentinel distinct from workspace targets", async () => {
|
||||
const h = harness(); h.options.crossRoot = { initialText: "outside sentinel" };
|
||||
h.current.targets["@cross-root"] = { absent: false, sha256: hash("outside sentinel"), parent: { dev: "1", ino: "outside" }, mutationCount: 0, complete: true };
|
||||
h.current.targets["@cross-root"]!.parent.ino = "42";
|
||||
const f = await bindRemoteNativeFixture(h.options);
|
||||
expect(f.outsideTarget).toMatch(/^\/tmp\/pc-native-[a-f0-9]{36}\/cross-root-target$/u);
|
||||
expect(f.outsideTarget!.startsWith(f.remoteCwd + "/")).toBe(false);
|
||||
expect(f.baseline.targets["@cross-root"]!.sha256).toBe(hash("outside sentinel"));
|
||||
});
|
||||
it("cleans only its admitted observer on a failed startup receipt and never publishes", async () => {
|
||||
const h = harness(); h.current.processes = { captured: false, root: null, journal: [], live: [] };
|
||||
await expect(bindRemoteNativeFixture(h.options)).rejects.toThrow("bootstrap_not_held");
|
||||
expect(h.calls.map(c => c.request.op)).toEqual(["install", "close"]);
|
||||
});
|
||||
it("refuses action publication without a confirmed receipt channel", async () => {
|
||||
const h = harness(); h.override(r => r.op === "arm" ? { exitCode: 0, result: JSON.stringify({ ok: true, result: { armed: false, sealed: false } }) } : undefined);
|
||||
await expect(bindRemoteNativeFixture(h.options)).rejects.toThrow("receipt_channel_not_armed");
|
||||
expect(h.calls.map(c => c.request.op)).toEqual(["install", "wait", "arm", "close"]);
|
||||
});
|
||||
it("bounds malformed command output and does not replay an uncertain publish", async () => {
|
||||
const h = harness(), f = await bindRemoteNativeFixture(h.options);
|
||||
h.override(r => r.op === "publish" ? { exitCode: 0, result: "x".repeat(262145) } : undefined);
|
||||
await expect(f.publishAction("action.txt", "task")).rejects.toThrow("output_bound");
|
||||
await expect(f.publishAction("action.txt", "task")).rejects.toThrow("publish_bound");
|
||||
expect(h.calls.filter(c => c.request.op === "publish")).toHaveLength(1);
|
||||
});
|
||||
it("ships syntactically valid closed Node programs with exact binary and no provider env", async () => {
|
||||
const h = harness(); await bindRemoteNativeFixture(h.options);
|
||||
const install = h.calls[0]!; const source = install.request.source as string;
|
||||
expect(() => new Script(source)).not.toThrow(); expect(source).not.toContain("__name(");
|
||||
expect(source).toContain("/proc/"); expect(source).toContain("workspaceWatch"); expect(source).toContain("finalReceipt.files");
|
||||
expect(install.command).toMatch(/^\/usr\/bin\/env -i PATH=\/usr\/bin:\/bin /u);
|
||||
expect(install.request.config.runnerdSha256).toBe(hash("runnerd"));
|
||||
expect(source).not.toMatch(/execSync|execFileSync/u);
|
||||
});
|
||||
});
|
||||
|
||||
describe("independent filesystem and process observations", () => {
|
||||
it("retains transient create/delete events and detects same-path parent replacement", async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), "remote-watch-")); const watcher = createRemoteTargetWatch(dir, "denied.txt");
|
||||
try {
|
||||
await writeFile(join(dir, "denied.txt"), "forbidden"); await rm(join(dir, "denied.txt"));
|
||||
await vi.waitFor(() => expect(watcher.snapshot().mutationCount).toBeGreaterThan(0));
|
||||
expect(watcher.snapshot().complete).toBe(true);
|
||||
await rename(dir, `${dir}-old`); await writeFile(dir, "replacement");
|
||||
expect(watcher.snapshot().complete).toBe(false);
|
||||
} finally { watcher.close(); await rm(dir, { recursive: true, force: true }); await rm(`${dir}-old`, { recursive: true, force: true }); }
|
||||
});
|
||||
it("marks lost filenames/watch errors incomplete", () => {
|
||||
let callback!: (_kind: string, filename: string | null) => void;
|
||||
const emitter = Object.assign(new EventEmitter(), { close: vi.fn() });
|
||||
const stat = { isDirectory: () => true, isSymbolicLink: () => false, dev: 1n, ino: 2n, mtimeNs: 3n, ctimeNs: 4n };
|
||||
const watcher = createRemoteTargetWatch("/test", "file", { watch: ((_path: unknown, cb: typeof callback) => { callback = cb; return emitter; }) as any, lstatSync: (() => stat) as any });
|
||||
callback("rename", null); expect(watcher.snapshot().complete).toBe(false); watcher.close();
|
||||
});
|
||||
it("parses Linux start ticks around unusual comm names and excludes shared/wrong-run daemons", () => {
|
||||
const fields = ["S", "1", "21", ...Array(16).fill("0"), "123456"];
|
||||
const p = parseRemoteProcStat(21, `21 (name with ) parens) ${fields.join(" ")}`, bootId);
|
||||
expect(p.startTicks).toBe("123456"); expect(isRemoteRunRoot(["/opt/paperclip-runnerd", "--run-id", "run", "--lifecycle-mode", "per_turn"], "run", p)).toBe(true);
|
||||
expect(isRemoteRunRoot(["/opt/paperclip-runnerd", "--run-id", "foreign", "--lifecycle-mode", "per_turn"], "run", p)).toBe(false);
|
||||
expect(isRemoteRunRoot(["/opt/paperclip-runnerd", "--run-id", "run", "--lifecycle-mode", "persistent"], "run", p)).toBe(false);
|
||||
expect(isRemoteRunRoot(["/opt/paperclip-runnerd", "--run-id", "run", "--lifecycle-mode", "per_turn"], "run", { ...p, group: 1 })).toBe(false);
|
||||
expect(parseRemoteProcStat(21, `21 (reused) ${fields.slice(0, -1).join(" ")} 999999`, bootId).startTicks).not.toBe(p.startTicks);
|
||||
});
|
||||
});
|
||||
|
||||
describe("actual generated observer state machine", () => {
|
||||
async function observerHarness() {
|
||||
const h = harness(); await bindRemoteNativeFixture(h.options);
|
||||
const { source, config } = h.calls[0]!.request;
|
||||
const intervals: Array<() => void> = [], timers: Array<{ fn: () => void; ms: number }> = [];
|
||||
const proc = new Map<number, { ppid: number; group: number; ticks: string; argv: string[] }>([[21, { ppid: 1, group: 21, ticks: "100", argv: ["/opt/bin/paperclip-runnerd", "--run-id", "run", "--environment-lease-id", "lease", "--lifecycle-mode", "per_turn"] }]]);
|
||||
const files = new Map<string, Buffer>([[`${config.root}/observer.cjs`, Buffer.from(source)], [config.sentinel.path, Buffer.from(JSON.stringify({ version: 1, provider: "daytona", token: config.sentinel.token, companyId: "company", environmentId: "environment" }))]]);
|
||||
const watches: Array<{ path: string; callback: (_kind: string, name: string | null) => void; closed: boolean }> = [];
|
||||
const handlers: Array<(socket: any) => void> = [], children: any[] = [];
|
||||
const missing = () => { throw Object.assign(new Error("missing"), { code: "ENOENT" }); };
|
||||
const fds = new Map<number, string>(); let nextFd = 50;
|
||||
const fs = {
|
||||
constants: { O_RDONLY: 0, O_NOFOLLOW: 131072 },
|
||||
openSync(path: string, flags: number) { expect(flags).toBe(131072); if (!files.has(path)) return missing(); const fd = nextFd++; fds.set(fd, path); return fd; },
|
||||
closeSync(fd: number) { fds.delete(fd); },
|
||||
fstatSync(fd: number) { const value = fs.lstatSync(fds.get(fd)!); return { ...value, size: BigInt(value.size) }; },
|
||||
readFileSync(path: string | number, encoding?: string) {
|
||||
if (typeof path === "number") path = fds.get(path)!;
|
||||
let value: Buffer | undefined;
|
||||
if (path === "/proc/sys/kernel/random/boot_id") value = Buffer.from(bootId);
|
||||
else if (path.startsWith("/proc/")) {
|
||||
const [, , raw, field] = path.split("/"), p = proc.get(Number(raw)); if (!p) return missing();
|
||||
if (field === "stat") value = Buffer.from(`${raw} (runner) S ${p.ppid} ${p.group} ${Array(16).fill("0").join(" ")} ${p.ticks}`);
|
||||
if (field === "cmdline") value = Buffer.from(p.argv.join("\0") + "\0");
|
||||
if (field === "exe") value = Buffer.from("runnerd");
|
||||
} else value = files.get(path);
|
||||
if (!value) return missing(); return encoding ? value.toString() : value;
|
||||
},
|
||||
lstatSync(path: string) {
|
||||
const directory = path === config.root || path === "/workspace";
|
||||
if (!directory && !files.has(path)) return missing();
|
||||
return { dev: 1n, ino: path === config.root ? 2n : 3n, mtimeNs: 4n, ctimeNs: 5n, isDirectory: () => directory, isFile: () => !directory, isSymbolicLink: () => false, size: files.get(path)?.length ?? 0 };
|
||||
},
|
||||
realpathSync: (path: string) => path,
|
||||
readdirSync(path: string) { if (path === "/proc") return [...proc.keys()].map(String); if (path === "/workspace") return [...files.keys()].filter(p => p.startsWith("/workspace/") && !p.slice(11).includes("/")).map(p => p.slice(11)); return []; },
|
||||
watch(path: string, options: unknown, callback?: (_kind: string, name: string | null) => void) {
|
||||
const entry = { path, callback: (callback ?? options) as (_kind: string, name: string | null) => void, closed: false }; watches.push(entry);
|
||||
return Object.assign(new EventEmitter(), { close: () => { entry.closed = true; } });
|
||||
},
|
||||
writeFileSync(path: string, content: string, opts: { flag: string }) {
|
||||
if (opts.flag === "wx" && files.has(path)) throw new Error("EEXIST"); files.set(path, Buffer.from(content));
|
||||
for (const w of watches) if (!w.closed && path.startsWith(w.path + "/")) w.callback("rename", path.slice(w.path.length + 1));
|
||||
},
|
||||
rmSync: vi.fn(),
|
||||
};
|
||||
const server = { listen: vi.fn(), close: vi.fn() };
|
||||
const net = { createServer(fn: (socket: any) => void) { handlers.push(fn); return server; } };
|
||||
const context = {
|
||||
require(name: string) { if (name === "node:fs") return fs; if (name === "node:net") return net; if (name === "node:child_process") return { spawn: vi.fn(() => { const child = Object.assign(new EventEmitter(), { pid: 88, exitCode: null, signalCode: null, kill: vi.fn() }); children.push(child); return child; }) }; if (name === "node:path") return { join: (...paths: string[]) => paths.join("/"), dirname: (path: string) => path.slice(0, path.lastIndexOf("/")), basename: (path: string) => path.slice(path.lastIndexOf("/") + 1) }; if (name === "node:crypto") return { createHash }; throw new Error("unexpected module"); },
|
||||
process: { argv: ["node", `${config.root}/observer.cjs`, Buffer.from(JSON.stringify(config)).toString("base64")], execPath: "/node", hrtime: { bigint: () => 12345n }, exit: vi.fn() },
|
||||
Buffer, __filename: `${config.root}/observer.cjs`,
|
||||
setInterval(fn: () => void) { intervals.push(fn); return 1; }, clearInterval: vi.fn(),
|
||||
setTimeout(fn: () => void, ms: number) { timers.push({ fn, ms }); return { unref() {} }; },
|
||||
};
|
||||
new Script(source).runInNewContext(context);
|
||||
function request(op: string, args: Record<string, unknown> = {}) {
|
||||
const replies: any[] = [], socket = Object.assign(new EventEmitter(), { end: (value: string) => replies.push(JSON.parse(value)), destroy: vi.fn() });
|
||||
handlers[0]!(socket); socket.emit("data", Buffer.from(JSON.stringify({ op, nonce: config.nonce, ...args }) + "\n")); return replies;
|
||||
}
|
||||
return { request, proc, files, watches, fs, intervals, timers, config, handlers, children };
|
||||
}
|
||||
it("acknowledges receipt-channel readiness only after the long waiter connects", async () => {
|
||||
const o = await observerHarness(); const arm = o.request("arm"); expect(arm).toHaveLength(0);
|
||||
o.request("wait"); expect(arm).toEqual([{ ok: true, result: { armed: true, sealed: false } }]);
|
||||
o.proc.set(1, { ppid: 0, group: 1, ticks: "1", argv: ["/sbin/init"] });
|
||||
expect(o.request("snapshot")[0].result.processes.captured).toBe(true);
|
||||
});
|
||||
it("arms exact remote root then seals drained no-live proof and retained bytes", async () => {
|
||||
const o = await observerHarness(); expect(o.request("snapshot")[0].result.processes.captured).toBe(true);
|
||||
const wait = o.request("wait"); expect(wait).toHaveLength(0);
|
||||
expect(o.request("publish", { path: "action.txt", text: "do task" })[0].ok).toBe(true);
|
||||
o.fs.writeFileSync("/workspace/result.txt", "exact 🪴\n", { flag: "wx" });
|
||||
o.proc.set(22, { ppid: 21, group: 21, ticks: "200", argv: ["/node"] }); o.intervals[0]!();
|
||||
o.proc.clear(); o.intervals[0]!(); expect(wait).toHaveLength(0); // Watch drain is mandatory.
|
||||
o.timers.find(t => t.ms === 100)!.fn();
|
||||
const final = wait[0].result;
|
||||
expect(final.complete).toBe(true); expect(final.processes.live).toEqual([]); expect(final.processes.journal.map((p: any) => p.pid)).toEqual([21, 22]);
|
||||
expect(Buffer.from(final.files["result.txt"], "base64").toString()).toBe("exact 🪴\n");
|
||||
expect(final.watcher.workspaceMutationCount).toBe(1); expect(final.workspace["action.txt"]).toBeUndefined(); expect(final.setup.sha256).toBe(hash("do task"));
|
||||
expect(o.watches.every(w => w.closed)).toBe(true);
|
||||
});
|
||||
it("links the attached client to the exact run and writes the marker only after independent child exit", async () => {
|
||||
const o = await observerHarness(); o.request("snapshot");
|
||||
const config = { marker: "result.txt", markerText: "settled", delayMs: 500, clientNonce: "test-client-nonce" };
|
||||
const fixture = o.request("attached", config)[0].result;
|
||||
o.proc.set(23, { ppid: 21, group: 21, ticks: "300", argv: ["/node", fixture.clientScript, fixture.clientSocket, config.clientNonce] });
|
||||
const replies: any[] = [], socket = Object.assign(new EventEmitter(), { end: (value: string) => replies.push(JSON.parse(value)), destroy: vi.fn() });
|
||||
o.handlers[1]!(socket); socket.emit("data", Buffer.from(JSON.stringify({ nonce: config.clientNonce, pid: 23 }) + "\n"));
|
||||
expect(o.children).toHaveLength(1); expect(replies).toHaveLength(0); expect(o.files.has("/workspace/result.txt")).toBe(false);
|
||||
o.children[0].exitCode = 0; o.children[0].emit("exit", 0, null);
|
||||
expect(replies).toEqual([{ code: 0 }]); expect(o.files.get("/workspace/result.txt")?.toString()).toBe("settled");
|
||||
o.proc.delete(23); const a = o.request("snapshot")[0].result.attached;
|
||||
expect(a.commandExit.code).toBe(0); expect(BigInt(a.commandExit.observedMonotonicNs)).toBeLessThanOrEqual(BigInt(a.markerWrittenMonotonicNs));
|
||||
expect(a.clientExitedAtMs).not.toBeNull(); expect(a.connections).toBe(1);
|
||||
});
|
||||
it("never starts attached work for a same-argv client outside the captured run", async () => {
|
||||
const o = await observerHarness(); o.request("snapshot");
|
||||
const config = { marker: "result.txt", markerText: "settled", delayMs: 500, clientNonce: "test-client-nonce" };
|
||||
const fixture = o.request("attached", config)[0].result;
|
||||
o.proc.set(23, { ppid: 1, group: 23, ticks: "300", argv: ["/node", fixture.clientScript, fixture.clientSocket, config.clientNonce] });
|
||||
const socket = Object.assign(new EventEmitter(), { end: vi.fn(), destroy: vi.fn() });
|
||||
o.handlers[1]!(socket); socket.emit("data", Buffer.from(JSON.stringify({ nonce: config.clientNonce, pid: 23 }) + "\n"));
|
||||
expect(o.children).toHaveLength(0); expect(socket.destroy).toHaveBeenCalled();
|
||||
expect(o.request("snapshot")[0].result.attached.failure).toBe("client_rejected");
|
||||
});
|
||||
it("marks PID reuse incomplete rather than mistaking a new process for retired authority", async () => {
|
||||
const o = await observerHarness(); o.request("snapshot"); const wait = o.request("wait");
|
||||
o.proc.set(21, { ppid: 1, group: 99, ticks: "999", argv: ["/unrelated"] }); o.intervals[0]!(); o.timers.find(t => t.ms === 100)!.fn();
|
||||
expect(wait[0].result.complete).toBe(false); expect(wait[0].result.processes.root.startTicks).toBe("100");
|
||||
});
|
||||
it("counts transient workspace create/delete and rejects changed setup-file bytes", async () => {
|
||||
const o = await observerHarness(); o.request("snapshot");
|
||||
o.fs.writeFileSync("/workspace/transient.txt", "not allowed", { flag: "wx" }); o.files.delete("/workspace/transient.txt");
|
||||
for (const w of o.watches) w.callback("rename", "transient.txt");
|
||||
const snapshot = o.request("snapshot")[0].result;
|
||||
expect(snapshot.workspace["transient.txt"]).toBeUndefined(); expect(snapshot.watcher.workspaceMutationCount).toBe(2);
|
||||
o.request("publish", { path: "action.txt", text: "approved" }); o.files.set("/workspace/action.txt", Buffer.from("replacement"));
|
||||
expect(o.request("snapshot")[0].ok).toBe(false);
|
||||
});
|
||||
it("rejects ambiguous run roots and mismatched lease flags", async () => {
|
||||
const o = await observerHarness(); const p = o.proc.get(21)!;
|
||||
o.proc.set(22, { ...p, group: 22, ticks: "200" }); expect(o.request("snapshot")[0].result.complete).toBe(false);
|
||||
const other = await observerHarness(); other.proc.get(21)!.argv[4] = "foreign-lease";
|
||||
expect(other.request("snapshot")[0].ok).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,361 @@
|
||||
import { createHash, randomBytes } from "node:crypto";
|
||||
import { watch, lstatSync } from "node:fs";
|
||||
import { posix } from "node:path";
|
||||
|
||||
export const REMOTE_FIXTURE_DAYTONA_SDK_VERSION = "0.203.0";
|
||||
const NODE = "/opt/paperclip-runner/provider-pack/node_modules/node/bin/node";
|
||||
const MAX_OUTPUT = 256 * 1024;
|
||||
const digest = (value: string) => `sha256:${createHash("sha256").update(value).digest("hex")}`;
|
||||
const record = (value: unknown): Record<string, unknown> => value !== null && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : {};
|
||||
const fail = (condition: unknown, code: string): void => { if (!condition) throw new Error(`remote_native_fixture:${code}`); };
|
||||
const quote = (value: string) => `'${value.replaceAll("'", "'\\''")}'`;
|
||||
const id = (value: string) => /^[a-zA-Z0-9_-]{1,128}$/u.test(value);
|
||||
function relative(value: string): string {
|
||||
fail(value.length <= 240 && /^[a-zA-Z0-9._/-]+$/u.test(value) && !value.startsWith("/") && value.split("/").every(part => part !== "" && part !== "." && part !== ".."), "unsafe_relative_target");
|
||||
return value;
|
||||
}
|
||||
export interface RemoteNativeAuthority { companyId: string; environmentId: string; runId: string; leaseId: string; sandboxId: string; image: string }
|
||||
export interface RemoteNativeBinding extends RemoteNativeAuthority { remoteCwd: string }
|
||||
export interface RemoteProcessIdentity { pid: number; ppid: number; startTicks: string; bootId: string }
|
||||
export interface RemoteNativeSnapshot {
|
||||
binding: RemoteNativeBinding;
|
||||
observedAtMs: number;
|
||||
observedMonotonicNs: string;
|
||||
receivedAtMs: number;
|
||||
complete: boolean;
|
||||
workspace: Record<string, string>;
|
||||
targets: Record<string, { absent: boolean; sha256: string | null; parent: { dev: string; ino: string }; mutationCount: number; complete: boolean }>;
|
||||
watcher: { complete: boolean; targetMutationCount: number; workspaceMutationCount: number };
|
||||
processes: { captured: boolean; root: RemoteProcessIdentity | null; journal: RemoteProcessIdentity[]; live: number[] };
|
||||
setup: { path: string; sha256: string | null; published: boolean };
|
||||
attached: { connections: number; failure: string | null; commandExit: { code: number; observedAtMs: number; observedMonotonicNs: string } | null; markerWrittenAtMs: number | null; markerWrittenMonotonicNs: string | null; clientExitedAtMs: number | null; clientExitedMonotonicNs: string | null } | null;
|
||||
}
|
||||
/** Structural subset of pinned SDK0.203.0; caller supplies its authenticated
|
||||
* client. No create/list/delete or arbitrary remote command API is exposed. */
|
||||
export interface RemoteFixtureDaytona {
|
||||
get(id: string): Promise<{ id: string; labels?: Record<string, string>; process: {
|
||||
executeCommand(command: string, cwd?: string, env?: Record<string, string>, timeout?: number): Promise<{ exitCode: number; result: string }>;
|
||||
} }>;
|
||||
}
|
||||
export interface RemoteFixtureApi { get<T>(path: string): Promise<T> }
|
||||
|
||||
/** Linux /proc identity uses boot ID + start ticks, never PID alone. */
|
||||
export function parseRemoteProcStat(pid: number, stat: string, bootId: string): RemoteProcessIdentity & { group: number; state: string } {
|
||||
const close = stat.lastIndexOf(") ");
|
||||
if (!Number.isSafeInteger(pid) || pid < 2 || close < 0 || !stat.startsWith(`${pid} (`) || !/^[a-f0-9-]{36}$/u.test(bootId)) throw new Error("invalid_proc_identity");
|
||||
const fields = stat.slice(close + 2).trim().split(/\s+/u);
|
||||
if (!/^\d+$/u.test(fields[19] ?? "") || !/^\d+$/u.test(fields[1] ?? "") || !/^\d+$/u.test(fields[2] ?? "")) throw new Error("invalid_proc_fields");
|
||||
return { pid, ppid: Number(fields[1]), group: Number(fields[2]), state: fields[0]!, startTicks: fields[19]!, bootId };
|
||||
}
|
||||
export function isRemoteRunRoot(argv: string[], runId: string, process: RemoteProcessIdentity & { group: number }): boolean {
|
||||
if (!/^[a-zA-Z0-9_-]{1,128}$/u.test(runId) || process.group !== process.pid || !argv[0]?.endsWith("/paperclip-runnerd")) return false;
|
||||
return [["--run-id", runId], ["--lifecycle-mode", "per_turn"]].every(([flag, value]) => {
|
||||
const at = argv.indexOf(flag!); return at > 0 && argv.lastIndexOf(flag!) === at && argv[at + 1] === value;
|
||||
});
|
||||
}
|
||||
/** Shared by the real remote observer and local deterministic tests. */
|
||||
export function createRemoteTargetWatch(directory: string, name: string, io = { watch, lstatSync }) {
|
||||
const before = io.lstatSync(directory, { bigint: true });
|
||||
if (!before.isDirectory() || before.isSymbolicLink()) throw new Error("unsafe_watch_parent");
|
||||
let complete = true, mutations = 0, events = 0, sealed = false;
|
||||
const watcher = io.watch(directory, (_kind, filename) => {
|
||||
if (sealed) return;
|
||||
events++;
|
||||
if (filename === null || events > 4096) complete = false;
|
||||
else if (String(filename) === name) mutations++;
|
||||
});
|
||||
watcher.on("error", () => { complete = false; });
|
||||
return {
|
||||
snapshot() {
|
||||
let deliveryComplete = true;
|
||||
try {
|
||||
const after = io.lstatSync(directory, { bigint: true });
|
||||
if (!after.isDirectory() || after.isSymbolicLink() || after.dev !== before.dev || after.ino !== before.ino) complete = false;
|
||||
if ((after.ctimeNs !== before.ctimeNs || after.mtimeNs !== before.mtimeNs) && events === 0) deliveryComplete = false;
|
||||
} catch { complete = false; }
|
||||
return { complete: complete && deliveryComplete, mutationCount: mutations, parent: { dev: String(before.dev), ino: String(before.ino) } };
|
||||
},
|
||||
close() { this.snapshot(); sealed = true; watcher.close(); },
|
||||
};
|
||||
}
|
||||
|
||||
// Only this closed Node program executes through the SDK. The observer has no
|
||||
// general command endpoint. Its control nonce is never given to the provider.
|
||||
const OBSERVER = String.raw`
|
||||
const fs=require('node:fs'),path=require('node:path'),crypto=require('node:crypto'),net=require('node:net'),cp=require('node:child_process');
|
||||
const config=JSON.parse(Buffer.from(process.argv[2],'base64').toString());
|
||||
const parseStat=PARSE_STAT;const runRoot=RUN_ROOT;const watchTarget=WATCH_TARGET;
|
||||
const hash=x=>'sha256:'+crypto.createHash('sha256').update(x).digest('hex');
|
||||
const cwdStat=fs.lstatSync(config.binding.remoteCwd,{bigint:true}),rootStat=fs.lstatSync(config.root,{bigint:true}),scriptStat=fs.lstatSync(__filename,{bigint:true}),scriptHash=hash(fs.readFileSync(__filename));
|
||||
const boot=fs.readFileSync('/proc/sys/kernel/random/boot_id','utf8').trim();
|
||||
const targets=new Map();let complete=true,sealed=false,root=null,attached=null,child=null,client=null,childTimer=null;
|
||||
const journal=new Map(),sockets=new Set(),waiters=new Set(),armWaiters=new Set();let observedRootCount=0,publishedHash=null,finalReceipt=null,retiring=false;
|
||||
function identity(pid){return parseStat(pid,fs.readFileSync('/proc/'+pid+'/stat','utf8'),boot)}
|
||||
function table(){const ids=fs.readdirSync('/proc').filter(x=>/^\d+$/.test(x)&&Number(x)>1);if(ids.length>4096)throw Error('process_bound');return ids.flatMap(x=>{try{return [identity(Number(x))]}catch(e){if(e.code==='ENOENT'||e.code==='ESRCH')return [];throw e}})}
|
||||
function sample(){
|
||||
const all=table();const candidates=[];
|
||||
for(const p of all){if(p.state==='Z')continue;let argv;try{argv=fs.readFileSync('/proc/'+p.pid+'/cmdline').toString().split('\0').filter(Boolean)}catch(e){if(e.code==='ENOENT'||e.code==='ESRCH')continue;throw e}
|
||||
if(runRoot(argv,config.binding.runId,p)){const at=argv.indexOf('--environment-lease-id');if(at<1||argv.lastIndexOf('--environment-lease-id')!==at||argv[at+1]!==config.binding.leaseId)throw Error('process_lease_identity');candidates.push(p);}}
|
||||
if(candidates.length>1)complete=false;
|
||||
if(!root&&candidates.length===1){if(hash(fs.readFileSync('/proc/'+candidates[0].pid+'/exe'))!==config.runnerdSha256)throw Error('runner_binary_identity');root=candidates[0];journal.set(root.pid,root);observedRootCount++;}
|
||||
if(root&&candidates.some(p=>p.pid!==root.pid||p.startTicks!==root.startTicks))complete=false;
|
||||
let changed=true;while(changed){changed=false;for(const p of all){const parent=journal.get(p.ppid);if(!journal.has(p.pid)&&((parent&&all.some(q=>q.pid===parent.pid&&q.startTicks===parent.startTicks))||(root&&p.group===root.pid))){if(journal.size>=512)throw Error('journal_bound');journal.set(p.pid,p);changed=true;}}}
|
||||
if(all.some(p=>journal.has(p.pid)&&journal.get(p.pid).startTicks!==p.startTicks))complete=false;
|
||||
const live=all.filter(p=>journal.get(p.pid)?.startTicks===p.startTicks&&p.state!=='Z').map(p=>p.pid);
|
||||
if(client&&!all.some(p=>p.pid===client.pid&&p.startTicks===client.startTicks&&p.state!=='Z')&&attached&&!attached.clientExitedAtMs){attached.clientExitedAtMs=Date.now();attached.clientExitedMonotonicNs=process.hrtime.bigint().toString();}
|
||||
return {captured:root!==null,root,journal:[...journal.values()],live};
|
||||
}
|
||||
function guard(){const s=fs.lstatSync(config.root,{bigint:true});if(s.dev!==rootStat.dev||s.ino!==rootStat.ino||!s.isDirectory()||s.isSymbolicLink()||hash(fs.readFileSync(__filename))!==scriptHash||fs.lstatSync(__filename,{bigint:true}).ino!==scriptStat.ino)throw Error('observer_identity_changed');
|
||||
const st=fs.lstatSync(config.sentinel.path);if(!st.isFile()||st.isSymbolicLink()||st.size>16384||fs.realpathSync(config.sentinel.path)!==config.sentinel.path)throw Error('sentinel_type');const sentinel=JSON.parse(fs.readFileSync(config.sentinel.path,'utf8'));if(sentinel.version!==1||sentinel.provider!=='daytona'||sentinel.token!==config.sentinel.token||sentinel.companyId!==config.binding.companyId||sentinel.environmentId!==config.binding.environmentId)throw Error('sentinel_changed');
|
||||
const c=fs.lstatSync(config.binding.remoteCwd,{bigint:true});if(c.dev!==cwdStat.dev||c.ino!==cwdStat.ino||!c.isDirectory()||c.isSymbolicLink()||fs.realpathSync(config.binding.remoteCwd)!==config.binding.remoteCwd)throw Error('workspace_replaced');}
|
||||
function readSafe(p){const fd=fs.openSync(p,fs.constants.O_RDONLY|fs.constants.O_NOFOLLOW);try{const before=fs.fstatSync(fd,{bigint:true});if(!before.isFile()||before.size>65536n)throw Error('file_bound_or_type');const bytes=fs.readFileSync(fd),after=fs.fstatSync(fd,{bigint:true}),named=fs.lstatSync(p,{bigint:true});if(before.dev!==named.dev||before.ino!==named.ino||named.isSymbolicLink()||before.size!==after.size||before.mtimeNs!==after.mtimeNs||BigInt(bytes.length)!==before.size)throw Error('file_changed');return bytes}finally{fs.closeSync(fd)}}
|
||||
function file(p){try{const s=fs.lstatSync(p);if(s.isSymbolicLink()||!s.isFile()||s.size>65536)throw Error('file_bound_or_type');return {absent:false,sha256:hash(readSafe(p))}}catch(e){if(e.code==='ENOENT')return {absent:true,sha256:null};throw e}}
|
||||
function workspace(){const result={};let count=0,bytes=0;function visit(dir,prefix){for(const name of fs.readdirSync(dir).sort()){if(++count>512)throw Error('workspace_entry_bound');const full=path.join(dir,name),rel=prefix+name,s=fs.lstatSync(full);if(rel===config.actionFile){if(!publishedHash||file(full).sha256!==publishedHash)throw Error('setup_file_changed');continue;}if(s.isSymbolicLink())throw Error('workspace_symlink');if(s.isDirectory()){result[rel]='directory';visit(full,rel+'/')}else if(s.isFile()){bytes+=s.size;if(s.size>65536||bytes>4194304)throw Error('workspace_byte_bound');result[rel]=hash(readSafe(full))}else throw Error('workspace_special_file')}}visit(config.binding.remoteCwd,'');return result}
|
||||
function snapshot(){guard();const processes=sample(),out={};let watchComplete=complete,total=0;for(const [name,t] of targets){const status=t.watch.snapshot();out[name]={...file(t.path),...status};watchComplete&&=status.complete;total+=status.mutationCount}return {binding:config.binding,observedAtMs:Date.now(),observedMonotonicNs:process.hrtime.bigint().toString(),complete:complete&&watchComplete,workspace:workspace(),targets:out,watcher:{complete:watchComplete,targetMutationCount:total,workspaceMutationCount},processes,setup:{path:config.actionFile,sha256:publishedHash,published:publishedHash!==null},attached}}
|
||||
for(const name of config.targets){const p=path.join(config.binding.remoteCwd,name);if(fs.realpathSync(path.dirname(p))!==path.dirname(p))throw Error('target_parent_symlink');targets.set(name,{path:p,watch:watchTarget(path.dirname(p),path.basename(p),{watch:fs.watch,lstatSync:fs.lstatSync})})}
|
||||
if(config.crossRoot){const p=path.join(config.root,'cross-root-target');fs.writeFileSync(p,config.crossRoot.initialText,{flag:'wx',mode:0o600});targets.set('@cross-root',{path:p,watch:watchTarget(config.root,'cross-root-target',{watch:fs.watch,lstatSync:fs.lstatSync})})}
|
||||
let workspaceMutationCount=0;const workspaceWatch=fs.watch(config.binding.remoteCwd,{recursive:true},(_kind,name)=>{if(name!==null&&String(name)===config.actionFile){try{if(!publishedHash||file(path.join(config.binding.remoteCwd,config.actionFile)).sha256!==publishedHash)complete=false}catch{complete=false}return}workspaceMutationCount++;if(name===null||workspaceMutationCount>4096)complete=false;});workspaceWatch.on('error',()=>{complete=false});
|
||||
function publicSnapshot(){const result=snapshot();if(result.attached)result.attached={connections:attached.connections,failure:attached.failure,commandExit:attached.commandExit,markerWrittenAtMs:attached.markerWrittenAtMs,markerWrittenMonotonicNs:attached.markerWrittenMonotonicNs,clientExitedAtMs:attached.clientExitedAtMs,clientExitedMonotonicNs:attached.clientExitedMonotonicNs};return result}
|
||||
function seal(){if(sealed)return;sealed=true;workspaceWatch.close();for(const t of targets.values())t.watch.close();clearInterval(observer);finalReceipt=publicSnapshot();finalReceipt.files={};for(const [name,t] of targets){if(!file(t.path).absent)finalReceipt.files[name]=readSafe(t.path).toString('base64');}if(Buffer.byteLength(JSON.stringify(finalReceipt))>250000){finalReceipt.complete=false;finalReceipt.files={};}if(child&&child.exitCode===null&&child.signalCode===null)finalReceipt.complete=false;for(const socket of waiters)socket.end(JSON.stringify({ok:true,result:finalReceipt})+'\n');waiters.clear();setTimeout(()=>shutdown(null),250);}
|
||||
const observer=setInterval(()=>{try{const p=sample();if(p.captured&&p.live.length===0&&!retiring){retiring=true;setTimeout(()=>{try{const end=sample();if(end.live.length===0)seal();else retiring=false}catch{complete=false}},100)}}catch{complete=false}},25);
|
||||
const server=net.createServer(socket=>{sockets.add(socket);socket.on('close',()=>sockets.delete(socket));socket.on('error',()=>{});let buffer='';socket.on('data',data=>{buffer+=data;if(Buffer.byteLength(buffer)>8192){socket.destroy();complete=false;return}if(!buffer.includes('\n'))return;socket.removeAllListeners('data');try{const r=JSON.parse(buffer);if(r.nonce!==config.nonce)throw Error('control_identity');guard();let result;
|
||||
if(r.op==='snapshot')result=finalReceipt??publicSnapshot();
|
||||
else if(r.op==='wait'){if(finalReceipt)result=finalReceipt;else{if(waiters.size)throw Error('duplicate_receipt_channel');waiters.add(socket);for(const arm of armWaiters)arm.end(JSON.stringify({ok:true,result:{armed:true,sealed:false}})+'\n');armWaiters.clear();socket.on('close',()=>waiters.delete(socket));return}}
|
||||
else if(r.op==='arm'){if(waiters.size)result={armed:true,sealed};else{armWaiters.add(socket);socket.on('close',()=>armWaiters.delete(socket));return}}
|
||||
else if(r.op==='publish'){if(sealed||publishedHash||r.path!==config.actionFile||typeof r.text!=='string'||Buffer.byteLength(r.text)>16384)throw Error('publish_bound');const p=path.join(config.binding.remoteCwd,config.actionFile);if(fs.realpathSync(path.dirname(p))!==path.dirname(p))throw Error('publish_parent');publishedHash=hash(r.text);try{fs.writeFileSync(p,r.text,{flag:'wx',mode:0o600})}catch(e){publishedHash=null;throw e}result={path:r.path,sha256:publishedHash,published:true};}
|
||||
else if(r.op==='read'){const p=r.path==='@cross-root'&&config.crossRoot?path.join(config.root,'cross-root-target'):path.join(config.binding.remoteCwd,r.path);if(r.path!=='@cross-root'&&!config.targets.includes(r.path))throw Error('unregistered_read');const status=file(p);if(status.absent)throw Error('file_missing');result={...status,base64:readSafe(p).toString('base64')};}
|
||||
else if(r.op==='attached'){if(attached||sealed)throw Error('attached_already_configured');if(!config.targets.includes(r.marker)||!Number.isInteger(r.delayMs)||r.delayMs<100||r.delayMs>8000)throw Error('attached_bounds');
|
||||
attached={connections:0,failure:null,commandExit:null,markerWrittenAtMs:null,markerWrittenMonotonicNs:null,clientExitedAtMs:null,clientExitedMonotonicNs:null};
|
||||
const clientScript=path.join(config.root,'client.cjs'),clientSocket=path.join(config.root,'attached.sock');
|
||||
fs.writeFileSync(clientScript,ATTACHED_CLIENT,{flag:'wx',mode:0o400});
|
||||
const srv=net.createServer(s=>{sockets.add(s);s.on('close',()=>sockets.delete(s));s.on('error',()=>{});let b='';s.on('data',d=>{b+=d;if(b.length>1024){s.destroy();attached.failure='client_bound';return}if(!b.includes('\n'))return;s.removeAllListeners('data');try{const q=JSON.parse(b);attached.connections++;if(q.nonce!==r.clientNonce||attached.connections!==1)throw Error('client_identity');client=identity(q.pid);sample();if(journal.get(q.pid)?.startTicks!==client.startTicks)throw Error('client_not_owned_by_run');const argv=fs.readFileSync('/proc/'+q.pid+'/cmdline').toString().split('\0');if(argv[1]!==clientScript||argv[2]!==clientSocket||argv[3]!==r.clientNonce)throw Error('client_argv');
|
||||
child=cp.spawn(process.execPath,['-e','setTimeout(()=>process.exit(0),'+r.delayMs+')'],{env:{PATH:'/usr/bin:/bin'},stdio:'ignore'});child.once('error',()=>{attached.failure='child_start';s.destroy()});child.once('exit',(code,signal)=>{attached.commandExit={code:code??-1,observedAtMs:Date.now(),observedMonotonicNs:process.hrtime.bigint().toString()};if(code!==0||signal){attached.failure='child_failed';s.destroy();return}try{fs.writeFileSync(path.join(config.binding.remoteCwd,r.marker),r.markerText,{flag:'wx'});attached.markerWrittenAtMs=Date.now();attached.markerWrittenMonotonicNs=process.hrtime.bigint().toString();s.end(JSON.stringify({code:0})+'\n')}catch{attached.failure='marker_failed';s.destroy()}});
|
||||
}catch{attached.failure='client_rejected';s.destroy()}})});srv.listen(clientSocket);attached.server=srv;result={clientScript,clientSocket};
|
||||
}
|
||||
else if(r.op==='finish'){result=snapshot();if(!result.complete||!result.processes.captured||result.processes.live.length)throw Error('retirement_unproven');sealed=true;workspaceWatch.close();for(const t of targets.values())t.watch.close();}
|
||||
else if(r.op==='close'){shutdown(socket);return;}
|
||||
else throw Error('unknown_operation');
|
||||
// Socket handles never cross the evidence boundary.
|
||||
if(result?.attached)result={...result,attached:{connections:attached.connections,failure:attached.failure,commandExit:attached.commandExit,markerWrittenAtMs:attached.markerWrittenAtMs,markerWrittenMonotonicNs:attached.markerWrittenMonotonicNs,clientExitedAtMs:attached.clientExitedAtMs,clientExitedMonotonicNs:attached.clientExitedMonotonicNs}};
|
||||
socket.end(JSON.stringify({ok:true,result})+'\n');
|
||||
}catch{socket.end(JSON.stringify({ok:false,error:'observer_evidence_incomplete'})+'\n')}})});
|
||||
let closing=false;async function shutdown(reply){if(closing){reply?.end(JSON.stringify({ok:false,error:'closing'})+'\n');return}closing=true;sealed=true;workspaceWatch.close();for(const t of targets.values())t.watch.close();clearInterval(observer);let settled=true;try{if(child&&child.exitCode===null&&child.signalCode===null){child.kill('SIGTERM');const exited=new Promise(resolve=>child.once('exit',resolve));await Promise.race([exited,new Promise(resolve=>setTimeout(resolve,2000))]);if(child.exitCode===null&&child.signalCode===null){child.kill('SIGKILL');await Promise.race([exited,new Promise(resolve=>setTimeout(resolve,2000))]);}settled=child.exitCode!==null||child.signalCode!==null;}if(reply)reply.end(JSON.stringify({ok:settled,result:{closed:settled}})+'\n');}finally{for(const s of sockets)if(s!==reply)s.destroy();server.close();if(attached?.server)attached.server.close();setTimeout(()=>{reply?.destroy();if(settled){const current=fs.lstatSync(config.root,{bigint:true});if(current.dev===rootStat.dev&¤t.ino===rootStat.ino&&!current.isSymbolicLink())fs.rmSync(config.root,{recursive:true,force:true});else settled=false;}process.exit(settled?0:2)},100)}}
|
||||
server.listen(path.join(config.root,'control.sock'));setTimeout(()=>{complete=false;shutdown(null)},300000).unref();
|
||||
`;
|
||||
const ATTACHED_CLIENT = String.raw`const net=require('node:net');const s=net.connect(process.argv[2]);let b='';s.setTimeout(15000,()=>process.exit(3));s.on('error',()=>process.exit(4));s.on('connect',()=>s.write(JSON.stringify({nonce:process.argv[3],pid:process.pid})+'\n'));s.on('data',x=>{b+=x;if(b.length>1024)process.exit(6);if(b.includes('\n')){const r=JSON.parse(b);s.end();process.exit(r.code===0?0:5)}});`;
|
||||
function observerSource() {
|
||||
return OBSERVER.replace("PARSE_STAT", () => parseRemoteProcStat.toString()).replace("RUN_ROOT", () => isRemoteRunRoot.toString())
|
||||
.replace("WATCH_TARGET", () => createRemoteTargetWatch.toString()).replace("ATTACHED_CLIENT", () => JSON.stringify(ATTACHED_CLIENT));
|
||||
}
|
||||
const RPC = String.raw`const fs=require('node:fs'),net=require('node:net'),cp=require('node:child_process'),crypto=require('node:crypto');const r=JSON.parse(Buffer.from(process.argv[1],'base64').toString());const hash=x=>'sha256:'+crypto.createHash('sha256').update(x).digest('hex');
|
||||
(async()=>{if(hash(fs.readFileSync(process.execPath))!==r.nodeSha256)throw Error('node_identity');if(r.op==='install'){const c=r.config;const st=fs.lstatSync(c.sentinel.path);if(!st.isFile()||st.isSymbolicLink()||st.size>16384||fs.realpathSync(c.sentinel.path)!==c.sentinel.path)throw Error('sentinel_type');const s=JSON.parse(fs.readFileSync(c.sentinel.path,'utf8'));if(s.version!==1||s.provider!=='daytona'||s.token!==c.sentinel.token||s.companyId!==c.binding.companyId||s.environmentId!==c.binding.environmentId)throw Error('sentinel');if(fs.realpathSync(c.binding.remoteCwd)!==c.binding.remoteCwd)throw Error('cwd');fs.mkdirSync(c.root,{mode:0o700});fs.writeFileSync(c.root+'/observer.cjs',r.source,{flag:'wx',mode:0o400});const child=cp.spawn(process.execPath,[c.root+'/observer.cjs',Buffer.from(JSON.stringify(c)).toString('base64')],{detached:true,stdio:'ignore',env:{PATH:'/usr/bin:/bin'}});child.unref();r.root=c.root;r.nonce=c.nonce;r.op='snapshot';}
|
||||
for(let i=0;!fs.existsSync(r.root+'/control.sock')&&i<200;i++)await new Promise(resolve=>setTimeout(resolve,10));const socket=net.connect(r.root+'/control.sock');let output='';socket.setTimeout(r.op==='wait'?290000:5000);socket.on('timeout',()=>{socket.destroy();process.exitCode=2});socket.on('error',()=>{process.exitCode=2});socket.on('connect',()=>socket.write(JSON.stringify(r)+'\n'));socket.on('data',b=>{output+=b;if(Buffer.byteLength(output)>262144){socket.destroy();process.exitCode=2}});socket.on('end',()=>{if(!process.exitCode)process.stdout.write(output)});
|
||||
})().catch(()=>{process.exitCode=2});`;
|
||||
|
||||
export interface RemoteNativeFixture {
|
||||
readonly binding: RemoteNativeBinding;
|
||||
readonly remoteCwd: string;
|
||||
readonly actionFile: string;
|
||||
readonly outsideTarget: string | null;
|
||||
/** Watchers and exact run-root identity are already armed at return from bind. */
|
||||
readonly baseline: RemoteNativeSnapshot;
|
||||
snapshot(label: string): Promise<RemoteNativeSnapshot>;
|
||||
readFile(path: string): Promise<Buffer>;
|
||||
publishAction(path: string, text: string): Promise<void>;
|
||||
setupAttachedCommand(input: { marker: string; markerText: string; delayMs: number }): Promise<{ command: string; commandSha256: string }>;
|
||||
/** Consumes the host-held terminal receipt; never queries a deleted lease. */
|
||||
finish(): Promise<RemoteNativeSnapshot>;
|
||||
/** Only fixture-owned observer cleanup. Never deletes a sandbox or lease. */
|
||||
close(): Promise<void>;
|
||||
}
|
||||
export interface RemoteNativeFixtureOptions {
|
||||
api: RemoteFixtureApi;
|
||||
daytona: RemoteFixtureDaytona;
|
||||
sdkVersion: typeof REMOTE_FIXTURE_DAYTONA_SDK_VERSION;
|
||||
authority: RemoteNativeAuthority;
|
||||
nodeSha256: string;
|
||||
runnerdSha256: string;
|
||||
targets: string[];
|
||||
actionFile: string;
|
||||
crossRoot?: { initialText: string };
|
||||
}
|
||||
const sha = (value: unknown): value is string => typeof value === "string" && /^sha256:[a-f0-9]{64}$/u.test(value);
|
||||
function readSnapshot(value: unknown, binding: RemoteNativeBinding, names: string[], actionFile: string): RemoteNativeSnapshot {
|
||||
const row = record(value), watcher = record(row.watcher), processes = record(row.processes), setup = record(row.setup);
|
||||
fail(JSON.stringify(row.binding) === JSON.stringify(binding), "receipt_binding");
|
||||
fail(Number.isSafeInteger(row.observedAtMs) && (row.observedAtMs as number) > 0 && typeof row.observedMonotonicNs === "string" && /^\d+$/u.test(row.observedMonotonicNs) && typeof row.complete === "boolean", "receipt_shape");
|
||||
fail(Object.keys(record(row.workspace)).length <= 512 && Object.entries(record(row.workspace)).every(([path, hash]) => {
|
||||
try { relative(path); return hash === "directory" || sha(hash); } catch { return false; }
|
||||
}), "workspace_shape");
|
||||
fail(Object.keys(record(row.targets)).sort().join("\0") === [...names].sort().join("\0"), "target_set");
|
||||
for (const target of Object.values(record(row.targets))) {
|
||||
const t = record(target), parent = record(t.parent);
|
||||
fail(typeof t.absent === "boolean" && (t.absent ? t.sha256 === null : sha(t.sha256)) && typeof t.complete === "boolean"
|
||||
&& Number.isSafeInteger(t.mutationCount) && (t.mutationCount as number) >= 0
|
||||
&& /^\d+$/u.test(String(parent.dev)) && /^\d+$/u.test(String(parent.ino)), "target_shape");
|
||||
}
|
||||
fail(typeof watcher.complete === "boolean" && [watcher.targetMutationCount, watcher.workspaceMutationCount].every(n => Number.isSafeInteger(n) && (n as number) >= 0), "watcher_shape");
|
||||
const validProcess = (p: unknown) => {
|
||||
const i = record(p); return Number.isSafeInteger(i.pid) && (i.pid as number) >= 2 && Number.isSafeInteger(i.ppid) && (i.ppid as number) >= 0
|
||||
&& typeof i.startTicks === "string" && /^\d+$/u.test(i.startTicks) && typeof i.bootId === "string" && /^[a-f0-9-]{36}$/u.test(i.bootId);
|
||||
};
|
||||
fail(typeof processes.captured === "boolean" && (processes.captured ? validProcess(processes.root) : processes.root === null)
|
||||
&& Array.isArray(processes.journal) && processes.journal.length <= 512 && processes.journal.every(validProcess)
|
||||
&& Array.isArray(processes.live) && processes.live.length <= 512
|
||||
&& processes.live.every(pid => (processes.journal as unknown[]).some((p: unknown) => record(p).pid === pid)), "process_shape");
|
||||
fail(setup.path === actionFile && typeof setup.published === "boolean" && (setup.published ? sha(setup.sha256) : setup.sha256 === null), "setup_shape");
|
||||
if (row.attached !== null) {
|
||||
const a = record(row.attached), exit = record(a.commandExit);
|
||||
fail(Number.isSafeInteger(a.connections) && (a.connections as number) >= 0 && (a.failure === null || ["client_bound", "child_start", "child_failed", "marker_failed", "client_rejected"].includes(String(a.failure)))
|
||||
&& (a.commandExit === null || (Number.isSafeInteger(exit.code) && Number.isSafeInteger(exit.observedAtMs) && typeof exit.observedMonotonicNs === "string" && /^\d+$/u.test(exit.observedMonotonicNs)))
|
||||
&& [a.markerWrittenAtMs, a.clientExitedAtMs].every(n => n === null || Number.isSafeInteger(n))
|
||||
&& [a.markerWrittenMonotonicNs, a.clientExitedMonotonicNs].every(n => n === null || typeof n === "string" && /^\d+$/u.test(n)), "attached_shape");
|
||||
}
|
||||
const { files: _files, ...publicReceipt } = row;
|
||||
return { ...publicReceipt, receivedAtMs: Date.now() } as unknown as RemoteNativeSnapshot;
|
||||
}
|
||||
|
||||
/** Must be called while the initial native input-file bootstrap holds the run.
|
||||
* No effect-producing task is published until baseline proves watcher + PID
|
||||
* admission. Same-UID observer opacity is not an OS adversarial sandbox. */
|
||||
export async function bindRemoteNativeFixture(options: RemoteNativeFixtureOptions): Promise<RemoteNativeFixture> {
|
||||
const { authority, api, daytona } = options;
|
||||
fail(options.sdkVersion === REMOTE_FIXTURE_DAYTONA_SDK_VERSION, "sdk_pin");
|
||||
fail(Object.entries(authority).every(([key, value]) => key === "image" ? typeof value === "string" && /^[^\s]+@sha256:[a-f0-9]{64}$/u.test(value) : typeof value === "string" && id(value)), "authority_shape");
|
||||
fail(sha(options.nodeSha256) && sha(options.runnerdSha256), "binary_pins");
|
||||
fail(options.targets.length <= 8 && new Set(options.targets).size === options.targets.length, "target_bound");
|
||||
const targets = options.targets.map(relative), actionFile = relative(options.actionFile);
|
||||
fail(!targets.includes(actionFile), "setup_target_overlap");
|
||||
fail(!options.crossRoot || Buffer.byteLength(options.crossRoot.initialText) <= 4096, "cross_root_bound");
|
||||
const root = `/tmp/pc-native-${randomBytes(18).toString("hex")}`, nonce = randomBytes(32).toString("hex");
|
||||
let binding: RemoteNativeBinding | undefined, sentinel: { path: string; token: string } | undefined, identityKey: string | undefined;
|
||||
async function admittedSandbox() {
|
||||
const list = await api.get<unknown>(`/api/environments/${authority.environmentId}/leases`);
|
||||
const row = record(await api.get<unknown>(`/api/environment-leases/${authority.leaseId}`));
|
||||
fail(Array.isArray(list) && list.filter(item => record(item).id === authority.leaseId).length === 1, "lease_list");
|
||||
function admitted(row: Record<string, unknown>) {
|
||||
const metadata = record(row.metadata), sent = record(metadata.workspaceSentinel), cwd = metadata.remoteCwd;
|
||||
fail(row.id === authority.leaseId && row.companyId === authority.companyId && row.environmentId === authority.environmentId
|
||||
&& row.heartbeatRunId === authority.runId && row.provider === "daytona" && row.providerLeaseId === authority.sandboxId
|
||||
&& row.status === "active" && row.releasedAt == null, "lease_scope");
|
||||
fail(metadata.sandboxId === authority.sandboxId && metadata.image === authority.image && metadata.reuseLease === false, "lease_metadata");
|
||||
fail(typeof cwd === "string" && cwd.startsWith("/") && cwd.length < 512 && posix.normalize(cwd) === cwd && cwd !== "/" && !cwd.includes("\0"), "remote_cwd");
|
||||
fail(sent.path === `${cwd}/.paperclip-runtime/reusable-sandbox-lease.json` && sent.result === "written" && typeof sent.token === "string"
|
||||
&& sent.token.length >= 16 && sent.token.length <= 256 && sent.runId === authority.runId && sent.providerLeaseId === authority.sandboxId, "sentinel_binding");
|
||||
return { binding: { ...authority, remoteCwd: cwd as string }, sentinel: { path: sent.path as string, token: sent.token as string } };
|
||||
}
|
||||
const current = admitted(row), listed = admitted(record((list as unknown[]).find(item => record(item).id === authority.leaseId)));
|
||||
const key = JSON.stringify(current);
|
||||
fail(JSON.stringify(listed) === key && (!identityKey || identityKey === key), "lease_identity_drift");
|
||||
const sandbox = await daytona.get(authority.sandboxId);
|
||||
fail(sandbox.id === authority.sandboxId && sandbox.labels?.["paperclip-provider"] === "daytona"
|
||||
&& sandbox.labels?.["paperclip-company-id"] === authority.companyId && sandbox.labels?.["paperclip-environment-id"] === authority.environmentId
|
||||
&& sandbox.labels?.["paperclip-run-id"] === authority.runId && sandbox.labels?.["paperclip-reuse-lease"] === "false", "sandbox_labels");
|
||||
binding = current.binding; sentinel = current.sentinel; identityKey = key;
|
||||
return sandbox;
|
||||
}
|
||||
async function rpc(request: Record<string, unknown>, admitted?: Awaited<ReturnType<typeof admittedSandbox>>) {
|
||||
const sandbox = admitted ?? await admittedSandbox();
|
||||
const payload = Buffer.from(JSON.stringify({ ...request, root, nonce, nodeSha256: options.nodeSha256 })).toString("base64");
|
||||
const command = `/usr/bin/env -i PATH=/usr/bin:/bin ${quote(NODE)} -e ${quote(RPC)} ${quote(payload)}`;
|
||||
let deadline: ReturnType<typeof setTimeout> | undefined;
|
||||
let response: { exitCode: number; result: string };
|
||||
try {
|
||||
response = await Promise.race([
|
||||
sandbox.process.executeCommand(command, binding!.remoteCwd, {}, request.op === "wait" ? 295 : 10),
|
||||
new Promise<never>((_resolve, reject) => { deadline = setTimeout(() => reject(new Error("deadline")), request.op === "wait" ? 300_000 : 12_000); deadline.unref(); }),
|
||||
]);
|
||||
} catch { throw new Error("remote_native_fixture:remote_command_failed_or_deadline"); }
|
||||
finally { if (deadline) clearTimeout(deadline); }
|
||||
fail(response.exitCode === 0 && typeof response.result === "string" && Buffer.byteLength(response.result) <= MAX_OUTPUT, "command_failed_or_output_bound");
|
||||
let parsed: Record<string, unknown>;
|
||||
try { parsed = record(JSON.parse(response.result)); } catch { throw new Error("remote_native_fixture:invalid_observer_json"); }
|
||||
fail(parsed.ok === true, "observer_incomplete");
|
||||
return parsed.result;
|
||||
}
|
||||
const sandbox = await admittedSandbox();
|
||||
const names = [...targets, ...(options.crossRoot ? ["@cross-root"] : [])];
|
||||
const config = { root, nonce, binding, sentinel, targets, actionFile, crossRoot: options.crossRoot, runnerdSha256: options.runnerdSha256 };
|
||||
let baseline: RemoteNativeSnapshot;
|
||||
try {
|
||||
baseline = readSnapshot(await rpc({ op: "install", config, source: observerSource() }, sandbox), binding!, names, actionFile);
|
||||
fail(baseline.complete && baseline.processes.captured && baseline.processes.live.length > 0 && baseline.watcher.complete && !baseline.setup.published, "bootstrap_not_held");
|
||||
} catch (error) {
|
||||
// Only the nonce/inode-bound observer can acknowledge this cleanup. If
|
||||
// launch failed before its socket became available, TTL remains a bound,
|
||||
// not a claimed successful cleanup receipt.
|
||||
try { await rpc({ op: "close" }, sandbox); }
|
||||
catch { throw new Error("remote_native_fixture:startup_failed_cleanup_unproven", { cause: error }); }
|
||||
throw error;
|
||||
}
|
||||
// Start receiving while the lease is still authorized, before publishAction.
|
||||
// Rejection is retained (no unhandled rejection); finish reports it unchanged.
|
||||
const terminal = rpc({ op: "wait" }, sandbox).then(value => ({ value }), error => ({ error }));
|
||||
try {
|
||||
const armed = record(await rpc({ op: "arm" }, sandbox));
|
||||
fail(armed.armed === true && armed.sealed === false, "receipt_channel_not_armed");
|
||||
} catch (error) {
|
||||
try { await rpc({ op: "close" }, sandbox); }
|
||||
catch { throw new Error("remote_native_fixture:receipt_channel_failed_cleanup_unproven", { cause: error }); }
|
||||
throw error;
|
||||
}
|
||||
let closed = false, published = false, finished: RemoteNativeSnapshot | undefined;
|
||||
const retainedFiles = new Map<string, Buffer>();
|
||||
return {
|
||||
binding: binding!, remoteCwd: binding!.remoteCwd, actionFile, outsideTarget: options.crossRoot ? `${root}/cross-root-target` : null, baseline,
|
||||
async snapshot(label) {
|
||||
fail(typeof label === "string" && /^[a-zA-Z0-9_-]{1,80}$/u.test(label), "snapshot_label");
|
||||
fail(!closed && !finished, "fixture_closed");
|
||||
return readSnapshot(await rpc({ op: "snapshot" }), binding!, names, actionFile);
|
||||
},
|
||||
async readFile(path) {
|
||||
fail(!closed && names.includes(path), "unregistered_read");
|
||||
if (finished) { const bytes = retainedFiles.get(path); fail(bytes, "terminal_file_missing"); return Buffer.from(bytes!); }
|
||||
const data = record(await rpc({ op: "read", path }));
|
||||
fail(typeof data.base64 === "string" && data.base64.length <= 87384 && /^(?:[A-Za-z0-9+/]{4})*(?:[A-Za-z0-9+/]{2}==|[A-Za-z0-9+/]{3}=)?$/u.test(data.base64), "read_bound");
|
||||
const bytes = Buffer.from(data.base64 as string, "base64");
|
||||
fail(bytes.length <= 65536 && `sha256:${createHash("sha256").update(bytes).digest("hex")}` === data.sha256, "read_digest");
|
||||
return bytes;
|
||||
},
|
||||
async publishAction(path, text) {
|
||||
fail(!closed && !published && path === actionFile && typeof text === "string" && Buffer.byteLength(text) <= 16384, "publish_bound");
|
||||
// Set before RPC: uncertain delivery must never cause a second publish.
|
||||
published = true;
|
||||
const ack = record(await rpc({ op: "publish", path, text }));
|
||||
fail(ack.published === true && ack.path === path && ack.sha256 === digest(text), "publish_ack");
|
||||
},
|
||||
async setupAttachedCommand(input) {
|
||||
fail(!closed && !published && targets.includes(input.marker) && Number.isInteger(input.delayMs) && input.delayMs >= 100 && input.delayMs <= 8000
|
||||
&& typeof input.markerText === "string" && Buffer.byteLength(input.markerText) <= 4096, "attached_bound");
|
||||
const clientNonce = randomBytes(24).toString("hex");
|
||||
const result = record(await rpc({ op: "attached", ...input, clientNonce }));
|
||||
fail(result.clientScript === `${root}/client.cjs` && result.clientSocket === `${root}/attached.sock`, "attached_identity");
|
||||
const command = `${quote(NODE)} ${quote(result.clientScript as string)} ${quote(result.clientSocket as string)} ${quote(clientNonce)}`;
|
||||
return { command, commandSha256: digest(command) };
|
||||
},
|
||||
async finish() {
|
||||
if (finished) return finished;
|
||||
const receipt = await terminal;
|
||||
if ("error" in receipt) throw receipt.error;
|
||||
const result = readSnapshot(receipt.value, binding!, names, actionFile);
|
||||
fail(published && result.setup.published && result.complete && result.watcher.complete && result.processes.captured && result.processes.live.length === 0, "terminal_evidence_incomplete");
|
||||
const files = record(record(receipt.value).files);
|
||||
for (const name of names) {
|
||||
const target = result.targets[name]!;
|
||||
if (target.absent) { fail(files[name] === undefined, "terminal_file_presence"); continue; }
|
||||
fail(typeof files[name] === "string" && (files[name] as string).length <= 87384, "terminal_file_bound");
|
||||
const bytes = Buffer.from(files[name] as string, "base64");
|
||||
fail(bytes.length <= 65536 && `sha256:${createHash("sha256").update(bytes).digest("hex")}` === target.sha256, "terminal_file_digest");
|
||||
retainedFiles.set(name, bytes);
|
||||
}
|
||||
finished = result; return result;
|
||||
},
|
||||
async close() {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
// A finalized receipt survives ordinary public lease deletion. If the
|
||||
// lease still exists, clean our opaque observer; never discover by name.
|
||||
if (!finished) await rpc({ op: "close" });
|
||||
},
|
||||
};
|
||||
}
|
||||
Reference in new issue
Block a user