From 53aad90b9e83dc147707797bf224bec12600b171 Mon Sep 17 00:00:00 2001 From: Devin Foley Date: Mon, 28 Sep 2026 19:15:33 -0700 Subject: [PATCH] fix: retry sandbox ACP input delivery after gateway failures (#14485) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Thinking Path > - Paperclip coordinates agent work through execution adapters. > - Sandbox ACP sessions send ordered input through a remote file queue. > - A temporary provider 502 currently closes the session during an input upload. > - A lost response can occur after the sandbox has consumed the message, so a blind retry can duplicate input. > - This pull request retries gateway failures with the same sequence and drops consumed sequences at the receiver. > - The session can continue through a brief provider failure without repeating a tool call. ## Linked Issues or Issue Description **What happened?** A sandbox ACP run can fail with `ACP agent disconnected during request (connection_close, exit=null, signal=null)` when a provider input upload returns HTTP 502. The bridge destroys its local socket on the first failure and can discard the diagnostic before the proxy reads it. **Expected behavior** A temporary gateway failure should get a bounded retry. A lost response after successful delivery must not duplicate input or reorder later messages. Permanent failures must still close the session. **Steps to reproduce** 1. Run the real sandbox process bridge with an echo child and a local test runner. 2. Inject a provider 502 before preparation, after a chunk upload, or after final publication and consumption. 3. Send the next input message. Before this change, the connection closes instead of delivering it. Searched open and closed PRs for `ACP disconnect`, `bridge retry`, and `502 sandbox`. Related work: #13287 covers shutdown after bridge loss; #13793 covers large launch envelopes. This change covers ordered input delivery within a running legacy ACP session. ## What Changed - Retry input uploads up to three times for recognized Daytona and Cloudflare HTTP 502, 503, and 504 diagnostics, with 250 ms and 500 ms delays. - Give each upload separate temporary paths and discard already-consumed input sequences, including late publication from an earlier attempt. Clean failed attempts in the background without removing a published message or another attempt’s files. Cleanup cannot delay retries or shutdown. - Keep later input behind the retry. Stop queued input on permanent failure and flush a fixed diagnostic before closing the socket. Neither failure-diagnostic persistence nor shutdown-warning persistence can block teardown. - Add real-process regression tests for lost responses, late publication, retry exhaustion, immediate permanent failure, and diagnostic redaction. - Give accepted run-log file appends up to three seconds to drain before finalization computes the size, hash, and durable copy. Close the run handle to later appends. This waits only for file writes, independently of later DB progress or live-event persistence. If writes remain stalled, return null size/hash metadata and skip the final durable copy so the run can settle. Late writes cannot restart mirroring. - Preserve legacy comment attribution when final log size is unknown by reading existing entries within the unchanged 2 MB scan limit. Storage errors or a three-second read deadline return the evidence already read instead of failing the comment listing; pagination stops at the deadline. The deadline requests cancellation of the underlying local stream or S3 HEAD, GET, and response stream. A separate response timeout returns partial evidence even when filesystem I/O delays cancellation; late reads cannot append evidence or start another page. Each listing retains its existing batches of eight reads, without a shared admission cap that skips readable logs under contention. - Document the retry and log-finalization boundaries in the development guide. ## Verification - Final commit `347daa564b`: [Linux CI](https://github.com/paperclipai/paperclip/actions/runs/36506995168/attempts/2) passed. Greptile Apex review 13 scored this commit 5/5 with no new findings; all 12 review threads are resolved. - The final CI run initially hit a Cursor test timeout and four Discord credential-lock contention failures. All five cases passed in isolation. The two failed shards and their aggregate gate passed on retry without a code change. Those intermittent failures are not claimed fixed by this PR. - `pnpm --filter @paperclipai/adapter-utils typecheck` passed. - `pnpm exec vitest run packages/adapter-utils/src/execution-target-stdin-race.test.ts packages/adapter-utils/src/execution-target-sandbox.test.ts packages/adapter-utils/src/sandbox-callback-bridge.test.ts`: 262 tests passed on the final implementation, including 21 new regressions. The original three fault-injection cases failed before the fix. - The regressions cover failed and indefinitely stalled cleanup, Cloudflare gateway responses and retry exhaustion, permanent errors that must not retry, and teardown while failure logging remains indefinitely stalled. Seven Apex regression cases failed before the review fixes. Adapter-utils typecheck and build passed again after the final review change. - `pnpm exec vitest run server/src/services/run-log-store.test.ts server/src/services/run-log-store-cancellation.test.ts`: all 25 tests passed, including four new regressions that failed before the finalization fix. They cover delayed and failed appends, late-write admission, agreement between the local bytes/summary/durable copy, and a stalled append that exhausts the three-second budget. The timeout case verifies unknown metadata, no final upload, and no mirror restart after late completion. New cancellation tests use the real AWS SDK against a local HTTP server. They verify that stalled HEAD, GET, and response-body connections close on abort and that a subsequent read succeeds. Local range and already-aborted read cases also pass. - `pnpm exec vitest run server/src/__tests__/issues-service.test.ts -t 'readIssueCommentRunLogText|deriveIssueCommentRunLogAttribution'`: 14 targeted tests passed. The null-size reader case, both storage-error cases, the stalled-read case, and the cancellation/concurrent-listing cases failed before their fixes. The new regressions verify that timed-out reads are cancelled, subsequent listings recover, and two concurrent listings both retain their attribution markers. A read that ignores cancellation still returns partial evidence at three seconds and cannot resume pagination when it finishes; this regression failed before the response-timeout fix. - `pnpm --filter @paperclipai/server typecheck` and `pnpm --filter @paperclipai/server build` passed after the response-timeout change. - Full `pnpm -r typecheck` and `pnpm build` passed earlier in this PR; the affected packages were rechecked after review fixes. - Full local `pnpm test:run` failed in the general-server group: 511 files passed, 40 failed, and 158 were skipped. Failures include embedded PostgreSQL initialization, read-only cache directory renames, a macOS long-path fixture, and a workspace exposure assertion. The PostgreSQL, cache-permission, and long-path failures also reproduce with both changed implementation files restored to baseline commit `24c58e479a`. The exposure suite passes in isolation both on baseline and the fixed branch (28 passed, 3 skipped). CI runs the full suite on Linux. Later local test groups were not reached. - An earlier CI run hit the Telegram retry-timing failure fixed upstream in #14501. The branch includes that master fix. The selected recovery test passed against a fresh, migrated PostgreSQL 16 database. The embedded PostgreSQL runner is unavailable on this Mac; the isolated database was stopped and removed afterward. - No live agent turn was replayed. The tests use local child processes and injected provider failures. ## Risks Retries are restricted to recognized Daytona SDK and Cloudflare bridge gateway-error messages, which survive plugin RPC serialization. Other errors fail immediately. Temporary upload paths are now unique for all command-managed queue writes. Receiver sequence checks prevent duplicate input; retries do not restart an agent turn. Cleanup and failure logging are nonblocking and best effort; session teardown remains the final cleanup boundary. Log finalization now drains accepted local file writes for at most three seconds and ignores later appends on the closed run handle. A timeout leaves final size/hash unknown and skips the final durable upload; an existing partial mirror may remain available, but it is not claimed as a verified final snapshot. It does not wait for later DB progress or live-event persistence. Optional attribution keeps partial evidence when a read fails or times out. Cancellation closes S3 requests and response streams. Local filesystem I/O may finish after the caller deadline, but a late read cannot change the returned evidence or continue pagination. Later listings can retry after storage recovers. There is no schema, authentication, or permission change. Revert this commit to restore the previous behavior. ## Model Used OpenAI GPT-6 through Codex, with reasoning, repository inspection, code editing, and local test execution. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally; targeted tests pass and full-suite limitations are documented above - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip --- doc/DEVELOPING.md | 29 +++ .../src/execution-target-sandbox.test.ts | 2 +- .../src/execution-target-stdin-race.test.ts | 245 +++++++++++++++++- .../adapter-utils/src/execution-target.ts | 72 +++-- .../src/sandbox-callback-bridge.ts | 45 +++- server/src/__tests__/issues-service.test.ts | 153 ++++++++++- server/src/services/heartbeat.ts | 8 +- server/src/services/issues.ts | 121 +++++---- .../run-log-store-cancellation.test.ts | 72 +++++ server/src/services/run-log-store.test.ts | 121 +++++++-- server/src/services/run-log-store.ts | 115 +++++--- server/src/storage/s3-provider.ts | 8 +- server/src/storage/types.ts | 2 + 13 files changed, 851 insertions(+), 142 deletions(-) create mode 100644 server/src/services/run-log-store-cancellation.test.ts diff --git a/doc/DEVELOPING.md b/doc/DEVELOPING.md index 0bca5513e5..c830de4bcb 100644 --- a/doc/DEVELOPING.md +++ b/doc/DEVELOPING.md @@ -1109,6 +1109,35 @@ agent workspace. The host `HOME` itself, a directory that contains it, a filesystem root, a `CODEX_HOME` overlap, or a canonical path outside the assigned workspace is rejected before provider startup. +### Sandbox ACP input delivery + +The legacy sandbox process bridge retries recognized Daytona and Cloudflare +HTTP 502, 503, and 504 failures while writing an input message, with at most +three attempts and a short backoff. +Retries keep the message sequence and use separate temporary upload files. +The remote wrapper discards already-consumed sequences, so a lost provider +response cannot send the same input bytes twice. Messages remain ordered. +This does not restart an agent turn or replay a tool call. Authentication and +shell errors fail immediately; exhausted input delivery closes the bridge and +records a fixed diagnostic without logging the input payload. Persisting that +failure diagnostic does not block bridge teardown. +Run-log finalization closes its write handle and waits for accepted file +appends before computing the size, hash, and durable copy. Writes submitted +after finalization starts are ignored; later progress persistence is not part +of that file-write barrier. If accepted writes remain stalled after three +seconds, finalization returns unknown size/hash metadata and skips the final +durable copy so the run can reach a terminal state. A late write cannot restart +mirroring or produce a claimed verified snapshot. +Readers still attempt bounded reads when size is unknown. Legacy comment +attribution retains its existing 2 MB scan limit and allows three seconds per +log. Storage errors or timeouts preserve any evidence already read and leave +the comments available without additional derived attribution. +The read deadline requests cancellation of local file streams, S3 HEAD and GET +requests, and S3 response streams. The listing stops waiting at the deadline +even if filesystem I/O delays cancellation. Late results cannot add evidence +or start another page. Each listing retains its existing batches of eight reads; +concurrent listings do not skip healthy logs because another listing is busy. + ### Preinstalled remote runner runtime For fast sandbox startup, bake `paperclip-runnerd` and the latest stable agent diff --git a/packages/adapter-utils/src/execution-target-sandbox.test.ts b/packages/adapter-utils/src/execution-target-sandbox.test.ts index 36695e85c1..bf600fdd74 100644 --- a/packages/adapter-utils/src/execution-target-sandbox.test.ts +++ b/packages/adapter-utils/src/execution-target-sandbox.test.ts @@ -609,7 +609,7 @@ describe("sandbox adapter execution targets", () => { target: { kind: "remote", transport: "sandbox", providerKey: "local-test", remoteCwd: rootDir, runner: { execute: async (input) => { - if (input.args?.[1]?.includes("command.b64.paperclip-upload.b64") && input.args[1].includes(">>")) { + if (/command\.b64\.[^/]+\.paperclip-upload\.b64/.test(input.args?.[1] ?? "") && input.args![1].includes(">>")) { throw new Error("Upload interrupted"); } return delegate.execute(input); diff --git a/packages/adapter-utils/src/execution-target-stdin-race.test.ts b/packages/adapter-utils/src/execution-target-stdin-race.test.ts index 9341358c38..a87deb0e89 100644 --- a/packages/adapter-utils/src/execution-target-stdin-race.test.ts +++ b/packages/adapter-utils/src/execution-target-stdin-race.test.ts @@ -122,10 +122,10 @@ describe("stdin file race (parent PAP-4037)", () => { return new Promise((resolve) => setTimeout(resolve, ms)); } - async function waitFor(check: () => boolean, timeoutMs = 4_000): Promise { + async function waitFor(check: () => boolean | Promise, timeoutMs = 4_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { - if (check()) return; + if (await check()) return; await delay(20); } throw new Error("Timed out waiting for condition."); @@ -492,8 +492,247 @@ describe("stdin file race (parent PAP-4037)", () => { } }); + it.each([ + ...["prepare", "append", "finalize", "late-finalize"].map((stage) => + [stage, "Request failed with status code 502"] as const), + ...[502, 503, 504].map((status) => + ["finalize", `Cloudflare sandbox bridge request failed with HTTP ${status}.`] as const), + ])( + "recovers a transient %s failure (%s) without repeating or reordering stdin", + async (stage, failure) => { + const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-stdin-retry-")); + cleanupDirs.push(rootDir); + const childPath = path.join(rootDir, "echo-child.mjs"); + await writeFile(childPath, "process.stdin.on('data', (c) => process.stdout.write(c));\n", "utf8"); + const first = "first-" + "x".repeat(70_000); + let delivered = ""; + let injected = false; + let lateFinalize: (() => Promise) | undefined; + const local = createLocalSandboxRunner(); + const runner = { + execute: async (input: Parameters[0]) => { + const script = input.args?.[1] ?? ""; + const matches = script.includes("/stdin/000000000001.json") && ( + stage === "prepare" ? script.includes("mkdir -p") : + stage === "append" ? script.startsWith("printf") : script.startsWith("base64 -d") + ); + if (matches && !injected) { + injected = true; + if (stage === "late-finalize") lateFinalize = () => local.execute(input); + else if (stage !== "prepare") await local.execute(input); + // The provider can lose the response after the receiver consumed + // the file. A retry must not repeat those bytes on the ACP stream. + if (stage === "finalize") await waitFor(() => delivered === first, 8_000); + throw new Error(failure); + } + return local.execute(input); + }, + }; + const bridge = await startAdapterExecutionTargetProcessSessionBridge({ + runId: "run-stdin-retry", + target: { kind: "remote", transport: "sandbox", remoteCwd: rootDir, runner }, + runtimeRootDir: path.join(rootDir, "runtime"), + adapterKey: "acpx", command: process.execPath, args: [childPath], cwd: rootDir, env: {}, + }); + let peer: net.Socket | undefined; + try { + const source = await readFile(bridge!.agentCommand, "utf8"); + const port = Number(/port: (\d+)/.exec(source)![1]); + const token = JSON.parse(/const token = (".*?");/.exec(source)![1]) as string; + peer = net.createConnection({ host: "127.0.0.1", port }); + peer.on("error", () => {}); + peer.setEncoding("utf8"); + let buffer = ""; + peer.on("data", (chunk) => { + buffer += chunk; + const lines = buffer.split("\n"); + buffer = lines.pop()!; + for (const line of lines) { + const frame = JSON.parse(line) as DeliveredFrame; + delivered += collectDelivered([frame]); + } + }); + await new Promise((resolve) => peer!.once("connect", resolve)); + for (const text of [first, "-second"]) + peer.write(JSON.stringify({ token, type: "stdin", data: Buffer.from(text).toString("base64") }) + "\n"); + await waitFor(() => delivered.endsWith("-second"), 10_000); + expect(injected).toBe(true); + expect(delivered).toBe(first + "-second"); + if (lateFinalize) { + // A provider can return 502 while its original finalize still runs. + // Cleanup may invalidate its private upload, but it cannot touch + // the retry's data or repeat input after newer messages arrived. + await lateFinalize(); + peer.write(JSON.stringify({ token, type: "stdin", data: Buffer.from("-third").toString("base64") }) + "\n"); + await waitFor(() => delivered.endsWith("-third"), 8_000); + expect(delivered).toBe(first + "-second-third"); + } + await waitFor(async () => { + const files = await readdir(path.join(rootDir, "runtime", "process-sessions"), { recursive: true }); + return files.every((file) => !file.endsWith(".paperclip-upload.b64") && !file.endsWith(".paperclip-upload.decoded")); + }); + } finally { + peer?.destroy(); + await bridge?.stop(); + } + }, + 20_000, + ); + + it.each([ + ["Request failed with status code 502", 3], + ["Request failed with status code 503", 3], + ["Request failed with status code 504", 3], + ["Cloudflare sandbox bridge request failed with HTTP 502.", 3], + ["Cloudflare sandbox bridge request failed with HTTP 503.", 3], + ["Cloudflare sandbox bridge request failed with HTTP 504.", 3], + ["Request failed with status code 403", 1], + ["Cloudflare sandbox bridge request failed with HTTP 403.", 1], + ["Remote command failed: Request failed with status code 502", 1], + ["Cloudflare sandbox bridge request failed with HTTP 502. sensitive-input", 1], + ["Remote command failed: sensitive-input", 1], + ] as const)("bounds input failure %s to %i attempts and stops later writes", async (failure, expectedAttempts) => { + const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-stdin-failed-")); + cleanupDirs.push(rootDir); + let attempts = 0; + let laterWrite = false; + let stderr = ""; + const runner = createLocalSandboxRunner(async (script) => { + if (!script.startsWith("mkdir -p")) return; + if (script.includes("/stdin/000000000002.json")) laterWrite = true; + if (script.includes("/stdin/000000000001.json")) { + attempts += 1; + throw new Error(failure); + } + }); + const bridge = await startAdapterExecutionTargetProcessSessionBridge({ + runId: "run-stdin-failed", + target: { kind: "remote", transport: "sandbox", remoteCwd: rootDir, runner }, + runtimeRootDir: path.join(rootDir, "runtime"), + adapterKey: "acpx", command: "cat", args: [], cwd: rootDir, env: {}, + onLog: async (stream, chunk) => { if (stream === "stderr") stderr += chunk; }, + }); + let peer: net.Socket | undefined; + try { + const source = await readFile(bridge!.agentCommand, "utf8"); + const port = Number(/port: (\d+)/.exec(source)![1]); + const token = JSON.parse(/const token = (".*?");/.exec(source)![1]) as string; + peer = net.createConnection({ host: "127.0.0.1", port }); + peer.setEncoding("utf8"); + peer.on("error", () => {}); + let output = ""; + peer.on("data", (chunk) => { output += chunk; }); + const closed = new Promise((resolve) => peer!.on("close", () => resolve())); + await new Promise((resolve) => peer!.once("connect", resolve)); + for (const text of ["first", "second"]) + peer.write(JSON.stringify({ token, type: "stdin", data: Buffer.from(text).toString("base64") }) + "\n"); + await closed; + expect(attempts).toBe(expectedAttempts); + expect(laterWrite).toBe(false); + expect(JSON.parse(output)).toEqual({ type: "error", message: "ACP process session input delivery failed." }); + expect(stderr).toContain("ACP process session input delivery failed."); + expect(stderr).not.toContain(failure); + } finally { + peer?.destroy(); + await bridge?.stop(); + } + }, 15_000); + + it("stops after exhausted input retries even when failure logging stalls", async () => { + const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-stdin-log-stall-")); + cleanupDirs.push(rootDir); + let attempts = 0; + let loggingStarted = false; + let releaseLog!: () => void; + const stalledLog = new Promise((resolve) => { releaseLog = resolve; }); + const runner = createLocalSandboxRunner(async (script) => { + if (script.startsWith("mkdir -p") && script.includes("/stdin/000000000001.json")) { + attempts += 1; + throw new Error("Request failed with status code 502"); + } + }); + const bridge = await startAdapterExecutionTargetProcessSessionBridge({ + runId: "run-stdin-log-stall", + target: { kind: "remote", transport: "sandbox", remoteCwd: rootDir, runner }, + runtimeRootDir: path.join(rootDir, "runtime"), + adapterKey: "acpx", command: "cat", args: [], cwd: rootDir, env: {}, + onLog: async (stream) => { + if (stream === "stderr") { + loggingStarted = true; + await stalledLog; + } + }, + }); + let peer: net.Socket | undefined; + let stop: Promise | undefined; + try { + const source = await readFile(bridge!.agentCommand, "utf8"); + const port = Number(/port: (\d+)/.exec(source)![1]); + const token = JSON.parse(/const token = (".*?");/.exec(source)![1]) as string; + peer = net.createConnection({ host: "127.0.0.1", port }); + peer.setEncoding("utf8"); + peer.on("error", () => {}); + let output = ""; + peer.on("data", (chunk) => { output += chunk; }); + const closed = new Promise((resolve) => peer!.once("close", resolve)); + await new Promise((resolve) => peer!.once("connect", resolve)); + peer.write(JSON.stringify({ token, type: "stdin", data: Buffer.from("input").toString("base64") }) + "\n"); + await closed; + expect(attempts).toBe(3); + expect(loggingStarted).toBe(true); + expect(JSON.parse(output)).toEqual({ type: "error", message: "ACP process session input delivery failed." }); + let stopped = false; + stop = bridge!.stop().then(() => { stopped = true; }); + // Teardown has a three-second acknowledgement budget. It must finish + // while the run-log promise remains unresolved, including local cleanup. + await waitFor(() => stopped, 6_000); + await expect(lstat(bridge!.agentCommand)).rejects.toMatchObject({ code: "ENOENT" }); + } finally { + releaseLog(); + peer?.destroy(); + await (stop ?? bridge?.stop()); + } + }, 15_000); + // ---- Host atomic-write tests ------------------------------------------ + it.each(["fails", "stalls"])("preserves the upload failure when best-effort cleanup %s", async (cleanupMode) => { + const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-upload-cleanup-")); + cleanupDirs.push(rootDir); + const local = createLocalSandboxRunner(); + const uploadFailure = new Error("Request failed with status code 502"); + let cleanupAttempted = false; + let rejectCleanup!: (error: Error) => void; + const stalledCleanup = new Promise((_resolve, reject) => { rejectCleanup = reject; }); + // Observe the test-owned promise even in the immediate-failure case. + void stalledCleanup.catch(() => {}); + const client = createCommandManagedSandboxCallbackBridgeQueueClient({ + remoteCwd: rootDir, + runner: { + execute: async (input) => { + const script = input.args?.[1] ?? ""; + if (script.startsWith("rm -f")) { + cleanupAttempted = true; + if (cleanupMode === "stalls") return stalledCleanup; + throw new Error("Request failed with status code 403"); + } + const result = await local.execute(input); + if (script.startsWith("printf")) throw uploadFailure; + return result; + }, + }, + }); + try { + await expect(Promise.race([ + client.writeTextFile(path.join(rootDir, "message.json"), "test input"), + delay(1_000).then(() => { throw new Error("Upload waited for stalled cleanup"); }), + ])).rejects.toBe(uploadFailure); + expect(cleanupAttempted).toBe(true); + } finally { + rejectCleanup(new Error("Cleanup unavailable")); + } + }); + // A runner that executes each bridge shell script on the local filesystem, // so the test exercises the real command-managed `writeTextFile` script. function createLocalShellRunner(scripts: string[]) { @@ -572,7 +811,7 @@ describe("stdin file race (parent PAP-4037)", () => { expect(finalizeScript).toBeDefined(); expect(finalizeScript).toContain(`mv `); expect(finalizeScript).not.toContain(`> '${jsonPath}'`); - expect(finalizeScript).toContain(`> '${jsonPath}.paperclip-upload.decoded'`); + expect(finalizeScript).toMatch(/> '[^']+\.paperclip-upload\.decoded'/); }); it("never exposes a partial .json file under a concurrent reader (command-managed host write)", async () => { diff --git a/packages/adapter-utils/src/execution-target.ts b/packages/adapter-utils/src/execution-target.ts index 54b055bac6..ca816ca25c 100644 --- a/packages/adapter-utils/src/execution-target.ts +++ b/packages/adapter-utils/src/execution-target.ts @@ -1985,6 +1985,11 @@ export async function startAdapterExecutionTargetProcessSessionBridge(input: { const target = input.target; const onLog = input.onLog ?? (async () => {}); + // Failure diagnostics are best effort: stalled or failed run-log persistence + // must not prevent sending shutdown or removing the bridge's session files. + const logFailureWithoutWaiting = (message: string) => { + void Promise.resolve().then(() => onLog("stderr", message)).catch(() => undefined); + }; const runner = requireSandboxRunner(target); // Run one unit of run-time work under its named wrapper span when a span // runner is injected. Without a runner, run the work under the current run @@ -2135,6 +2140,28 @@ export async function startAdapterExecutionTargetProcessSessionBridge(input: { // a big earlier chunk, so the wrapper reads the stdin bytes out of order and // corrupts a large prompt on the stdin path. let stdinWriteChain: Promise = Promise.resolve(); + let stdinDeliveryFailed = false; + const writeStdinFile = async (filePath: string, body: string) => { + // Retry the same sequence, never the ACP request or the tool itself. Each + // upload uses private temporary paths, and the wrapper drops sequences it + // already consumed when a provider loses the final rename's response. + for (let attempt = 1; ; attempt += 1) { + try { + await client.writeTextFile(filePath, body); + return; + } catch (error) { + // Plugin RPC preserves provider messages but not HTTP error classes. + // Match the Daytona SDK and Cloudflare bridge's gateway diagnostics + // exactly; shell failures and auth errors must still fail immediately. + const gatewayFailure = error instanceof Error && ( + /^Request failed with status code (502|503|504)$/.test(error.message) || + /^Cloudflare sandbox bridge request failed with HTTP (502|503|504)\.$/.test(error.message) + ); + if (!gatewayFailure || attempt >= 3) throw error; + await new Promise((resolve) => setTimeout(resolve, attempt * 250)); + } + } + }; let pollTimer: NodeJS.Timeout | null = null; const pendingRemoteEvents: Array<{ type?: string; @@ -2266,19 +2293,24 @@ export async function startAdapterExecutionTargetProcessSessionBridge(input: { // Chain this write after the previous one, so the atomic rename for // file N finishes before the write for file N+1 starts. Keep the // per-message `sandbox.agentSession.sendInput` span inside the chain. - const write = stdinWriteChain.then(() => - runRuntimeWork(AGENT_SESSION_SEND_INPUT_SPAN, () => - client.writeTextFile(filePath, jsonLine(stdinPayload)), - ), - ); - // The next message chains after this write on success or failure, so a - // failed write never blocks the chain. This mirrors the wrapper - // `writeChain` pattern for its event files. - stdinWriteChain = write.then(() => undefined, () => undefined); - // Keep the failure behavior: send one error line, then destroy the socket. - write.catch((error) => { - nextSocket.write(jsonLine({ type: "error", message: error instanceof Error ? error.message : String(error) })); - nextSocket.destroy(); + stdinWriteChain = stdinWriteChain.then(async () => { + if (stdinDeliveryFailed) return; + try { + await runRuntimeWork(AGENT_SESSION_SEND_INPUT_SPAN, () => + writeStdinFile(filePath, jsonLine(stdinPayload)), + ); + } catch { + stdinDeliveryFailed = true; + stopping = true; + const message = "ACP process session input delivery failed."; + // Flush the diagnostic before closing; destroy() can discard it + // and leave only ACP's generic connection_close error. Do not + // expose provider error text, which may contain a command payload. + nextSocket.end(jsonLine({ type: "error", message })); + // stop() awaits this input chain before sending shutdown. Run-log + // persistence must not hold teardown open when it stalls or fails. + logFailureWithoutWaiting(`[paperclip] ${message}\n`); + } }); } } @@ -2549,10 +2581,9 @@ export async function startAdapterExecutionTargetProcessSessionBridge(input: { ]); stopReadingForShutdownAck = true; if (!acknowledgedInTime) { - await onLog( - "stderr", + logFailureWithoutWaiting( `[paperclip] ACP process session wrapper did not acknowledge shutdown within ${DEFAULT_PROCESS_SESSION_SHUTDOWN_WAIT_MS}ms; removing the session directory anyway.\n`, - ).catch(() => undefined); + ); } // Unconditional: this removal runs whether or not the wrapper // acknowledged, and whether or not any event (real or forged) arrived @@ -2983,6 +3014,14 @@ async function pollStdin() { for (const name of entries) { if (shuttingDown) break; const entrySeq = Number.parseInt(name, 10); + const file = path.posix.join(stdinDir, name); + // A successful publication can be retried after its provider response + // was lost, even after we consumed it. Never send those bytes twice or + // move the expected sequence backwards. This also handles late uploads. + if (Number.isFinite(entrySeq) && entrySeq < stdinExpectedSeq) { + await fs.rm(file, { force: true }).catch(() => undefined); + continue; + } // Hold the send order when an earlier file has not appeared. Do not consume // this later file: wait for the missing file on a later cycle, bounded by // the retry budget. After the budget, fail loud and advance past the gap, @@ -3001,7 +3040,6 @@ async function pollStdin() { stdinGapRetries = 0; stdinExpectedSeq = entrySeq; } - const file = path.posix.join(stdinDir, name); let message; try { // Hardening (I3): open with O_NOFOLLOW where the platform defines it, diff --git a/packages/adapter-utils/src/sandbox-callback-bridge.ts b/packages/adapter-utils/src/sandbox-callback-bridge.ts index d9145c9c25..40776a3da5 100644 --- a/packages/adapter-utils/src/sandbox-callback-bridge.ts +++ b/packages/adapter-utils/src/sandbox-callback-bridge.ts @@ -701,23 +701,40 @@ export function createCommandManagedSandboxCallbackBridgeQueueClient(input: { // then moves the complete decoded content onto the final `.json` path. // A direct `> remotePath` redirect truncates the final path before the // decode writes it, so a reader can see an empty or partial file. - const tempPath = `${remotePath}.paperclip-upload.b64`; - const decodedPath = `${remotePath}.paperclip-upload.decoded`; - await runChecked( - `prepare upload ${remotePath}`, - `mkdir -p ${shellQuote(remoteDir)} && rm -f ${shellQuote(tempPath)} ${shellQuote(decodedPath)} && : > ${shellQuote(tempPath)}`, - ); - const base64Body = toBuffer(Buffer.from(body, "utf8")).toString("base64"); - for (const chunk of base64Chunks(base64Body)) { + // A failed provider response does not prove the remote command stopped. + // Keep concurrent or retried uploads from truncating each other's bytes. + const uploadPath = `${remotePath}.${randomUUID()}.paperclip-upload`; + const tempPath = `${uploadPath}.b64`; + const decodedPath = `${uploadPath}.decoded`; + try { await runChecked( - `append upload chunk ${remotePath}`, - `printf '%s' ${shellQuote(chunk)} >> ${shellQuote(tempPath)}`, + `prepare upload ${remotePath}`, + `mkdir -p ${shellQuote(remoteDir)} && rm -f ${shellQuote(tempPath)} ${shellQuote(decodedPath)} && : > ${shellQuote(tempPath)}`, ); + const base64Body = toBuffer(Buffer.from(body, "utf8")).toString("base64"); + for (const chunk of base64Chunks(base64Body)) { + await runChecked( + `append upload chunk ${remotePath}`, + `printf '%s' ${shellQuote(chunk)} >> ${shellQuote(tempPath)}`, + ); + } + await runChecked( + `finalize upload ${remotePath}`, + `base64 -d < ${shellQuote(tempPath)} > ${shellQuote(decodedPath)} && mv ${shellQuote(decodedPath)} ${shellQuote(remotePath)} && rm -f ${shellQuote(tempPath)}`, + ); + } catch (error) { + // Abandon only this attempt's intermediates, never the published file + // or another attempt. A late finalize may fail or finish publishing; + // either is safe for a sequence-aware caller. Preserve the original + // failure even when the provider is still unavailable for cleanup. + // Cleanup must not put another provider timeout on the retry/shutdown + // path. Its unique paths stay safe to remove after this call returns. + void runChecked( + `clean failed upload ${remotePath}`, + `rm -f ${shellQuote(tempPath)} ${shellQuote(decodedPath)}`, + ).catch(() => undefined); + throw error; } - await runChecked( - `finalize upload ${remotePath}`, - `base64 -d < ${shellQuote(tempPath)} > ${shellQuote(decodedPath)} && mv ${shellQuote(decodedPath)} ${shellQuote(remotePath)} && rm -f ${shellQuote(tempPath)}`, - ); }, writeResponseFile: async (responsePath, body, options = {}) => { const responseDir = path.posix.dirname(responsePath); diff --git a/server/src/__tests__/issues-service.test.ts b/server/src/__tests__/issues-service.test.ts index 9309731a6d..a3fa946571 100644 --- a/server/src/__tests__/issues-service.test.ts +++ b/server/src/__tests__/issues-service.test.ts @@ -1,6 +1,6 @@ import { randomUUID } from "node:crypto"; import { asc, eq } from "drizzle-orm"; -import { afterAll, afterEach, beforeAll, describe, expect, it } from "vitest"; +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; import { sql } from "drizzle-orm"; import { activityLog, @@ -41,7 +41,9 @@ import { deriveIssueCommentRunLogAttribution, ISSUE_LIST_MAX_LIMIT, issueService, + readIssueCommentRunLogText, } from "../services/issues.ts"; +import { getRunLogStore } from "../services/run-log-store.js"; import { WORKSPACE_WORKTREE_REQUIRES_PROJECT_CODE, WORKSPACE_WORKTREE_REQUIRES_PROJECT_MESSAGE, @@ -149,6 +151,155 @@ describeEmbeddedPostgres("issueService run attachment artifacts", () => { }, 20_000); }); +describe("readIssueCommentRunLogText", () => { + it("cancels timed-out storage reads so later listings can recover", async () => { + let active = 0; + const cleanups: Array<() => void> = []; + const read = vi.spyOn(getRunLogStore(), "read").mockImplementation((_handle, options) => + new Promise((_resolve, reject) => { + active += 1; + let settled = false; + const abort = () => { + if (settled) return; + settled = true; + active -= 1; + reject(new DOMException("Read aborted", "AbortError")); + }; + cleanups.push(abort); + options?.signal?.addEventListener("abort", abort, { once: true }); + }), + ); + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + const run = { runId: "run", logStore: "local_file", logRef: "test/run.ndjson", logBytes: null }; + const firstBatch = Array.from({ length: 8 }, () => readIssueCommentRunLogText(run)); + try { + await vi.advanceTimersByTimeAsync(3_000); + expect(await Promise.all(firstBatch)).toEqual(Array(8).fill("")); + expect(active).toBe(0); + read.mockResolvedValueOnce({ content: "storage recovered" }); + await expect(readIssueCommentRunLogText(run)).resolves.toBe("storage recovered"); + expect(read).toHaveBeenCalledTimes(9); + } finally { + for (const cleanup of cleanups) cleanup(); + await Promise.allSettled(firstBatch); + await vi.advanceTimersByTimeAsync(0); + vi.useRealTimers(); + read.mockRestore(); + } + }); + + it("keeps readable attribution evidence for concurrent listings", async () => { + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const read = vi.spyOn(getRunLogStore(), "read").mockImplementation(async () => { + await gate; + return { content: "comment id: legacy-comment" }; + }); + const run = { runId: "run", logStore: "local_file", logRef: "test/run.ndjson", logBytes: null }; + const listings = Array.from({ length: 2 }, () => + Promise.all(Array.from({ length: 8 }, () => readIssueCommentRunLogText(run))), + ); + try { + release(); + for (const listing of listings) { + expect(await listing).toEqual(Array(8).fill("comment id: legacy-comment")); + } + expect(read).toHaveBeenCalledTimes(16); + } finally { + release(); + await Promise.allSettled(listings); + read.mockRestore(); + } + }); + + it.each([null, 128])("keeps partial attribution evidence when storage fails with logBytes=%s", async (logBytes) => { + const read = vi.spyOn(getRunLogStore(), "read").mockRejectedValue(new Error("Storage gateway unavailable")); + const run = { runId: "run", logStore: "local_file", logRef: "test/run.ndjson", logBytes }; + try { + await expect(readIssueCommentRunLogText(run)).resolves.toBe(""); + read.mockResolvedValueOnce({ content: "earlier evidence", nextOffset: 16 }); + await expect(readIssueCommentRunLogText(run)).resolves.toBe("earlier evidence"); + } finally { + read.mockRestore(); + } + }); + + it("bounds reads that ignore cancellation and stops late pagination after the deadline", async () => { + let release!: (value: { content: string; nextOffset: number }) => void; + const stalled = new Promise<{ content: string; nextOffset: number }>((resolve) => { + release = resolve; + }); + const read = vi.spyOn(getRunLogStore(), "read") + .mockResolvedValueOnce({ content: "earlier evidence", nextOffset: 16 }) + .mockReturnValueOnce(stalled); + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + let result: string | undefined; + const pending = readIssueCommentRunLogText({ + runId: "run", logStore: "local_file", logRef: "test/run.ndjson", logBytes: null, + }).then((value) => { result = value; }); + try { + await vi.advanceTimersByTimeAsync(3_000); + expect(result).toBe("earlier evidence"); + expect(read.mock.calls[1]?.[1]?.signal?.aborted).toBe(true); + release({ content: "too late", nextOffset: 24 }); + await pending; + await vi.advanceTimersByTimeAsync(0); + expect(read).toHaveBeenCalledTimes(2); + expect(result).toBe("earlier evidence"); + } finally { + release({ content: "", nextOffset: 0 }); + await pending.catch(() => {}); + vi.useRealTimers(); + read.mockRestore(); + } + }); + + it.each([null, 128, 0])("reads existing attribution markers with logBytes=%s", async (logBytes) => { + const commentId = randomUUID(); + const runId = randomUUID(); + const agentId = randomUUID(); + const read = vi.spyOn(getRunLogStore(), "read") + .mockResolvedValueOnce({ content: "comment id: ", nextOffset: 12 }) + .mockResolvedValueOnce({ content: commentId + "\n" }); + try { + const logContent = await readIssueCommentRunLogText({ + runId, logStore: "local_file", logRef: "test/run.ndjson", logBytes, + }); + const derived = deriveIssueCommentRunLogAttribution( + [{ + id: commentId, + authorAgentId: null, + authorUserId: "local-board", + createdByRunId: null, + createdAt: new Date("2020-01-01T00:00:01Z"), + }], + [{ + runId, + agentId, + createdAt: new Date("2020-01-01T00:00:00Z"), + startedAt: new Date("2020-01-01T00:00:00Z"), + finishedAt: new Date("2020-01-01T00:00:02Z"), + logContent, + }], + ); + if (logBytes === 0) { + expect(read).not.toHaveBeenCalled(); + expect(derived.size).toBe(0); + } else { + expect(read).toHaveBeenCalledTimes(2); + expect(read.mock.calls[1]?.[1]?.offset).toBe(12); + expect(derived.get(commentId)).toEqual({ + derivedAuthorAgentId: agentId, + derivedCreatedByRunId: runId, + derivedAuthorSource: "run_log_comment_post", + }); + } + } finally { + read.mockRestore(); + } + }); +}); + describe("deriveIssueCommentRunLogAttribution", () => { it("recovers agent attribution from run logs that printed the posted comment id", () => { const commentId = randomUUID(); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index d722d3ea45..e923d178d0 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -24931,8 +24931,8 @@ export function heartbeatService( : null; let logSummary: { - bytes: number; - sha256?: string; + bytes: number | null; + sha256?: string | null; compressed: boolean; } | null = null; if (handle) { @@ -25650,8 +25650,8 @@ export function heartbeatService( logger.error({ err, runId }, "heartbeat execution failed"); let logSummary: { - bytes: number; - sha256?: string; + bytes: number | null; + sha256?: string | null; compressed: boolean; } | null = null; if (handle) { diff --git a/server/src/services/issues.ts b/server/src/services/issues.ts index 9feb50394a..ceb13379f7 100644 --- a/server/src/services/issues.ts +++ b/server/src/services/issues.ts @@ -231,6 +231,7 @@ const ISSUE_COMMENT_RUN_LOG_DERIVATION_MAX_LOG_BYTES = 2_000_000; const ISSUE_COMMENT_RUN_LOG_DERIVATION_CHUNK_BYTES = 256_000; const ISSUE_COMMENT_RUN_LOG_DERIVATION_END_SLACK_MS = 60_000; const ISSUE_COMMENT_RUN_LOG_DERIVATION_MAX_PARALLEL_READS = 8; +const ISSUE_COMMENT_RUN_LOG_DERIVATION_TIMEOUT_MS = 3_000; export const ISSUE_CREATE_IDEMPOTENCY_KEY_RETENTION_DAYS = 7; const ISSUE_CREATE_IDEMPOTENCY_KEY_RETENTION_MS = ISSUE_CREATE_IDEMPOTENCY_KEY_RETENTION_DAYS * 24 * 60 * 60 * 1000; @@ -6542,6 +6543,76 @@ async function countBlockedInboxIssues( }, 0); } +export async function readIssueCommentRunLogText(run: { + runId?: string | null; + logStore: string | null; + logRef: string | null; + logBytes: number | null; +}) { + if (run.logStore !== "local_file" || !run.logRef) return ""; + // A timed-out finalization leaves size unknown even when earlier entries + // exist. Read those logs within the same byte budget as a known-size log. + if (run.logBytes !== null && (!Number.isFinite(run.logBytes) || run.logBytes <= 0)) return ""; + + const logRef = run.logRef; + const store = getRunLogStore(); + let offset = 0; + let content = ""; + let nextOffset: number | undefined = 0; + const controller = new AbortController(); + let readTimer: NodeJS.Timeout | undefined; + + const readChunks = async () => { + while (nextOffset !== undefined) { + controller.signal.throwIfAborted(); + const remainingBytes = + ISSUE_COMMENT_RUN_LOG_DERIVATION_MAX_LOG_BYTES - + Buffer.byteLength(content, "utf8"); + if (remainingBytes <= 0) break; + const chunk = await store.read( + { store: "local_file", logRef }, + { + offset, + limitBytes: Math.min(ISSUE_COMMENT_RUN_LOG_DERIVATION_CHUNK_BYTES, remainingBytes), + signal: controller.signal, + }, + ); + controller.signal.throwIfAborted(); + content += chunk.content; + nextOffset = chunk.nextOffset; + offset = chunk.nextOffset ?? 0; + } + }; + + try { + await Promise.race([ + readChunks(), + new Promise((_resolve, reject) => { + readTimer = setTimeout(() => { + const reason = new DOMException("Attribution log read timed out", "TimeoutError"); + // Cancellation closes storage work where supported, but filesystem + // I/O can delay stream destruction. Keep the response deadline too. + reject(reason); + controller.abort(reason); + }, ISSUE_COMMENT_RUN_LOG_DERIVATION_TIMEOUT_MS); + readTimer.unref?.(); + }), + ]); + } catch (err) { + // Attribution enriches already-authorized comments. Missing, failed, or + // stalled storage must not prevent listing them; keep any evidence read. + // Do not log raw provider errors, which can contain credentialed URLs. + logger.warn( + { runId: run.runId ?? undefined, logRef, status: err instanceof HttpError ? err.status : undefined }, + "could not read heartbeat run log while deriving optional issue comment metadata", + ); + } finally { + clearTimeout(readTimer); + } + + return content; +} + export function issueService(db: Db) { const instanceSettings = instanceSettingsService(db); const treeControlSvc = issueTreeControlService(db); @@ -6760,54 +6831,6 @@ export function issueService(db: Db) { }; } - async function readRunLogText(run: { - runId?: string | null; - logStore: string | null; - logRef: string | null; - logBytes: number | null; - }) { - if (run.logStore !== "local_file" || !run.logRef) return ""; - const logBytes = Number(run.logBytes ?? 0); - if (!Number.isFinite(logBytes) || logBytes <= 0) return ""; - - const store = getRunLogStore(); - let offset = 0; - let content = ""; - let nextOffset: number | undefined = 0; - - try { - while (nextOffset !== undefined) { - const remainingBytes = - ISSUE_COMMENT_RUN_LOG_DERIVATION_MAX_LOG_BYTES - - Buffer.byteLength(content, "utf8"); - if (remainingBytes <= 0) break; - const chunk = await store.read( - { store: "local_file", logRef: run.logRef }, - { - offset, - limitBytes: Math.min( - ISSUE_COMMENT_RUN_LOG_DERIVATION_CHUNK_BYTES, - remainingBytes, - ), - }, - ); - content += chunk.content; - nextOffset = chunk.nextOffset; - offset = chunk.nextOffset ?? 0; - } - } catch (err) { - if (err instanceof HttpError && err.status === 404) { - logger.warn( - { err, runId: run.runId ?? undefined, logRef: run.logRef }, - "missing heartbeat run log while deriving issue comment metadata", - ); - return content; - } - throw err; - } - - return content; - } // Persist a resolved attribution so subsequent reads stop re-scanning run // logs (and old "Board" threads stay fixed durably). Best-effort: a write @@ -7027,7 +7050,7 @@ export function issueService(db: Db) { ); await Promise.all( batch.map(async (run) => { - logByRunId.set(run.runId, await readRunLogText(run)); + logByRunId.set(run.runId, await readIssueCommentRunLogText(run)); }), ); } diff --git a/server/src/services/run-log-store-cancellation.test.ts b/server/src/services/run-log-store-cancellation.test.ts new file mode 100644 index 0000000000..36d405d132 --- /dev/null +++ b/server/src/services/run-log-store-cancellation.test.ts @@ -0,0 +1,72 @@ +import { createServer } from "node:http"; +import { promises as fs } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { createDurableRunLogStore } from "./run-log-store.js"; +import { createS3StorageProvider } from "../storage/s3-provider.js"; + +afterEach(() => vi.unstubAllEnvs()); + +describe("run-log read cancellation", () => { + it.each(["head", "get", "body"])("closes a stalled S3 %s connection and permits a subsequent read", async (stage) => { + // Exercise the real SDK against an on-host server. No provider credentials + // or external network are used by this cancellation regression. + vi.stubEnv("AWS_ACCESS_KEY_ID", "test-access-key"); + vi.stubEnv("AWS_SECRET_ACCESS_KEY", "test-secret-key"); + vi.stubEnv("AWS_SESSION_TOKEN", ""); + const basePath = await fs.mkdtemp(path.join(os.tmpdir(), "run-log-abort-")); + let recover = false; + let closed = false; + let started!: () => void; + const stalled = new Promise((resolve) => { started = resolve; }); + const server = createServer((request, response) => { + const shouldStall = !recover && (stage === "head" ? request.method === "HEAD" : request.method === "GET"); + if (shouldStall) { + response.on("close", () => { closed = true; }); + if (stage === "body") { + response.writeHead(200, { "Content-Length": "4" }); + response.write("d"); + } + started(); + return; + } + response.writeHead(200, { "Content-Length": "4" }); + response.end(request.method === "HEAD" ? undefined : "data"); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address(); + if (!address || typeof address === "string") throw new Error("Test server did not bind"); + const provider = createS3StorageProvider({ + bucket: "test-bucket", region: "us-east-1", forcePathStyle: true, + endpoint: `http://127.0.0.1:${address.port}`, + }); + const store = createDurableRunLogStore({ basePath, s3: { provider } }); + const handle = { store: "local_file" as const, logRef: "missing.ndjson" }; + const controller = new AbortController(); + const read = store.read(handle, { signal: controller.signal }); + void read.catch(() => {}); + try { + await stalled; + controller.abort(); + await expect(read).rejects.toMatchObject({ name: "AbortError" }); + await vi.waitFor(() => expect(closed).toBe(true)); + recover = true; + expect(await store.read(handle)).toEqual({ content: "data", nextOffset: undefined }); + } finally { + controller.abort(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + await fs.rm(basePath, { recursive: true, force: true }); + } + }, 10_000); + + it("rejects an already-cancelled local read before opening a file", async () => { + const store = createDurableRunLogStore({ basePath: os.tmpdir() }); + const controller = new AbortController(); + controller.abort(); + await expect(store.read({ store: "local_file", logRef: "unused.ndjson" }, { + signal: controller.signal, + })).rejects.toMatchObject({ name: "AbortError" }); + }); +}); diff --git a/server/src/services/run-log-store.test.ts b/server/src/services/run-log-store.test.ts index ab12ae17b3..ccadb78c30 100644 --- a/server/src/services/run-log-store.test.ts +++ b/server/src/services/run-log-store.test.ts @@ -2,6 +2,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { promises as fs } from "node:fs"; import path from "node:path"; import os from "node:os"; +import { createHash } from "node:crypto"; import { Readable } from "node:stream"; import { createDurableRunLogStore } from "./run-log-store.js"; import type { StorageProvider } from "../storage/types.js"; @@ -100,6 +101,102 @@ describe("createDurableRunLogStore", () => { expect(objects.get(key)!.toString("utf8")).toContain("line-B"); }); + it.each(["finishes", "fails"])("waits for an accepted file append that %s before finalizing", async (outcome) => { + const { provider, objects } = createMemoryProvider(); + const store = createDurableRunLogStore({ basePath: baseDir, s3: { provider } }); + const handle = await store.begin(begin); + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const appendFile = fs.appendFile.bind(fs); + const spy = vi.spyOn(fs, "appendFile").mockImplementationOnce(async (...args) => { + await gate; + if (outcome === "fails") throw new Error("disk unavailable"); + await appendFile(...args); + }); + const append = store.append(handle, { stream: "stderr", chunk: "late diagnostic", ts: "t1" }); + // Observe the deliberate rejection independently of finalization. + void append.catch(() => {}); + let finalized = false; + const finalize = store.finalize(handle).then((summary) => { + finalized = true; + return summary; + }); + try { + await new Promise((resolve) => setTimeout(resolve, 100)); + expect(finalized).toBe(false); + release(); + if (outcome === "fails") await expect(append).rejects.toThrow("disk unavailable"); + else await append; + const summary = await finalize; + const local = await fs.readFile(path.join(baseDir, handle.logRef)); + expect(summary.bytes).toBe(local.length); + expect(summary.sha256).toBe(createHash("sha256").update(local).digest("hex")); + expect(objects.get(handle.logRef)).toEqual(local); + expect(local.toString()).toBe(outcome === "fails" ? "" : JSON.stringify({ + ts: "t1", stream: "stderr", chunk: "late diagnostic", + }) + "\n"); + } finally { + release(); + await append.catch(() => {}); + await finalize; + spy.mockRestore(); + } + }); + + it("ignores appends once finalization starts so the durable snapshot stays immutable", async () => { + const { provider, objects } = createMemoryProvider(); + const store = createDurableRunLogStore({ basePath: baseDir, s3: { provider } }); + const handle = await store.begin(begin); + await store.append(handle, { stream: "stdout", chunk: "accepted", ts: "t1" }); + const finalize = store.finalize(handle); + expect(await store.append(handle, { stream: "stderr", chunk: "too late", ts: "t2" })).toBe(0); + const summary = await finalize; + expect(await store.append(handle, { stream: "stderr", chunk: "also too late", ts: "t3" })).toBe(0); + const local = await fs.readFile(path.join(baseDir, handle.logRef)); + expect(local.toString()).not.toContain("too late"); + expect(summary.bytes).toBe(local.length); + expect(summary.sha256).toBe(createHash("sha256").update(local).digest("hex")); + expect(objects.get(handle.logRef)).toEqual(local); + }); + + it("leaves final metadata unknown when an append stalls and never mirrors its late completion", async () => { + const { provider, calls } = createMemoryProvider(); + const store = createDurableRunLogStore({ basePath: baseDir, s3: { provider, inflightMirrorMs: 10_000 } }); + const handle = await store.begin(begin); + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const appendFile = fs.appendFile.bind(fs); + const spy = vi.spyOn(fs, "appendFile").mockImplementationOnce(async (...args) => { + await gate; + await appendFile(...args); + }); + const warn = vi.spyOn(console, "warn").mockImplementation(() => {}); + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + const append = store.append(handle, { stream: "stderr", chunk: "stalled diagnostic", ts: "t1" }); + let summary: Awaited> | undefined; + const finalize = store.finalize(handle).then((result) => { summary = result; }); + try { + await vi.advanceTimersByTimeAsync(3_000); + expect(summary).toEqual({ bytes: null, sha256: null, compressed: false }); + expect(warn).toHaveBeenCalled(); + expect(calls.put).toBe(0); + release(); + await append; + await vi.advanceTimersByTimeAsync(20_000); + await store.flushInflightMirrors!(); + expect(calls.put).toBe(0); + expect(await store.finalize(handle)).toEqual(summary); + expect(await store.append(handle, { stream: "stderr", chunk: "too late", ts: "t2" })).toBe(0); + } finally { + release(); + await append; + await finalize; + vi.useRealTimers(); + spy.mockRestore(); + warn.mockRestore(); + } + }); + it("falls back to S3 when the local file is gone (the pod-roll case that caused 'Run log not found')", async () => { const { provider } = createMemoryProvider(); const store = createDurableRunLogStore({ basePath: baseDir, s3: { provider, keyPrefix: "run-logs" } }); @@ -163,25 +260,15 @@ describe("createDurableRunLogStore", () => { expect(caughtUp.nextOffset).toBeUndefined(); }); - it("falls back to S3 when the local file vanishes between stat() and open (TOCTOU race)", async () => { - const { provider } = createMemoryProvider(); - const store = createDurableRunLogStore({ basePath: baseDir, s3: { provider, keyPrefix: "run-logs" } }); + it("reads local pages without waiting for a separate metadata request", async () => { + const store = createDurableRunLogStore({ basePath: baseDir }); const handle = await store.begin(begin); - await store.append(handle, { stream: "stdout", chunk: "raced-line", ts: "t1" }); - await store.finalize(handle); - // Delete the local file DURING stat(), i.e. after it reports the file - // present but before createReadStream opens it -> the open hits ENOENT. - const realStat = fs.stat.bind(fs); - const statSpy = vi.spyOn(fs, "stat").mockImplementation(async (target, ...rest) => { - const result = await realStat(target as Parameters[0], ...(rest as [])); - if (String(target).endsWith(".ndjson")) { - await fs.rm(target as string, { force: true }); - } - return result; - }); + await fs.writeFile(path.join(baseDir, handle.logRef), "0123456789"); + const statSpy = vi.spyOn(fs, "stat").mockRejectedValue(new Error("Metadata unavailable")); try { - const res = await store.read(handle); - expect(res.content).toContain("raced-line"); + expect(await store.read(handle, { offset: 2, limitBytes: 4 })).toEqual({ content: "2345", nextOffset: 6 }); + expect(await store.read(handle, { offset: 6, limitBytes: 4 })).toEqual({ content: "6789", nextOffset: undefined }); + expect(await store.read(handle, { offset: 20, limitBytes: 4 })).toEqual({ content: "", nextOffset: undefined }); } finally { statSpy.mockRestore(); } diff --git a/server/src/services/run-log-store.ts b/server/src/services/run-log-store.ts index f07ede9abd..7e7b4597ea 100644 --- a/server/src/services/run-log-store.ts +++ b/server/src/services/run-log-store.ts @@ -1,6 +1,7 @@ import { createReadStream, promises as fs } from "node:fs"; import path from "node:path"; import { createHash } from "node:crypto"; +import { addAbortSignal } from "node:stream"; import { notFound } from "../errors.js"; import { resolvePaperclipInstanceRoot } from "../home-paths.js"; import { createS3StorageProvider } from "../storage/s3-provider.js"; @@ -16,6 +17,7 @@ export interface RunLogHandle { export interface RunLogReadOptions { offset?: number; limitBytes?: number; + signal?: AbortSignal; } export interface RunLogReadResult { @@ -24,8 +26,10 @@ export interface RunLogReadResult { } export interface RunLogFinalizeSummary { - bytes: number; - sha256?: string; + // Null means a stalled write prevented a verified final snapshot. Callers + // can still settle the run without recording a false byte count or hash. + bytes: number | null; + sha256?: string | null; compressed: boolean; } @@ -100,6 +104,13 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru const s3 = options.s3; const s3Prefix = normalizeKeyPrefix(s3?.keyPrefix); const inflightMirrorMs = s3?.inflightMirrorMs && s3.inflightMirrorMs > 0 ? s3.inflightMirrorMs : 0; + // A run owns the write handle returned by begin(). Diagnostics can be + // dispatched without awaiting the rest of onLog (DB progress/live events), + // but finalize must include every file append already accepted on that handle. + // Weak collections let completed handles disappear with their owning runs. + const pendingAppends = new WeakMap>>(); + const closingHandles = new WeakSet(); + const abandonedHandles = new WeakSet(); function s3Key(logRef: string): string { return s3Prefix ? `${s3Prefix}/${logRef}` : logRef; @@ -159,6 +170,7 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru } function scheduleInflightMirror(logRef: string, entry: InflightMirrorEntry): void { + if (inflightMirrors.get(logRef) !== entry) return; if (entry.timer || entry.upload) return; const delay = Math.max(0, inflightMirrorMs - (Date.now() - entry.lastMirrorAt)); entry.timer = setTimeout(() => { @@ -205,33 +217,28 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru filePath: string, offset: number, limitBytes: number, + signal?: AbortSignal, ): Promise { - const stat = await fs.stat(filePath).catch(() => null); - if (!stat) return null; - const start = Math.max(0, Math.min(offset, stat.size)); - // No lower clamp to `start`: when the reader is fully caught up - // (offset === size) that clamp made end === start and produced a - // 1-byte-past-EOF range instead of an empty read. - const end = Math.min(start + limitBytes - 1, stat.size - 1); - if (start > end) return { content: "", nextOffset: start < stat.size ? start : undefined }; - + signal?.throwIfAborted(); + const start = Math.max(0, offset); + // Read one extra byte to discover whether another page exists. A single + // abortable stream avoids an uncancellable stat before opening the file. const chunks: Buffer[] = []; try { - await new Promise((resolve, reject) => { - const stream = createReadStream(filePath, { start, end }); - stream.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk))); - stream.on("error", reject); - stream.on("end", () => resolve()); - }); + const stream = createReadStream(filePath, { start, end: start + limitBytes, signal }); + for await (const chunk of stream) { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + } } catch (err) { - // File deleted between stat() and open (pod-roll cleanup racing a read): - // treat as missing so the caller falls through to the S3 mirror instead - // of surfacing the very "Run log not found" this store exists to prevent. + signal?.throwIfAborted(); + // A missing file, including deletion before open, falls back to S3. if ((err as NodeJS.ErrnoException | null)?.code === "ENOENT") return null; throw err; } - const content = Buffer.concat(chunks).toString("utf8"); - const nextOffset = end + 1 < stat.size ? end + 1 : undefined; + signal?.throwIfAborted(); + const bytes = Buffer.concat(chunks); + const content = bytes.subarray(0, limitBytes).toString("utf8"); + const nextOffset = bytes.length > limitBytes ? start + limitBytes : undefined; return { content, nextOffset }; } @@ -239,10 +246,13 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru logRef: string, offset: number, limitBytes: number, + signal?: AbortSignal, ): Promise { + signal?.throwIfAborted(); if (!s3) throw notFound("Run log not found"); const key = s3Key(logRef); - const head = await s3.provider.headObject({ objectKey: key }); + const head = await s3.provider.headObject({ objectKey: key, signal }); + signal?.throwIfAborted(); if (!head.exists) throw notFound("Run log not found"); const total = head.contentLength ?? 0; const start = Math.max(0, Math.min(offset, total)); @@ -253,13 +263,15 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru const end = Math.min(start + limitBytes - 1, total - 1); if (total === 0 || start > end) return { content: "", nextOffset: start < total ? start : undefined }; - const result = await s3.provider.getObject({ objectKey: key, range: { start, end } }); + const result = await s3.provider.getObject({ objectKey: key, range: { start, end }, signal }); + // Destroy a body that stalls after headers arrive, including a body + // returned just after the caller cancelled the request. + if (signal) addAbortSignal(signal, result.stream); const chunks: Buffer[] = []; - await new Promise((resolve, reject) => { - result.stream.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk))); - result.stream.on("error", reject); - result.stream.on("end", () => resolve()); - }); + for await (const chunk of result.stream) { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + } + signal?.throwIfAborted(); const content = Buffer.concat(chunks).toString("utf8"); const nextOffset = end + 1 < total ? end + 1 : undefined; return { content, nextOffset }; @@ -289,7 +301,7 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru }, async append(handle, event) { - if (handle.store !== "local_file") return 0; + if (handle.store !== "local_file" || closingHandles.has(handle)) return 0; const absPath = resolveWithin(basePath, handle.logRef); const line = JSON.stringify({ ts: event.ts, @@ -301,13 +313,47 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru ...(typeof event.seq === "number" && Number.isFinite(event.seq) ? { seq: event.seq } : {}), }); const persisted = `${line}\n`; - await fs.appendFile(absPath, persisted, "utf8"); - noteInflightAppend(handle.logRef); + let pending = pendingAppends.get(handle); + if (!pending) { + pending = new Set(); + pendingAppends.set(handle, pending); + } + const write = fs.appendFile(absPath, persisted, "utf8").then(() => { + if (!closingHandles.has(handle)) noteInflightAppend(handle.logRef); + }); + pending.add(write); + try { + await write; + } finally { + pending.delete(write); + } return Buffer.byteLength(persisted, "utf8"); }, async finalize(handle) { if (handle.store !== "local_file") return { bytes: 0, compressed: false }; + if (abandonedHandles.has(handle)) return { bytes: null, sha256: null, compressed: false }; + // Close admission before the first await. Drain file writes only, not + // heartbeat's later DB/live-event persistence, then freeze one consistent + // byte count, hash, and durable copy. Late diagnostics cannot mutate it. + closingHandles.add(handle); + let drainTimer: NodeJS.Timeout | undefined; + const drained = await Promise.race([ + Promise.allSettled(pendingAppends.get(handle) ?? []).then(() => true), + new Promise((resolve) => { + drainTimer = setTimeout(() => resolve(false), 3_000); + drainTimer.unref?.(); + }), + ]).finally(() => clearTimeout(drainTimer)); + if (!drained) { + // An in-flight fs append cannot be cancelled safely. Do not hash or + // mirror a file it may still change, and never re-arm mirroring when + // that write eventually finishes. Terminal run status can still settle. + abandonedHandles.add(handle); + void retireInflightMirror(handle.logRef).catch(() => undefined); + console.warn("[run-log-store] Pending log writes did not settle within 3000ms; final log size and hash are unknown."); + return { bytes: null, sha256: null, compressed: false }; + } await retireInflightMirror(handle.logRef); const absPath = resolveWithin(basePath, handle.logRef); const stat = await fs.stat(absPath).catch(() => null); @@ -344,14 +390,15 @@ export function createDurableRunLogStore(options: DurableRunLogStoreOptions): Ru }, async read(handle, opts) { + opts?.signal?.throwIfAborted(); if (handle.store !== "local_file") throw notFound("Run log not found"); const absPath = resolveWithin(basePath, handle.logRef); const offset = opts?.offset ?? 0; const limitBytes = opts?.limitBytes ?? 256_000; - const local = await readLocalRange(absPath, offset, limitBytes); + const local = await readLocalRange(absPath, offset, limitBytes, opts?.signal); if (local) return local; // Local file gone (pod rolled) -> serve from the S3 mirror if configured. - return readS3Range(handle.logRef, offset, limitBytes); + return readS3Range(handle.logRef, offset, limitBytes, opts?.signal); }, async flushInflightMirrors() { diff --git a/server/src/storage/s3-provider.ts b/server/src/storage/s3-provider.ts index 517ccf23d9..4939bb9b25 100644 --- a/server/src/storage/s3-provider.ts +++ b/server/src/storage/s3-provider.ts @@ -6,7 +6,7 @@ import { PutObjectCommand, } from "@aws-sdk/client-s3"; import { putS3Multipart } from "./s3-multipart.js"; -import { Readable } from "node:stream"; +import { addAbortSignal, Readable } from "node:stream"; import type { StorageProvider, GetObjectResult, HeadObjectResult } from "./types.js"; import { notFound, unprocessable } from "../errors.js"; @@ -106,10 +106,13 @@ export function createS3StorageProvider(config: S3ProviderConfig): StorageProvid Key: key, Range: input.range ? `bytes=${input.range.start}-${input.range.end}` : undefined, }), + { abortSignal: input.signal }, ); + const stream = await toReadableStream(output.Body); + if (input.signal) addAbortSignal(input.signal, stream); return { - stream: await toReadableStream(output.Body), + stream, contentType: output.ContentType, contentLength: output.ContentLength, etag: output.ETag, @@ -130,6 +133,7 @@ export function createS3StorageProvider(config: S3ProviderConfig): StorageProvid Bucket: bucket, Key: key, }), + { abortSignal: input.signal }, ); return { diff --git a/server/src/storage/types.ts b/server/src/storage/types.ts index 15289f77a5..30f0845b5b 100644 --- a/server/src/storage/types.ts +++ b/server/src/storage/types.ts @@ -12,6 +12,8 @@ export interface PutObjectInput { export interface GetObjectInput { objectKey: string; + // S3 reads cancel pending requests and their response streams. + signal?: AbortSignal; range?: { start: number; end: number;