diff --git a/packages/plugins/sandbox-providers/daytona/README.md b/packages/plugins/sandbox-providers/daytona/README.md index 967c545207..b384c57549 100644 --- a/packages/plugins/sandbox-providers/daytona/README.md +++ b/packages/plugins/sandbox-providers/daytona/README.md @@ -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. diff --git a/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts b/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts index 63369685be..ac3a7c1dd0 100644 --- a/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts +++ b/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts @@ -2437,6 +2437,116 @@ describe("Daytona sandbox provider plugin", () => { const streamExecParams = (overrides: Record = {}) => 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((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: "" }; diff --git a/packages/plugins/sandbox-providers/daytona/src/plugin.ts b/packages/plugins/sandbox-providers/daytona/src/plugin.ts index 347564f2f0..4dfa65c9d5 100644 --- a/packages/plugins/sandbox-providers/daytona/src/plugin.ts +++ b/packages/plugins/sandbox-providers/daytona/src/plugin.ts @@ -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 { return new Promise((resolve) => setTimeout(resolve, ms)); @@ -1907,7 +1910,7 @@ function createSessionStreamBuffer( type SessionLogStreamResult = { exitCode: number | null; - closedStream: boolean; + pollUntilExit: boolean; buffer: ReturnType; }; @@ -1945,8 +1948,8 @@ async function observeSessionCommand(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, ): Promise { let reconnects = 0; - let closedStream = false; + let pollUntilExit = false; while (true) { let streamClosed = false; let acceptingOutput = true; + let idleTimer: ReturnType | 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 (action: () => Promise): Promise => { 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) {