From 8de942552f278039263f15f27a9aa1daf9bb5654 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Mon, 5 Oct 2026 10:08:17 +0200 Subject: [PATCH] fix(sandbox): bind Docker controls and cleanup to run ownership --- server/routes/simulation.ws.ts | 44 +++- server/services/sandbox-runner-pool.ts | 65 ++++-- server/services/sandbox-runner.ts | 135 ++++++++----- server/services/sandbox/docker-manager.ts | 4 + server/services/sandbox/execution-manager.ts | 6 +- .../sandbox/execution-phases/cleanup-phase.ts | 44 +++- .../sandbox/execution-phases/timeout-phase.ts | 4 +- .../docker-security-contract.test.ts | 42 ++++ tests/server/pause-resume-timing.test.ts | 6 +- .../websocket-lifecycle-handlers.test.ts | 35 ++++ .../sandbox-runner-control-ownership.test.ts | 191 ++++++++++++++++++ .../services/sandbox-runner-pool.test.ts | 33 +++ .../services/sandbox/cleanup-phase.test.ts | 2 +- .../services/sandbox/timeout-phase.test.ts | 2 +- 14 files changed, 527 insertions(+), 86 deletions(-) create mode 100644 tests/server/services/sandbox-runner-control-ownership.test.ts diff --git a/server/routes/simulation.ws.ts b/server/routes/simulation.ws.ts index 686ac7814..ab3ac3ca6 100644 --- a/server/routes/simulation.ws.ts +++ b/server/routes/simulation.ws.ts @@ -68,16 +68,40 @@ function rejectAdmission( } } +/** Preserve the existing serial diagnostic and stopped-status protocol. */ +async function releaseFailedControl( + ws: WebSocket, clientState: ClientState, sessionManager: WsSessionManager, + operation: "pause" | "resume", +): Promise { + sendMessageToClient(ws, { + type: WSMessageType.SERIAL_OUTPUT, + data: `[ERR] Docker ${operation} failed; simulation stopped.\n`, + }); + const reservation = clientState.reservation; + await sessionManager.safeReleaseRunner(clientState, `${operation}_failed`, reservation); + // A release can finish after a new run has already reserved this session. + if (clientState.reservation && clientState.reservation !== reservation) return; + sendMessageToClient(ws, { type: WSMessageType.SIMULATION_STATUS, status: "stopped" }); +} + /** * Handle "pause_simulation" WebSocket message */ -function handlePauseSimulation( +async function handlePauseSimulation( _ws: WebSocket, clientState: ClientState, sessionManager: WsSessionManager, -): void { +): Promise { if (clientState?.runner && clientState.isRunning) { - const paused = clientState.runner.pause(); + const runner = clientState.runner; + const reservation = clientState.reservation; + if (!runner.pause()) return; + const paused = runner.controlResult ? await runner.controlResult : true; + if (clientState.runner !== runner || clientState.reservation !== reservation) return; + if (!paused) { + await releaseFailedControl(_ws, clientState, sessionManager, "pause"); + return; + } if (paused) { clientState.isPaused = true; sessionManager.markSessionPaused(clientState); @@ -96,13 +120,21 @@ function handlePauseSimulation( /** * Handle "resume_simulation" WebSocket message */ -function handleResumeSimulation( +async function handleResumeSimulation( _ws: WebSocket, clientState: ClientState, sessionManager: WsSessionManager, -): void { +): Promise { if (clientState?.runner && clientState.isPaused) { - const resumed = clientState.runner.resume(); + const runner = clientState.runner; + const reservation = clientState.reservation; + if (!runner.resume()) return; + const resumed = runner.controlResult ? await runner.controlResult : true; + if (clientState.runner !== runner || clientState.reservation !== reservation) return; + if (!resumed) { + await releaseFailedControl(_ws, clientState, sessionManager, "resume"); + return; + } if (resumed) { clientState.isPaused = false; clientState.isRunning = true; diff --git a/server/services/sandbox-runner-pool.ts b/server/services/sandbox-runner-pool.ts index 6204eea5c..bfcc7703e 100644 --- a/server/services/sandbox-runner-pool.ts +++ b/server/services/sandbox-runner-pool.ts @@ -6,6 +6,7 @@ interface PooledRunner { runner: SandboxRunner; inUse: boolean; resetting: boolean; + quarantined?: boolean; lastReleasedTime: number; idleTimer: ReturnType | null; } @@ -36,6 +37,7 @@ export class SandboxRunnerPool { private readonly acquireTimeoutMs: number; private readonly resetTimeoutMs: number; private initialized = false; + private shuttingDown = false; constructor(options: SandboxRunnerPoolOptions = {}) { this.minRunners = options.minRunners ?? config.sandbox.pool.minRunners; @@ -179,7 +181,7 @@ export class SandboxRunnerPool { return; } - if (!pooledRunner.inUse) { + if (!pooledRunner.inUse && !pooledRunner.quarantined) { this.logger.warn( "[SandboxRunnerPool] Attempt to release already-released runner (ignored)", ); @@ -187,7 +189,12 @@ export class SandboxRunnerPool { } // Always mark as free FIRST, even if reset hangs — prevents permanent pool deadlock + if (pooledRunner.idleTimer !== null) { + clearTimeout(pooledRunner.idleTimer); + pooledRunner.idleTimer = null; + } pooledRunner.inUse = false; + pooledRunner.quarantined = false; pooledRunner.resetting = true; pooledRunner.lastReleasedTime = Date.now(); @@ -206,23 +213,9 @@ export class SandboxRunnerPool { pooledRunner.resetting = false; } catch (error) { this.logger.error( - `[SandboxRunnerPool] Runner reset failed or timed out: ${error}. Force-replacing runner.`, + `[SandboxRunnerPool] Runner reset failed or timed out: ${error}.`, ); - // Replace the stuck runner with a fresh one - const index = this.runners.indexOf(pooledRunner); - if (index !== -1) { - const freshRunner = new SandboxRunner(); - this.runners[index] = { - runner: freshRunner, - inUse: false, - resetting: false, - lastReleasedTime: Date.now(), - idleTimer: null, - }; - this.logger.info( - `[SandboxRunnerPool] Replaced stuck runner at index ${index} with fresh instance`, - ); - } + if (this.recoverFailedReset(pooledRunner)) return; } finally { if (resetTimer !== undefined) { clearTimeout(resetTimer); @@ -257,6 +250,41 @@ export class SandboxRunnerPool { } } + private recoverFailedReset(pooledRunner: PooledRunner): boolean { + if (pooledRunner.runner.hasPendingContainerCleanup) { + this.quarantineForCleanup(pooledRunner); + return true; + } + // Replace the stuck runner with a fresh one + const index = this.runners.indexOf(pooledRunner); + if (index !== -1) { + const freshRunner = new SandboxRunner(); + this.runners[index] = { + runner: freshRunner, + inUse: false, + resetting: false, + lastReleasedTime: Date.now(), + idleTimer: null, + }; + this.logger.info( + `[SandboxRunnerPool] Replaced stuck runner at index ${index} with fresh instance`, + ); + } + return false; + } + + private quarantineForCleanup(pooledRunner: PooledRunner): void { + // Retain ownership in this bounded slot and retry using the existing + // idle-maintenance interval. Never offer it until reset confirms cleanup. + pooledRunner.quarantined = true; + if (this.shuttingDown) return; + pooledRunner.idleTimer = setTimeout(() => { + pooledRunner.idleTimer = null; + void this.releaseRunner(pooledRunner.runner); + }, this.idleTimeoutMs); + pooledRunner.idleTimer.unref(); + } + /** * Schedule idle cleanup for logical runners above the derived warm floor. * If the runner is re-acquired before the timer fires, the timer is cancelled. @@ -317,6 +345,7 @@ export class SandboxRunnerPool { } async shutdown(): Promise { + this.shuttingDown = true; this.logger.info("[SandboxRunnerPool] Shutting down..."); for (const entry of this.queue) { @@ -335,7 +364,7 @@ export class SandboxRunnerPool { for (const { runner } of this.runners) { try { - if (runner.isRunning) { + if (runner.isRunning || runner.hasPendingContainerCleanup) { await runner.stop(); } } catch (error) { diff --git a/server/services/sandbox-runner.ts b/server/services/sandbox-runner.ts index ca8e7b760..e75b20e27 100644 --- a/server/services/sandbox-runner.ts +++ b/server/services/sandbox-runner.ts @@ -20,7 +20,7 @@ import { DockerManager } from "./sandbox/docker-manager"; import { StreamHandler } from "./sandbox/stream-handler"; import { FilesystemHelper } from "./sandbox/filesystem-helper"; import { ExecutionManager, type ExecutionState, SimulationState, SANDBOX_CONFIG } from "./sandbox/execution-manager"; -import { flushMessageQueue } from "./sandbox/execution-phases/cleanup-phase"; +import { cleanupExecutionContainer, flushMessageQueue } from "./sandbox/execution-phases/cleanup-phase"; import { config } from "../config"; export class SandboxRunner { @@ -40,6 +40,18 @@ export class SandboxRunner { private readonly executionState: ExecutionState; private readonly processExecutor: ProcessExecutor; + private controlCompletion?: Promise; + private controlPending = false; + + /** Completion of the last accepted control; Stop never waits for it. */ + get controlResult(): Promise | undefined { + return this.controlCompletion; + } + + get hasPendingContainerCleanup(): boolean { + return Boolean(this.executionState.currentContainerName); + } + private dockerAvailable = false; private dockerImageBuilt = false; private dockerChecked = false; @@ -92,6 +104,7 @@ export class SandboxRunner { this.executionState.backpressurePaused = streamState.backpressurePaused; } }, + () => { this.executionState.terminationRequested = true; }, ); this.streamHandler = new StreamHandler(this.processController); this.filesystemHelper = new FilesystemHelper(this.fileBuilder, this.localCompiler); @@ -283,21 +296,42 @@ export class SandboxRunner { private async cleanupDockerContainer(containerName?: string): Promise { if (!containerName) return; + const s = this.executionState; + await cleanupExecutionContainer(s, { + processExecutor: this.processExecutor, logger: this.logger, + }); + } - try { - const result = await this.processExecutor.execute("docker", ["rm", "-f", containerName], { - timeout: 5000, - stdio: "pipe", - }); - this.logger.info(`Docker cleanup for ${containerName} finished (code ${result.code})`); - } catch (error) { - this.logger.debug(`Docker cleanup for ${containerName} failed: ${error}`); - } + private controlDocker(operation: "pause" | "unpause", commit: () => void): void { + const s = this.executionState; + const containerName = s.currentContainerName!; + const generation = s.runGeneration; + const expectedState = this.state; + const abort = s.runAbort; + const isCurrent = () => s.runGeneration === generation && s.runAbort === abort && + !abort?.signal.aborted && !s.processKilled && s.currentContainerName === containerName && + !s.terminationRequested && this.state === expectedState && this.processController.hasProcess(); + this.controlPending = true; + const finish = async (result: Awaited>): Promise => { + if (!isCurrent()) return false; + this.controlPending = false; + if (result.code !== 0 || result.error) { + this.logger.warn(`Docker ${operation} failed for ${containerName} (code ${result.code}): ${String(result.error ?? "nonzero exit")}`); + await this.stop(); + return false; + } + commit(); + return true; + }; + this.controlCompletion = this.processExecutor.execute("docker", [operation, containerName], { + timeout: 5000, stdio: "pipe", + }).then(finish, (error: unknown) => finish({ + code: -1, error: error instanceof Error ? error : new Error(String(error)), + })); } - pause(): boolean { + private commitPause(): void { const s = this.executionState; - if (this.state !== SimulationState.RUNNING || !this.processController.hasProcess()) return false; this.state = SimulationState.PAUSED; this.timeoutManager.pause(); s.pinStateBatcher?.pause(); @@ -306,47 +340,26 @@ export class SandboxRunner { if (!s.processKilled) this.processController.writeStdin("[[PAUSE_TIME]]\n"); s.pauseStartTime = Date.now(); this.registryManager.markPauseTime(s.pauseStartTime); - // SIGSTOP only suspends the local `docker run` client process; the - // container itself would continue producing output into the pipe. Pause - // the container when running in Docker so the sketch really stops at the - // exact instruction boundary and no ten-second backlog accumulates. - if (s.currentContainerName) { - const containerName = s.currentContainerName; - void this.processExecutor.execute("docker", ["pause", containerName], { - timeout: 5000, - stdio: "pipe", - }).catch((error) => { - this.logger.warn(`Docker pause failed for ${containerName}: ${error instanceof Error ? error.message : String(error)}`); - }); + this.logger.info("Simulation paused"); + } + + pause(): boolean { + if (this.controlPending || this.executionState.terminationRequested || this.state !== SimulationState.RUNNING || !this.processController.hasProcess()) return false; + if (this.executionState.currentContainerName) { + this.controlDocker("pause", () => this.commitPause()); } else { this.processController.kill("SIGSTOP"); + this.commitPause(); + this.controlCompletion = undefined; } - this.logger.info("Simulation paused (SIGSTOP)"); return true; } - resume(): boolean { + private commitResume(): void { const s = this.executionState; - if (this.state !== SimulationState.PAUSED || !this.processController.hasProcess()) return false; - const containerName = s.currentContainerName; - if (!containerName) this.processController.kill("SIGCONT"); const pauseDuration = Date.now() - (this.pauseStartTime ?? Date.now()); s.totalPausedTime += pauseDuration; - const resumeContainer = containerName - ? this.processExecutor.execute("docker", ["unpause", containerName], { - timeout: 5000, - stdio: "pipe", - }) - : Promise.resolve(); - void resumeContainer - .catch((error) => { - this.logger.warn(`Docker resume failed for ${containerName}: ${error instanceof Error ? error.message : String(error)}`); - }) - .finally(() => { - if (!s.processKilled) { - this.processController.writeStdin(`[[RESUME_TIME:${pauseDuration}]]\n`); - } - }); + if (!s.processKilled) this.processController.writeStdin(`[[RESUME_TIME:${pauseDuration}]]\n`); s.pauseStartTime = null; this.registryManager.markPauseTime(null); this.state = SimulationState.RUNNING; @@ -354,15 +367,24 @@ export class SandboxRunner { s.pinStateBatcher?.resume(); s.serialOutputBatcher?.resume(); this.registryManager.resumeTelemetry(); - this.logger.info(`Simulation resumed after ${pauseDuration}ms (SIGCONT)`); - if (!containerName && !s.processKilled) this.processController.writeStdin("\n"); + this.logger.info(`Simulation resumed after ${pauseDuration}ms`); + if (!s.currentContainerName && !s.processKilled) this.processController.writeStdin("\n"); if (s.outputBuffer.length > 0 && s.onOutputCallback && !s.isSendingOutput) { this.sendOutputWithDelay(s.onOutputCallback); } - return true; } - + resume(): boolean { + if (this.controlPending || this.executionState.terminationRequested || this.state !== SimulationState.PAUSED || !this.processController.hasProcess()) return false; + if (this.executionState.currentContainerName) { + this.controlDocker("unpause", () => this.commitResume()); + } else { + this.processController.kill("SIGCONT"); + this.commitResume(); + this.controlCompletion = undefined; + } + return true; + } sendSerialInput(input: string): void { const s = this.executionState; @@ -400,7 +422,11 @@ export class SandboxRunner { const s = this.executionState; // Cancels a run that is still preparing or waiting for a start slot. s.runAbort?.abort(); - if (this.state === SimulationState.STOPPED || s.processKilled) return; + this.controlPending = false; + if (this.state === SimulationState.STOPPED || s.processKilled) { + await this.cleanupDockerContainer(s.currentContainerName); + return; + } this.state = SimulationState.STOPPED; s.processKilled = true; s.pendingCleanup = true; @@ -444,7 +470,6 @@ export class SandboxRunner { s.outputBuffer = ""; s.outputBufferIndex = 0; s.isSendingOutput = false; const containerName = s.currentContainerName; - s.currentContainerName = undefined; if (s.flushTimer) { clearTimeout(s.flushTimer); s.flushTimer = null; } await this.cleanupDockerContainer(containerName); @@ -464,10 +489,18 @@ export class SandboxRunner { } } - this.processController.clearListeners(); const s = this.executionState; + s.runAbort?.abort(); + this.controlPending = false; + this.controlCompletion = undefined; + if (s.currentContainerName) { + await this.cleanupDockerContainer(s.currentContainerName); + if (s.currentContainerName) throw new Error(`Docker cleanup unconfirmed for ${s.currentContainerName}`); + } + this.processController.clearListeners(); s.state = SimulationState.STOPPED; s.processKilled = false; + s.terminationRequested = false; s.pauseStartTime = null; s.totalPausedTime = 0; s.pinStateBatcher = null; diff --git a/server/services/sandbox/docker-manager.ts b/server/services/sandbox/docker-manager.ts index 8d2c27696..e1bd2fe30 100644 --- a/server/services/sandbox/docker-manager.ts +++ b/server/services/sandbox/docker-manager.ts @@ -81,12 +81,14 @@ export class DockerManager { private readonly stderrParser: ArduinoOutputParser, private readonly timeoutManager: SimulationTimeoutManager, private readonly handleParsedLine: HandleParsedLineDelegate, + private readonly onTerminationRequested?: () => void, ) {} private consumeOutputBudget(state: OutputBudgetState, data: Buffer | string, callbacks: DockerManagerCallbacks): boolean { const counter = state.totalOutputBytes; counter.value += Buffer.byteLength(data); if (counter.value <= this.SANDBOX_CONFIG.maxOutputBytes) return true; + this.onTerminationRequested?.(); this.processController.kill("SIGKILL"); callbacks.onError("Output size limit exceeded"); return false; @@ -100,6 +102,7 @@ export class DockerManager { const timeoutSec = normalizeSimulationTimeout(executionTimeout); const handleTimeout = () => { + this.onTerminationRequested?.(); this.processController.kill("SIGKILL"); callbacks.onOutput(`--- Simulation timeout (${timeoutSec}s) ---`, true); this.logger.info(`Docker runtime timeout after ${timeoutSec}s`); @@ -115,6 +118,7 @@ export class DockerManager { private setupDockerStartupTimeout(callbacks: DockerManagerCallbacks): void { const startupTimeoutSec = this.SANDBOX_CONFIG.maxExecutionTimeSec; this.timeoutManager.schedule(startupTimeoutSec * 1000, () => { + this.onTerminationRequested?.(); this.processController.kill("SIGKILL"); callbacks.onOutput(`--- Sandbox startup timeout (${startupTimeoutSec}s) ---`, true); this.logger.warn(`Docker sandbox startup timeout after ${startupTimeoutSec}s`); diff --git a/server/services/sandbox/execution-manager.ts b/server/services/sandbox/execution-manager.ts index be4bf22ec..bc7fe40c6 100644 --- a/server/services/sandbox/execution-manager.ts +++ b/server/services/sandbox/execution-manager.ts @@ -24,7 +24,7 @@ import { normalizeBaudrate, normalizeSimulationTimeout } from "@shared/input-lim import { canTransition } from "../simulation-state-machine"; import { createProcessExecutionPort, type ProcessExecution } from "../process-execution-port"; import { OutputCollector } from "../output-collector"; -import { flushMessageQueue, flushBatchers, cleanupDockerContainer } from "./execution-phases/cleanup-phase"; +import { flushMessageQueue, flushBatchers, cleanupExecutionContainer } from "./execution-phases/cleanup-phase"; import { scheduleExecutionTimeout } from "./execution-phases/timeout-phase"; import { createStreamCallbacks, delegateParsedLineToStreamHandler, handleStderrFallbackData } from "./execution-phases/stream-phase"; import { runLocalStart, runDockerStart, type LocalStartContext, type DockerStartContext, type DockerStartParams, type TransitionToFn } from "./execution-phases/start-phase"; @@ -130,6 +130,7 @@ export interface ExecutionState { pendingCleanup: boolean; processController: IProcessController; currentContainerName?: string; + terminationRequested?: boolean; dockerAvailable?: boolean; dockerImageBuilt?: boolean; outputCollector?: OutputCollector; @@ -332,6 +333,7 @@ export class ExecutionManager { const generation = (state.runGeneration ?? 0) + 1; state.runAbort = abort; state.runGeneration = generation; + state.terminationRequested = false; return { signal: abort.signal, isStale: () => abort.signal.aborted || state.runGeneration !== generation, @@ -571,7 +573,7 @@ export class ExecutionManager { clearTimeout(state.flushTimer); state.flushTimer = null; } - void cleanupDockerContainer(state.currentContainerName, { + void cleanupExecutionContainer(state, { processExecutor: this.processExecutor, logger: this.logger, }); diff --git a/server/services/sandbox/execution-phases/cleanup-phase.ts b/server/services/sandbox/execution-phases/cleanup-phase.ts index be0e00916..034b17351 100644 --- a/server/services/sandbox/execution-phases/cleanup-phase.ts +++ b/server/services/sandbox/execution-phases/cleanup-phase.ts @@ -53,18 +53,56 @@ export function flushBatchers(state: ExecutionState): void { export async function cleanupDockerContainer( containerName: string | undefined, deps: CleanupDependencies, -): Promise { +): Promise { if (!containerName) { - return; + return true; } try { - await deps.processExecutor.execute("docker", ["rm", "-f", containerName], { + const result = await deps.processExecutor.execute("docker", ["rm", "-f", containerName], { timeout: 5000, stdio: "pipe", }); + if (result.code !== 0 || result.error) { + // docker run uses --rm. A failed rm can mean auto-removal, but only a + // successful daemon listing can establish absence (not an error string). + // A name filter can conservatively match more than the exact name; + // an empty successful listing still proves this container is absent. + const remaining = await deps.processExecutor.execute("docker", [ + "container", "ls", "--all", "--filter", `name=${containerName}`, "--quiet", + ], { timeout: 5000, stdio: "pipe" }); + if (remaining.code !== 0 || remaining.error || remaining.stdout?.trim() !== "") { + deps.logger.warn(`Docker cleanup failed for ${containerName} (code ${result.code})`); + return false; + } + } deps.logger.info(`Docker container cleanup: ${containerName}`); + return true; } catch (error) { deps.logger.debug(`Docker cleanup failed for ${containerName}: ${error}`); + return false; } } + + +// Natural close, deadline and explicit Stop must join the same removal. The +// weak key retains no completed sessions; each runner owns at most one flight. +const cleanupFlights = new WeakMap }>(); + +export function cleanupExecutionContainer(state: ExecutionState, deps: CleanupDependencies): Promise { + const name = state.currentContainerName; + if (!name) return Promise.resolve(); + const pending = cleanupFlights.get(state); + if (pending?.name === name) return pending.completion; + const generation = state.runGeneration; + const flight = { name, completion: Promise.resolve() }; + flight.completion = cleanupDockerContainer(name, deps).then((removed) => { + if (removed && state.runGeneration === generation && state.currentContainerName === name) { + state.currentContainerName = undefined; + } + }).finally(() => { + if (cleanupFlights.get(state) === flight) cleanupFlights.delete(state); + }); + cleanupFlights.set(state, flight); + return flight.completion; +} diff --git a/server/services/sandbox/execution-phases/timeout-phase.ts b/server/services/sandbox/execution-phases/timeout-phase.ts index 61e04ef3a..3f6e971ed 100644 --- a/server/services/sandbox/execution-phases/timeout-phase.ts +++ b/server/services/sandbox/execution-phases/timeout-phase.ts @@ -8,7 +8,7 @@ import type { Logger } from "@shared/logger"; import type { ProcessExecution } from "../../process-execution-port"; import type { ExecutionState } from "../execution-manager"; -import { cleanupDockerContainer } from "./cleanup-phase"; +import { cleanupExecutionContainer } from "./cleanup-phase"; interface TimeoutScheduler { schedule(timeoutMs: number | null, callback: () => void): void; @@ -42,7 +42,7 @@ export function handleExecutionTimeout( abortExecution(state); callbacks.onOutput(`--- Simulation timeout (${executionTimeout}s) ---`, true); - void cleanupDockerContainer(state.currentContainerName, deps); + void cleanupExecutionContainer(state, deps); } /** diff --git a/tests/integration/docker-security-contract.test.ts b/tests/integration/docker-security-contract.test.ts index 43bead071..d290329fc 100644 --- a/tests/integration/docker-security-contract.test.ts +++ b/tests/integration/docker-security-contract.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it } from "vitest"; import { existsSync } from "node:fs"; import { join } from "node:path"; import { SandboxRunner } from "../../server/services/sandbox-runner"; +import { SandboxRunnerPool } from "../../server/services/sandbox-runner-pool"; import { ProcessExecutor } from "../../server/services/process-executor"; import { extractPlainText, @@ -47,6 +48,47 @@ async function runSecurityProbe(sketch: string): Promise { } maybeDescribe("Docker sandbox security contract", () => { + it("confirms pause/resume and reuses the same runner after natural auto-removal", async () => { + const pool = new SandboxRunnerPool({ minRunners: 1, maxRunners: 1 }); + await pool.initialize(); + const runner = await pool.acquireRunner(); + const executor = new ProcessExecutor(); + try { + expect(await runner.runSketch({ + code: "void setup() {} void loop() { delay(100); }", timeoutSec: 10, + onOutput: () => {}, onError: () => {}, onExit: () => {}, + })).toBe(true); + const name = await waitForContainerName(runner); + expect(runner.pause()).toBe(true); + expect(await runner.controlResult).toBe(true); + const paused = await executor.execute("docker", ["inspect", name], { timeout: 5000 }); + expect(paused.code).toBe(0); + expect(JSON.parse(paused.stdout ?? "[]")[0].State.Paused).toBe(true); + expect(runner.resume()).toBe(true); + expect(await runner.controlResult).toBe(true); + const resumed = await executor.execute("docker", ["inspect", name], { timeout: 5000 }); + expect(resumed.code).toBe(0); + expect(JSON.parse(resumed.stdout ?? "[]")[0].State.Paused).toBe(false); + await runner.stop(); + await pool.releaseRunner(runner); + expect(await pool.acquireRunner()).toBe(runner); + + let exit!: () => void; + const exited = new Promise((resolve) => { exit = resolve; }); + await runner.runSketch({ + code: "#include \nvoid setup() {} void loop() { exit(0); }", timeoutSec: 5, + onOutput: () => {}, onError: () => {}, onExit: () => exit(), + }); + await exited; + await pool.releaseRunner(runner); + expect(runner.hasPendingContainerCleanup).toBe(false); + expect(await pool.acquireRunner()).toBe(runner); + } finally { + await runner.stop(); + await pool.shutdown(); + } + }, 45_000); + it("applies isolation options to a real running container", async () => { const runner = new SandboxRunner(); const runPromise = runner.runSketch({ diff --git a/tests/server/pause-resume-timing.test.ts b/tests/server/pause-resume-timing.test.ts index f53bfcbbb..dfeb13e93 100644 --- a/tests/server/pause-resume-timing.test.ts +++ b/tests/server/pause-resume-timing.test.ts @@ -41,7 +41,7 @@ maybeDescribe("SandboxRunner - Pause/Resume Timing", () => { runner.runSketch({ code, - onOutput: (line) => { + onOutput: async (line) => { const matchRe = /TIME:(\d+)/; const match = matchRe.exec(line); if (match) { @@ -50,7 +50,8 @@ maybeDescribe("SandboxRunner - Pause/Resume Timing", () => { if (timeValues.length === 5) { runner.pause(); - const valAtPause = t; + await runner.controlResult; + const valAtPause = timeValues.at(-1) ?? t; // Wir warten 500ms in der "echten" Welt setTimeout(() => { @@ -171,6 +172,7 @@ maybeDescribe("SandboxRunner - Pause/Resume Timing", () => { clearInterval(check); try { runner.pause(); + await runner.controlResult; expect(runner.isPaused).toBe(true); await runner.stop(); expect(runner.isPaused).toBe(false); diff --git a/tests/server/routes/websocket-lifecycle-handlers.test.ts b/tests/server/routes/websocket-lifecycle-handlers.test.ts index 51ed33256..d5844e4a1 100644 --- a/tests/server/routes/websocket-lifecycle-handlers.test.ts +++ b/tests/server/routes/websocket-lifecycle-handlers.test.ts @@ -164,6 +164,41 @@ describe("WebSocket lifecycle through the production route", () => { }); }); + it("waits for confirmed Docker pause before publishing paused", async () => { + let complete!: (value: boolean) => void; + Object.assign(runner, { controlResult: new Promise((resolve) => { complete = resolve; }) }); + client.send(JSON.stringify({ type: "pause_simulation" })); + await waitFor(() => runner.pause.mock.calls.length === 1, "runner pause"); + expect(messages.filter((m) => m.type === "simulation_status").map((m) => m.status)).not.toContain("paused"); + complete(true); + await waitFor(() => messages.some((m) => m.type === "simulation_status" && m.status === "paused"), "confirmed pause"); + }); + + it("releases a failed control through the existing stopped status", async () => { + Object.assign(runner, { controlResult: Promise.resolve(false) }); + client.send(JSON.stringify({ type: "pause_simulation" })); + await waitFor(() => pool.releaseRunner.mock.calls.length === 1, "failed control release"); + await waitFor(() => messages.some((m) => m.type === "simulation_status" && m.status === "stopped"), "stopped status"); + expect(messages.some((m) => m.type === "operation_error")).toBe(false); + expect(messages.some((m) => m.type === "serial_output" && m.data.startsWith("[ERR]"))).toBe(true); + }); + + it("ignores delayed control acknowledgement after stop and reuse of the same runner", async () => { + let complete!: (value: boolean) => void; + Object.assign(runner, { controlResult: new Promise((resolve) => { complete = resolve; }) }); + client.send(JSON.stringify({ type: "pause_simulation" })); + await waitFor(() => runner.pause.mock.calls.length === 1, "runner pause"); + client.send(JSON.stringify({ type: "stop_simulation" })); + await waitFor(() => pool.releaseRunner.mock.calls.length === 1, "stop release"); + client.send(JSON.stringify({ type: "start_simulation", code: "void setup() {} void loop() {}" })); + await waitFor(() => runner.runSketch.mock.calls.length === 2, "successor run"); + complete(true); + // Round trip through the same socket drains all earlier acknowledgements. + client.send(JSON.stringify({ type: "serial_input", data: "barrier" })); + await waitFor(() => runner.sendSerialInput.mock.calls.length === 1, "successor input"); + expect(messages.filter((m) => m.type === "simulation_status").map((m) => m.status)).not.toContain("paused"); + }); + it("forwards serial input to the active runner", async () => { client.send(JSON.stringify({ type: "serial_input", data: "hello Uno\n" })); diff --git a/tests/server/services/sandbox-runner-control-ownership.test.ts b/tests/server/services/sandbox-runner-control-ownership.test.ts new file mode 100644 index 000000000..889baed9d --- /dev/null +++ b/tests/server/services/sandbox-runner-control-ownership.test.ts @@ -0,0 +1,191 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { SandboxRunner } from "../../../server/services/sandbox-runner"; +import { cleanupDockerContainer } from "../../../server/services/sandbox/execution-phases/cleanup-phase"; +import { wsMessageSchema } from "../../../shared/schema"; + +const failure = { code: 1, stdout: "", stderr: "synthetic Docker failure", error: null }; +const success = { code: 0, stdout: "", stderr: "", error: null }; +async function flush() { for (let i = 0; i < 8; i++) await Promise.resolve(); } + +function fixture(stateName = "running") { + const pc = { + spawn: vi.fn(), onStdout: vi.fn(), onStderr: vi.fn(), onStderrLine: vi.fn(), + supportsStderrLineStreaming: vi.fn(() => true), onClose: vi.fn(), onError: vi.fn(), + writeStdin: vi.fn(() => true), kill: vi.fn(), destroySockets: vi.fn(), + hasProcess: vi.fn(() => true), clearListeners: vi.fn(), getPid: vi.fn(() => null), + }; + const runner = new SandboxRunner({ processController: pc }); + const internals = runner as any; + const s = internals.executionState; + Object.assign(s, { + state: stateName, currentContainerName: "unosim-synthetic-run-A", runGeneration: 1, + runAbort: new AbortController(), processKilled: false, pauseStartTime: Date.now() - 100, + }); + const execute = vi.spyOn(internals.processExecutor, "execute").mockResolvedValue(failure); + const warn = vi.spyOn(internals.logger, "warn"); + return { runner, pc, s, execute, warn, timeout: internals.timeoutManager }; +} + +describe("Docker lifecycle ownership", () => { + beforeEach(() => vi.useFakeTimers()); + afterEach(() => { vi.restoreAllMocks(); vi.clearAllTimers(); vi.useRealTimers(); }); + + it("must not remain falsely paused with a suspended deadline after nonzero docker pause", async () => { + const f = fixture(); + f.timeout.schedule(30_000, vi.fn()); + f.runner.pause(); + await flush(); + expect.soft(f.runner.isPaused).toBe(false); + expect.soft(f.timeout.isTimeoutPaused()).toBe(false); + expect.soft(f.warn).toHaveBeenCalled(); + }); + + it("must not write a resume clock marker after nonzero docker unpause", async () => { + const f = fixture("paused"); + f.runner.resume(); + await flush(); + expect(f.pc.writeStdin).not.toHaveBeenCalledWith(expect.stringMatching(/^\[\[RESUME_TIME:/)); + }); + + it("must not forward run A's delayed resume result to reused run B", async () => { + const f = fixture("paused"); + let complete!: (value: typeof success) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })) + .mockResolvedValue(success); + f.runner.resume(); + await f.runner.stop(); + await f.runner.resetForReuse(); + Object.assign(f.s, { state: "running", currentContainerName: "unosim-synthetic-run-B", + runGeneration: 2, runAbort: new AbortController(), processKilled: false }); + f.pc.writeStdin.mockClear(); + complete(success); + await flush(); + expect(f.pc.writeStdin).not.toHaveBeenCalled(); + }); + + it("must not report a successful removal when docker rm resolves with code 1", async () => { + const logger = { info: vi.fn(), warn: vi.fn(), debug: vi.fn() }; + await cleanupDockerContainer("unosim-synthetic-cleanup", { + processExecutor: { execute: vi.fn().mockResolvedValue(failure) }, logger: logger as any, + }); + expect.soft(logger.info).not.toHaveBeenCalled(); + expect.soft(logger.warn).toHaveBeenCalled(); + }); + + it("must retain failed container cleanup ownership for retry or quarantine", async () => { + const f = fixture(); + await f.runner.stop(); + expect(f.s.currentContainerName).toBe("unosim-synthetic-run-A"); + }); + + it.each(["pause", "resume"])("commits %s only after confirmed Docker success", async (operation) => { + const f = fixture(operation === "pause" ? "running" : "paused"); + let complete!: (value: typeof success) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })); + expect(f.runner[operation]()).toBe(true); + expect(f.runner.isPaused).toBe(operation === "resume"); + complete(success); + await flush(); + expect(f.runner.isPaused).toBe(operation === "pause"); + }); + + it.each([ + ["pause", "spawn failed"], ["resume", "spawn failed"], + ["pause", "Process timeout after 5000ms"], ["resume", "Process timeout after 5000ms"], + ] as const)("stops after %s failure: %s", async (operation, message) => { + const f = fixture(operation === "pause" ? "running" : "paused"); + f.execute.mockRejectedValueOnce(new Error(message)) + .mockResolvedValue(success); + f.runner[operation](); + await flush(); + expect(f.runner.simulationState).toBe("stopped"); + expect(f.pc.kill).toHaveBeenCalledWith("SIGKILL"); + expect(f.warn).toHaveBeenCalled(); + }); + + it("ignores a stale failed control instead of stopping a successor", async () => { + const f = fixture(); + let complete!: (value: typeof failure) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })) + .mockResolvedValue(success); + f.runner.pause(); + await f.runner.stop(); + await f.runner.resetForReuse(); + Object.assign(f.s, { state: "running", currentContainerName: "unosim-synthetic-run-B", + runGeneration: 2, runAbort: new AbortController(), processKilled: false }); + f.pc.kill.mockClear(); + complete(failure); + await flush(); + expect(f.runner.simulationState).toBe("running"); + expect(f.pc.kill).not.toHaveBeenCalled(); + }); + + it("retries failed removal on repeated stop", async () => { + const f = fixture(); + await f.runner.stop(); + f.execute.mockResolvedValue(success); + await f.runner.stop(); + expect(f.execute.mock.calls.filter(([, args]) => args[0] === "rm")).toHaveLength(2); + expect(f.s.currentContainerName).toBeUndefined(); + }); + + it("refuses reuse while removal remains unconfirmed", async () => { + const f = fixture(); + await f.runner.stop(); + await expect(f.runner.resetForReuse()).rejects.toThrow(/cleanup/i); + expect(f.s.currentContainerName).toBe("unosim-synthetic-run-A"); + }); + + it("does not clear a successor container after delayed removal", async () => { + const f = fixture(); + let complete!: (value: typeof success) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })); + const stopped = f.runner.stop(); + Object.assign(f.s, { currentContainerName: "unosim-synthetic-run-B", runGeneration: 2 }); + complete(success); + await stopped; + expect(f.s.currentContainerName).toBe("unosim-synthetic-run-B"); + }); + + it("does not resurrect a naturally closed run from delayed unpause", async () => { + const f = fixture("paused"); + let complete!: (value: typeof success) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })); + f.runner.resume(); + f.s.state = "stopped"; + f.pc.hasProcess.mockReturnValue(false); + complete(success); + await flush(); + expect(f.runner.simulationState).toBe("stopped"); + expect(f.pc.writeStdin).not.toHaveBeenCalled(); + }); + + it("accepts auto-removal only after Docker confirms no matching container exists", async () => { + const logger = { info: vi.fn(), warn: vi.fn(), debug: vi.fn() }; + const removed = await cleanupDockerContainer("unosim-synthetic-auto-removed", { + processExecutor: { execute: vi.fn().mockResolvedValueOnce(failure).mockResolvedValue(success) }, + logger: logger as any, + }); + expect(removed).toBe(true); + }); + + it("does not commit pending pause after the execution deadline", async () => { + const f = fixture(); + let complete!: (value: typeof success) => void; + f.execute.mockImplementationOnce(() => new Promise((resolve) => { complete = resolve; })) + .mockResolvedValue(success); + f.runner.pause(); + (f.runner as any).dockerManager.setupDockerTimeout(1, { + onOutput: vi.fn(), onError: vi.fn(), onPinState: vi.fn(), + }); + await vi.advanceTimersByTimeAsync(1000); + complete(success); + await flush(); + expect(f.runner.isPaused).toBe(false); + }); + + it.each(["pause_simulation", "resume_simulation"])("confirms the current error contract rejects %s", (operation) => { + expect(wsMessageSchema.safeParse({ type: "operation_error", operation, + code: "SIMULATION_CONTROL_FAILED", message: "synthetic error" }).success).toBe(false); + }); +}); diff --git a/tests/server/services/sandbox-runner-pool.test.ts b/tests/server/services/sandbox-runner-pool.test.ts index c77c9a283..01a886524 100644 --- a/tests/server/services/sandbox-runner-pool.test.ts +++ b/tests/server/services/sandbox-runner-pool.test.ts @@ -702,3 +702,36 @@ describe("SandboxRunnerPool – scalability proof (20 concurrent)", () => { await pool.shutdown(); }); }); + + +describe("SandboxRunnerPool cleanup quarantine", () => { + it("retains unresolved cleanup and recovers the same slot after confirmed retry", async () => { + vi.useFakeTimers(); + const pool = new SandboxRunnerPool({ minRunners: 1, maxRunners: 1, idleTimeoutMs: 5000 }); + await pool.initialize(); + const runner = await pool.acquireRunner(); + runner.hasPendingContainerCleanup = true; + runner.resetForReuse.mockRejectedValue(new Error("Docker cleanup unconfirmed")); + await pool.releaseRunner(runner); + expect(pool.getRunnerIndex(runner)).toBe(0); + expect(pool.getStats().availableRunners).toBe(0); + expect(pool.getStats().resettingRunners).toBe(1); + runner.resetForReuse.mockImplementation(async () => { runner.hasPendingContainerCleanup = false; }); + await vi.advanceTimersByTimeAsync(5000); + expect(await pool.acquireRunner()).toBe(runner); + await pool.shutdown(); + }); + + it("retries stopped cleanup owners during shutdown and cancels recovery timers", async () => { + vi.useFakeTimers(); + const pool = new SandboxRunnerPool({ minRunners: 1, maxRunners: 1, idleTimeoutMs: 5000 }); + await pool.initialize(); + const runner = await pool.acquireRunner(); + runner.hasPendingContainerCleanup = true; + runner.resetForReuse.mockRejectedValue(new Error("Docker cleanup unconfirmed")); + await pool.releaseRunner(runner); + await pool.shutdown(); + expect(runner.stop).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + }); +}); diff --git a/tests/server/services/sandbox/cleanup-phase.test.ts b/tests/server/services/sandbox/cleanup-phase.test.ts index 9f56d6ccc..9fab10e06 100644 --- a/tests/server/services/sandbox/cleanup-phase.test.ts +++ b/tests/server/services/sandbox/cleanup-phase.test.ts @@ -24,7 +24,7 @@ const createMockLogger = () => ({ }); const createMockProcessExecutor = () => ({ - execute: vi.fn(), + execute: vi.fn().mockResolvedValue({ code: 0, stdout: "", stderr: "", error: null }), }); const createMockPinStateBatcher = () => ({ diff --git a/tests/server/services/sandbox/timeout-phase.test.ts b/tests/server/services/sandbox/timeout-phase.test.ts index 4ee362caf..f4a99f255 100644 --- a/tests/server/services/sandbox/timeout-phase.test.ts +++ b/tests/server/services/sandbox/timeout-phase.test.ts @@ -43,7 +43,7 @@ const createBaseState = (): ExecutionState => ({ const createDependencies = () => ({ processExecutor: { - execute: vi.fn(), + execute: vi.fn().mockResolvedValue({ code: 0, stdout: "", stderr: "", error: null }), }, logger: { info: vi.fn(),