mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
fix: retry sandbox ACP input delivery after gateway failures (#14485)
## 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 <noreply@paperclip.ing>
This commit is contained in:
1 parent
ea371b9684
commit
53aad90b9e
13 files changed
+851
-142
No files matched your search
@@ -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);
|
||||
|
||||
@@ -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<void> {
|
||||
async function waitFor(check: () => boolean | Promise<boolean>, timeoutMs = 4_000): Promise<void> {
|
||||
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<RunProcessResult>) | undefined;
|
||||
const local = createLocalSandboxRunner();
|
||||
const runner = {
|
||||
execute: async (input: Parameters<typeof local.execute>[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<void>((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<void>((resolve) => peer!.on("close", () => resolve()));
|
||||
await new Promise<void>((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<void>((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<void> | 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<void>((resolve) => peer!.once("close", resolve));
|
||||
await new Promise<void>((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<never>((_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 () => {
|
||||
|
||||
@@ -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<void> = 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,
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in new issue
Block a user