diff --git a/tests/runner-e2e/remote-native-fixtures.test.ts b/tests/runner-e2e/remote-native-fixtures.test.ts new file mode 100644 index 0000000000..1be821a782 --- /dev/null +++ b/tests/runner-e2e/remote-native-fixtures.test.ts @@ -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 = { 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 = { "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; 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) => unknown) | undefined; + const executeCommand = vi.fn(async (command: string, _cwd?: string, _env?: Record, 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) { 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([[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([[`${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(); 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 = {}) { + 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); + }); +}); diff --git a/tests/runner-e2e/remote-native-fixtures.ts b/tests/runner-e2e/remote-native-fixtures.ts new file mode 100644 index 0000000000..b79f59a5ea --- /dev/null +++ b/tests/runner-e2e/remote-native-fixtures.ts @@ -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 => value !== null && typeof value === "object" && !Array.isArray(value) ? value as Record : {}; +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; + targets: Record; + 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; process: { + executeCommand(command: string, cwd?: string, env?: Record, timeout?: number): Promise<{ exitCode: number; result: string }>; + } }>; +} +export interface RemoteFixtureApi { get(path: string): Promise } + +/** 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; + readFile(path: string): Promise; + publishAction(path: string, text: string): Promise; + 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; + /** Only fixture-owned observer cleanup. Never deletes a sandbox or lease. */ + close(): Promise; +} +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 { + 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(`/api/environments/${authority.environmentId}/leases`); + const row = record(await api.get(`/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) { + 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, admitted?: Awaited>) { + 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 | 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((_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; + 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(); + 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" }); + }, + }; +}