mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-11 14:10:50 +02:00
fix(daytona): recover output from stalled log streams (#14889)
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - The Daytona driver streams sandbox command output to the host. > - A log socket can stop delivering bytes without closing or rejecting. > - Missing permission requests and cancellation results can leave a run active and block queued work. > - This pull request switches an idle log stream to saved-output polling for the same command. > - The host can receive the missing output without another command dispatch. ## Linked Issues or Issue Description **What happened?** The driver waits for the SDK log-stream promise before it can start recovery. If that promise never settles, the host can miss new output that is already in the provider's saved logs. A run can remain active after a watch completes. Interrupt can then time out with “Execution is still stopping; termination has not been verified.” **Expected behavior** Recover command observation when the live stream stalls. Forward new permission requests and cancellation results. Require a recorded command exit before reporting completion. Keep quiet commands running under the caller's existing lifetime controls. **Steps to reproduce** 1. Start a session command and leave the callback log promise pending. 2. Put new output in the snapshot API without invoking the stream callback. 3. Keep the command running until the host receives that output, then expose its cancellation output and exit code. 4. Verify that the host receives each byte once and dispatches the command once. **Paperclip version or commit** Base commit `479554120d`. **Agent adapter(s) involved** Daytona session commands, including sandbox ACP agent sessions. Related: #14799 handles closed streams; this change handles sockets that never close. #14485 handles input delivery retries. #13262 adds permission diagnostics. ## What Changed - After 15 seconds with no stdout or stderr, switch directly to the existing status and log-snapshot polling path. - Ignore callbacks and delayed failures from the abandoned stream. Clear its idle timer on every exit path. - Preserve byte-offset deduplication, one command dispatch, and independent timeouts for recovery reads. - Add regressions for an initial stream stall, a stall after UTF-8 output, cancellation output, late callbacks, hung recovery reads, and quiet commands that outlive an operation timeout. - Update the provider documentation and keep hour-long healthy-stream coverage active with periodic output. ## Verification - `pnpm vitest run packages/plugins/sandbox-providers/daytona/src/plugin.test.ts`: 245 passed. - `pnpm exec vitest run --project @paperclipai/plugin-daytona`: 339 passed; 14 gated live tests skipped. - Both new stalled-stream regressions fail on the unchanged base driver because it never starts snapshot recovery. Both pass with this change. - `pnpm -r typecheck`: passed. - `pnpm build`: passed. - `pnpm test:run`: the local run did not pass. It was stopped after confirmed local skill-path and macOS skill-cache failures, once complete PR CI was green. Four chat/email tests could not load connector skill files from an ancestor directory outside the checkout. Three company-skills tests hit macOS `EACCES` during cache publication. One unrelated wakeup test timed out in the full run and passed on a focused rerun (`1 passed`, `27 skipped`). No source or test assertions were changed for these failures. This is not a complete local-suite pass. - Complete PR CI on `11ac4e030b`: 53 successful checks, 2 expected skips, no pending or failed checks. The clean CI run includes the full test suite. - Greptile reviewed `11ac4e030b` at 5/5 with no findings or unresolved threads. The branch has no merge conflicts with `master`. - `git diff --check` and a local scan for secrets and private identifiers passed. ## Risks - Quiet healthy commands also switch to polling. Full snapshots can increase bandwidth as output grows; polling remains limited to one snapshot per second. - The SDK exposes no stream cancellation handle. The old socket remains owned by session teardown, and its callbacks cannot publish after fallback. - This change recovers a stalled output stream. It does not claim to identify every cause of an unanswered permission request or change the requirement to verify termination before releasing work. - No schema migration or command replay. ## Model Used OpenAI GPT-6 through Codex, with tool use and code execution. The exact serving model ID and context window are not exposed in this session. ## 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: #` / `Refs: #` 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 (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run local change-specific tests; the full suite passes in CI, with local-suite limitations 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
261c24ccf9
commit
4e52463203
3 files changed
+160
-23
No files matched your search
@@ -30,7 +30,7 @@ Notes:
|
||||
- Each cold create uses a unique provider name and ownership labels. If creation fails after Daytona has allocated a sandbox, the driver looks up that exact name, checks every ownership label, and waits for deletion. Failed, missing, or timed-out lookups and failed deletions report unconfirmed cleanup with the provider name; they never return a usable lease. A missing name lookup after an uncertain create is not proof that a delayed provider request cannot create a resource.
|
||||
- For plugin-backed sandbox lease acquisition, unresolved creation cleanup crosses the worker RPC as a validated ownership envelope. The host records `pending_cleanup` before retrying deletion. If only a creation name is known, the provider first returns an ownership-verified sandbox ID without deleting it. The host saves that observation before authorizing deletion. Its existing cleanup sweep and durable spool preserve retries across controller restarts and environment deletion. A missing observed ID confirms cleanup after a lost deletion reply or database update; a name that has never been observed remains unresolved. Provider exceptions and resolved credentials are excluded from that envelope. This requires the matching host and plugin SDK update; probe and custom-image interactive-setup calls still report unconfirmed immediate cleanup without this lease-recovery path.
|
||||
- Reusable leases map to Daytona stop/start semantics. Non-reusable leases are deleted on release. A provider-resolved `target` does not change the identity of an existing sandbox. Release closes the same scoped lease that a later sentinel-verified resume reopens.
|
||||
- A session log socket can close while its command still runs. The driver requires a recorded command exit before returning a completed execution. It reconnects the log stream once, then polls status and log snapshots at most once per second. Replayed output is removed by byte offset, new output continues to reach the host, and the command is never dispatched again. Healthy initial and reconnected streams retain their existing lifetime under the caller’s RPC/run guard. After a clean close, each status or snapshot read gets a fresh operation timeout; successful recovery polling does not impose a new command lifetime. When both stream attempts fail without a clean close, fallback retains its original timeout budget starting after the stream attempts.
|
||||
- A session log socket can close while its command still runs. The driver requires a recorded command exit before returning a completed execution. It reconnects the log stream once, then polls status and log snapshots at most once per second. Replayed output is removed by byte offset, new output continues to reach the host, and the command is never dispatched again. A socket that delivers no output for 15 seconds switches directly to polling, even if the socket stays open or its initial connection never completes. Late callbacks from that socket are ignored. This recovers pending permission requests and cancellation results from saved output without replaying the command. Healthy initial and reconnected streams retain their existing lifetime under the caller’s RPC/run guard. After a clean close or an idle socket, each status or snapshot read gets a fresh operation timeout; successful recovery polling does not impose a new command lifetime. When both stream attempts fail without a clean close, fallback retains its original timeout budget starting after the stream attempts.
|
||||
- The SDK returns full log snapshots, so fallback bandwidth grows with retained output; it has no paged snapshot or stream cancellation API. A recovery-read timeout keeps partial output and explicitly reports whether command exit remains unconfirmed; it does not prove the remote command stopped. Provider/session teardown remains responsible for closing an outstanding socket. A final snapshot reconciles bytes missed by the stream before returning a recorded exit.
|
||||
- A sandbox record can survive the loss of its underlying container. Resume treats it as expired only when a fresh provider read confirms the exact missing-container error for that sandbox and marks it unrecoverable. Unknown errors and failed confirmation reads preserve the lease. The host still requires a verified native-runner backup before replacement.
|
||||
|
||||
|
||||
@@ -2437,6 +2437,116 @@ describe("Daytona sandbox provider plugin", () => {
|
||||
const streamExecParams = (overrides: Record<string, unknown> = {}) =>
|
||||
sessionExecParams(overrides);
|
||||
|
||||
// Long-lived healthy streams keep producing output. Quiet streams now
|
||||
// switch to snapshots, which have separate lifetime coverage below.
|
||||
async function streamFor(ms: number, onStderr?: (chunk: string) => void) {
|
||||
const timer = setInterval(() => onStderr?.("."), 10_000);
|
||||
try { await new Promise((resolve) => setTimeout(resolve, ms)); }
|
||||
finally { clearInterval(timer); }
|
||||
}
|
||||
|
||||
it.each(["", "prefix;🙂;"])("recovers unseen output and cancellation from a stuck socket after prefix %j", async (prefix) => {
|
||||
process.env.DAYTONA_API_KEY = "host-key";
|
||||
const executionLog = vi.fn();
|
||||
const restore = __setDaytonaPluginContextForTest(
|
||||
{ execution: { log: executionLog } } as unknown as PluginContext,
|
||||
);
|
||||
const sandbox = createMockSandbox();
|
||||
let lateOutput!: (chunk: string) => void;
|
||||
let rejectStream!: (error: Error) => void;
|
||||
let snapshots = 0;
|
||||
let cancelled = false;
|
||||
sandbox.process.getSessionCommandLogs.mockImplementation(
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void) => {
|
||||
if (onStdout) {
|
||||
lateOutput = onStdout;
|
||||
onStdout(prefix);
|
||||
return new Promise((_resolve, reject) => { rejectStream = reject; });
|
||||
}
|
||||
snapshots += 1;
|
||||
return { stdout: prefix + "permission;" + (cancelled ? "cancelled;" : ""), stderr: "" };
|
||||
},
|
||||
);
|
||||
sandbox.process.getSessionCommand.mockImplementation(async () => ({ exitCode: cancelled ? 0 : undefined }));
|
||||
mockGet.mockResolvedValue(sandbox);
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const promise = plugin.definition.onEnvironmentExecute?.(streamExecParams({ timeoutMs: 400 }));
|
||||
let settled = false;
|
||||
void promise?.then(() => { settled = true; });
|
||||
await vi.advanceTimersByTimeAsync(14_999);
|
||||
expect(snapshots).toBe(0);
|
||||
expect(settled).toBe(false);
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
expect(snapshots).toBe(1);
|
||||
expect(executionLog).toHaveBeenCalledWith("stdout", "permission;");
|
||||
// The abandoned socket cannot corrupt snapshot cursors or publish
|
||||
// output after the timeout, including a delayed rejection.
|
||||
lateOutput("stale bytes");
|
||||
rejectStream(new Error("late socket failure"));
|
||||
await vi.advanceTimersByTimeAsync(2_000);
|
||||
expect(settled).toBe(false);
|
||||
expect(snapshots).toBe(3);
|
||||
cancelled = true;
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
expect(await promise).toMatchObject({ exitCode: 0, timedOut: false, stdout: prefix + "permission;cancelled;" });
|
||||
lateOutput("after settlement");
|
||||
expect(executionLog.mock.calls).toEqual([
|
||||
...(prefix ? [["stdout", prefix]] : []), ["stdout", "permission;"], ["stdout", "cancelled;"],
|
||||
]);
|
||||
expect(sandbox.process.executeSessionCommand).toHaveBeenCalledTimes(1);
|
||||
expect(sandbox.process.getSessionCommandLogs.mock.calls.filter((call) => call[2])).toHaveLength(1);
|
||||
expect(vi.getTimerCount()).toBe(0);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
restore();
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["status", "snapshot"])("bounds a stuck %s read after an idle socket without claiming exit", async (stalledRead) => {
|
||||
process.env.DAYTONA_API_KEY = "host-key";
|
||||
const sandbox = createMockSandbox();
|
||||
let releaseRead!: () => void;
|
||||
const stalled = new Promise<void>((resolve) => { releaseRead = resolve; });
|
||||
const executionLog = vi.fn();
|
||||
const restore = __setDaytonaPluginContextForTest(
|
||||
{ execution: { log: executionLog } } as unknown as PluginContext,
|
||||
);
|
||||
sandbox.process.getSessionCommand.mockImplementation(async () => {
|
||||
if (stalledRead === "status") await stalled;
|
||||
return { exitCode: undefined };
|
||||
});
|
||||
sandbox.process.getSessionCommandLogs.mockImplementation(
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void) => {
|
||||
if (onStdout) {
|
||||
onStdout("prefix;");
|
||||
return new Promise(() => {});
|
||||
}
|
||||
await stalled;
|
||||
return { stdout: "prefix;late;", stderr: "" };
|
||||
},
|
||||
);
|
||||
mockGet.mockResolvedValue(sandbox);
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const promise = plugin.definition.onEnvironmentExecute?.(streamExecParams({ timeoutMs: 400 }));
|
||||
await vi.advanceTimersByTimeAsync(15_400);
|
||||
expect(await promise).toMatchObject({
|
||||
exitCode: null, timedOut: true, stdout: "prefix;",
|
||||
metadata: { timeoutScope: "session_log_observation", commandExitConfirmed: false },
|
||||
});
|
||||
releaseRead();
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
expect(executionLog.mock.calls).toEqual([["stdout", "prefix;"]]);
|
||||
expect(sandbox.process.executeSessionCommand).toHaveBeenCalledTimes(1);
|
||||
expect(vi.getTimerCount()).toBe(0);
|
||||
} finally {
|
||||
releaseRead();
|
||||
vi.useRealTimers();
|
||||
restore();
|
||||
}
|
||||
});
|
||||
|
||||
it("streams ordered stdout and stderr from the callback log form (test_log_stream_delivers_ordered_chunks)", async () => {
|
||||
process.env.DAYTONA_API_KEY = "host-key";
|
||||
const sandbox = createMockSandbox();
|
||||
@@ -2575,11 +2685,11 @@ describe("Daytona sandbox provider plugin", () => {
|
||||
const sandbox = createMockSandbox();
|
||||
let streams = 0;
|
||||
sandbox.process.getSessionCommandLogs.mockImplementation(
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void) => {
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void, onStderr?: (chunk: string) => void) => {
|
||||
if (!onStdout) return { stdout: "first;last;", stderr: "" };
|
||||
streams += 1;
|
||||
onStdout("first;");
|
||||
await new Promise((resolve) => setTimeout(resolve, 3_600_000));
|
||||
await streamFor(3_600_000, onStderr);
|
||||
if (streams === 1 && ending === "reject") throw new Error("socket error");
|
||||
if (streams === 2) onStdout("last;");
|
||||
},
|
||||
@@ -2613,10 +2723,10 @@ describe("Daytona sandbox provider plugin", () => {
|
||||
let snapshots = 0;
|
||||
let exited = false;
|
||||
sandbox.process.getSessionCommandLogs.mockImplementation(
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void) => {
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void, onStderr?: (chunk: string) => void) => {
|
||||
if (!onStdout) { snapshots += 1; return { stdout: "partial;", stderr: "" }; }
|
||||
onStdout("partial;");
|
||||
if (++streams === 1) await new Promise((resolve) => setTimeout(resolve, 3_600_000));
|
||||
if (++streams === 1) await streamFor(3_600_000, onStderr);
|
||||
},
|
||||
);
|
||||
sandbox.process.getSessionCommand.mockImplementation(async () => ({ exitCode: exited ? 0 : undefined }));
|
||||
@@ -2644,10 +2754,10 @@ describe("Daytona sandbox provider plugin", () => {
|
||||
const sandbox = createMockSandbox();
|
||||
let streams = 0;
|
||||
sandbox.process.getSessionCommandLogs.mockImplementation(
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void) => {
|
||||
async (_sid: string, _cmdId: string, onStdout?: (chunk: string) => void, onStderr?: (chunk: string) => void) => {
|
||||
if (onStdout) {
|
||||
onStdout("partial;");
|
||||
if (++streams === 1 && streamLifetime > 0) await new Promise((resolve) => setTimeout(resolve, streamLifetime));
|
||||
if (++streams === 1 && streamLifetime > 0) await streamFor(streamLifetime, onStderr);
|
||||
throw new Error("socket error");
|
||||
}
|
||||
return { stdout: "partial;", stderr: "" };
|
||||
|
||||
@@ -1775,6 +1775,9 @@ async function executeOneShot(
|
||||
// Fallback reads full log snapshots as well as status. Limit the heavier
|
||||
// polling path while keeping output available to interactive commands.
|
||||
const SESSION_LOG_POLL_INTERVAL_MS = 1000;
|
||||
// A socket can stay open after it stops delivering output. Switch to saved
|
||||
// logs after this much silence; silence never proves that the command exited.
|
||||
const SESSION_LOG_STREAM_IDLE_MS = 15_000;
|
||||
|
||||
function sleep(ms: number): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
@@ -1907,7 +1910,7 @@ function createSessionStreamBuffer(
|
||||
|
||||
type SessionLogStreamResult = {
|
||||
exitCode: number | null;
|
||||
closedStream: boolean;
|
||||
pollUntilExit: boolean;
|
||||
buffer: ReturnType<typeof createSessionStreamBuffer>;
|
||||
};
|
||||
|
||||
@@ -1945,8 +1948,8 @@ async function observeSessionCommand<T>(timeoutMs: number, action: () => Promise
|
||||
|
||||
// Stream stdout and stderr of one session command from the callback log form.
|
||||
// The stream buffer drops a replayed prefix by byte offset on a reconnect.
|
||||
// Both a rejected stream and a clean close without a confirmed process exit
|
||||
// use bounded reconnects, then polling.
|
||||
// Rejected streams and clean closes use bounded reconnects, then polling.
|
||||
// An idle socket goes directly to polling, including one that never settles.
|
||||
async function runSessionLogStream(
|
||||
sandbox: Sandbox,
|
||||
sessionId: string,
|
||||
@@ -1955,32 +1958,56 @@ async function runSessionLogStream(
|
||||
buffer: ReturnType<typeof createSessionStreamBuffer>,
|
||||
): Promise<SessionLogStreamResult> {
|
||||
let reconnects = 0;
|
||||
let closedStream = false;
|
||||
let pollUntilExit = false;
|
||||
while (true) {
|
||||
let streamClosed = false;
|
||||
let acceptingOutput = true;
|
||||
let idleTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
let signalIdle!: () => void;
|
||||
const idle = new Promise<"idle">((resolve) => { signalIdle = () => resolve("idle"); });
|
||||
const armIdleTimer = () => {
|
||||
clearTimeout(idleTimer);
|
||||
idleTimer = setTimeout(signalIdle, SESSION_LOG_STREAM_IDLE_MS);
|
||||
};
|
||||
const receive = (stream: "stdout" | "stderr", chunk: string) => {
|
||||
if (!acceptingOutput || !chunk) return;
|
||||
armIdleTimer();
|
||||
if (stream === "stdout") buffer.onStdout(chunk);
|
||||
else buffer.onStderr(chunk);
|
||||
};
|
||||
try {
|
||||
await sandbox.process.getSessionCommandLogs(
|
||||
sessionId,
|
||||
commandId,
|
||||
(chunk) => { if (acceptingOutput) buffer.onStdout(chunk); },
|
||||
(chunk) => { if (acceptingOutput) buffer.onStderr(chunk); },
|
||||
);
|
||||
armIdleTimer();
|
||||
const outcome = await Promise.race([
|
||||
sandbox.process.getSessionCommandLogs(
|
||||
sessionId,
|
||||
commandId,
|
||||
(chunk) => receive("stdout", chunk),
|
||||
(chunk) => receive("stderr", chunk),
|
||||
).then(() => "closed" as const),
|
||||
idle,
|
||||
]);
|
||||
if (outcome === "idle") {
|
||||
// Do not wait for the abandoned socket or open another one. The SDK
|
||||
// exposes no cancellation handle. Fence its callbacks before the
|
||||
// snapshot path resets the buffer's byte cursors.
|
||||
return { exitCode: null, pollUntilExit: true, buffer };
|
||||
}
|
||||
streamClosed = true;
|
||||
closedStream = true;
|
||||
pollUntilExit = true;
|
||||
} catch {
|
||||
// The command can still be running after a log transport failure.
|
||||
} finally {
|
||||
// The SDK exposes no cancellation handle for this socket. Session
|
||||
// teardown owns closing it; late data cannot change a settled execution.
|
||||
acceptingOutput = false;
|
||||
clearTimeout(idleTimer);
|
||||
}
|
||||
if (streamClosed) {
|
||||
const exitCode = await readSessionExitCode(sandbox, sessionId, commandId, timeoutMs);
|
||||
if (exitCode !== null) return { exitCode, closedStream, buffer };
|
||||
if (exitCode !== null) return { exitCode, pollUntilExit, buffer };
|
||||
}
|
||||
if (reconnects >= MAX_SESSION_STREAM_RECONNECTS) {
|
||||
return { exitCode: null, closedStream, buffer };
|
||||
return { exitCode: null, pollUntilExit, buffer };
|
||||
}
|
||||
reconnects += 1;
|
||||
buffer.resetConnectionCursors();
|
||||
@@ -2110,14 +2137,14 @@ async function executeInSession(
|
||||
// Keep forwarding output while waiting for the actual command exit. ACP
|
||||
// peers may need a tool response before they can finish the command. Waiting
|
||||
// for exit before reading logs would strand those peers after a stream loss.
|
||||
// An EOF after a live stream does not end the command's lifetime. Bound
|
||||
// A closed or idle log socket does not end the command's lifetime. Bound
|
||||
// each recovery read independently; the caller still owns stop/teardown.
|
||||
// Preserve the legacy budget starting at fallback entry when both stream
|
||||
// attempts fail, including a stream that fails after running for hours.
|
||||
const deadlineMs = Date.now() + effectiveTimeoutMs;
|
||||
const observe = async <T>(action: () => Promise<T>): Promise<T> => {
|
||||
try {
|
||||
return await (streamResult.closedStream
|
||||
return await (streamResult.pollUntilExit
|
||||
? observeSessionCommand(effectiveTimeoutMs, action)
|
||||
: beforeSessionDeadline(deadlineMs, action));
|
||||
} catch (error) {
|
||||
@@ -2146,7 +2173,7 @@ async function executeInSession(
|
||||
// Full log snapshots are heavier than a status read. Limit fallback
|
||||
// traffic to one snapshot per second. Failed-stream fallback retains
|
||||
// its original deadline, including the time spent waiting between reads.
|
||||
await sleep(streamResult.closedStream ? SESSION_LOG_POLL_INTERVAL_MS
|
||||
await sleep(streamResult.pollUntilExit ? SESSION_LOG_POLL_INTERVAL_MS
|
||||
: Math.max(0, Math.min(SESSION_LOG_POLL_INTERVAL_MS, deadlineMs - Date.now())));
|
||||
}
|
||||
} catch (error) {
|
||||
|
||||
Reference in new issue
Block a user