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;