From 441e1a4adbdd1ab99a19273857024260e7f6f263 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sat, 3 Oct 2026 23:26:34 +0200 Subject: [PATCH] fix(sandbox): keep a stopped run from starting or touching the runner's next run A start waiting for a sandbox slot survived stop(); when the pool handed the runner to the next user, the stale run started first, with its own removed sketch directory, on the next user's session. Runs now carry a generation and an abort signal, the slot wait is cancellable, and the process controller drops events of replaced children. Co-Authored-By: Claude Opus 5.5 --- docs/UNOSIM_REFACTORING_OPL.md | 8 +- server/services/process-controller.ts | 16 +- server/services/sandbox-runner.ts | 2 + .../sandbox/docker-compile-semaphore.ts | 23 ++- server/services/sandbox/execution-manager.ts | 118 +++++++++--- .../process-controller-stale-child.test.ts | 25 +++ .../docker-compile-semaphore-abort.test.ts | 33 ++++ .../sandbox/runner-reuse-race.test.ts | 168 ++++++++++++++++++ 8 files changed, 358 insertions(+), 35 deletions(-) create mode 100644 tests/server/services/process-controller-stale-child.test.ts create mode 100644 tests/server/services/sandbox/docker-compile-semaphore-abort.test.ts create mode 100644 tests/server/services/sandbox/runner-reuse-race.test.ts diff --git a/docs/UNOSIM_REFACTORING_OPL.md b/docs/UNOSIM_REFACTORING_OPL.md index fe514d7b8..a6bce2581 100644 --- a/docs/UNOSIM_REFACTORING_OPL.md +++ b/docs/UNOSIM_REFACTORING_OPL.md @@ -18,7 +18,7 @@ Abweichungen werden unter „Reihenfolge-Änderungen“ begründet. | R2 | `lastCompiledCode`-Fallback und Sketch-CRUD pro Identität absichern (nicht entfernen) | Isolation | S2, S3 [Code] | mittel–hoch | Nutzerisolation ohne Bruch des REST-/WS-Vertrags | klein | – | fix/isolate-legacy-global-state | DONE | RED→GREEN: `simulation-last-compiled-code.test.ts` (fremder Code wird nie ausgeführt, eigener Fallback bleibt), `last-compiled-code-store.test.ts`, `sketches.routes.test.ts` (Seed read-only, fremde Sketches 404); zwei Bestandstests an Identität angepasst (Begründung im PR) | PR-Merge siehe Verlauf | | R6 | Einheitlicher Compile-Hash inkl. Header (Worker-Identität); kein Binary in REST-Payload/LRU | Korrektheit/Performance/Security | A4, A5, S1-INCBIN [Code] | mittel | keine veralteten Cache-Treffer; Worker und direkter Pfad teilen Cache-Einträge; Antwort 59.653 → 562 Byte (Blink-Sketch, Worker-Pfad) | klein | – | fix/compile-cache-key-and-payload | DONE | RED→GREEN: `arduino-compiler-cache-key.test.ts` (Header-Änderung kompiliert neu; gleiche Identität wie der Worker), `compiler-binary-payload.test.ts` (frisch, gecacht, LRU ohne `binary`) | PR-Merge siehe Verlauf | | R7 | Globaler API-Limiter im Gateway-Modus nach authentifiziertem `subject` statt IP | Skalierung | P1 [Code] | mittel (topologieabhängig) | keine kursweiten 429 hinter Campus-NAT | klein | – | fix/api-rate-limit-identity | DONE | RED→GREEN: `api-rate-limit-key.test.ts` (zwei Subjects hinter einer IP mit getrenntem Budget; ungültiges Gateway-Secret und Local-Modus bleiben pro IP) | PR-Merge siehe Verlauf | -| R3a | Lauf-Generation + Abbruch im Runner-Lifecycle | Isolation/Lifecycle | S4 [plausibel] | hoch | keine fremde Ausgabe, keine verwaisten Container | mittel | – | fix/runner-run-generation | OPEN | deterministischer Race-Test A wartet → A stoppt → B übernimmt | – | +| R3a | Lauf-Generation + Abbruch im Runner-Lifecycle; `ProcessController` leitet nur Events des aktuellen Kindprozesses weiter | Isolation/Lifecycle | S4, S4-CHILD [Code, deterministisch reproduziert] | hoch | keine fremde Ausgabe, kein Start mit fremdem/aufgeräumtem Verzeichnis, kein Eingriff in den Container des Nachfolgers | mittel | – | fix/runner-run-generation | DONE | RED→GREEN: `runner-reuse-race.test.ts` (echter Pool/Runner/ExecutionManager/Semaphore; vorher startete A mit eigenem, bereits gelöschtem Verzeichnis für B), `docker-compile-semaphore-abort.test.ts`, `process-controller-stale-child.test.ts`; Unit 2714, Docker-Integration 26/26 | PR-Merge siehe Verlauf | | R3b | Reset-Ownership in `runner.resetForReuse()` | Kapselung | A6 [Code] | mittel | Reset an einer Stelle | mittel | R3a | refactor/runner-reset-ownership | OPEN | Pool-/Isolationstests | – | | R4a | Orphan-Sweep für Sandbox-Container | Lifecycle | A7 [Code] | mittel | Ressourcen nach Crash frei | klein–mittel | R3a | fix/sandbox-orphan-sweep | OPEN | Sweep-Test (Fake-Executor), Docker-Gate | – | | R4b | WS-Heartbeat, Serialisierung pro Verbindung, Nachrichtenlimit | Lifecycle | A7 [Code] | mittel | halb offene Verbindungen und Floods begrenzt | mittel | – | fix/ws-connection-lifecycle | OPEN | Lifecycle-Tests mit Fake-Timern | – | @@ -39,7 +39,8 @@ Abweichungen werden unter „Reihenfolge-Änderungen“ begründet. | S1-ASM | GAS `.include` im Inline-Assembler gibt die ersten ~10 Zeichen je Zeile einer beliebigen Datei als Fehlermeldung aus; per String-Konkatenation nicht robust textuell filterbar | Security | [Code] verifiziert (Marker-Präfix in der Assemblermeldung) | mittel (Teilinhalt; Default-Deployment hält Secrets nur in Env) | – | groß | – | – | BLOCKED_DECISION | – | Robuster Fix = REST-Compile ohne Zugriff auf Backend-Dateien (Sandbox/Namespace); Architekturentscheidung | | S2 | Globaler `lastCompiledCode` | Isolation | [Code] bestätigt | mittel–hoch | – | – | – | R2 | DONE | WS-Test zweier Subjects | Fallback jetzt pro Subject (LRU, 1000 Subjects) | | S3 | Sketch-CRUD ohne Besitzer | Isolation | [Code] bestätigt | mittel | – | – | – | R2 | DONE | Route-Test zweier Identitäten | Schreiben nur auf eigene Sketches, Seed schreibgeschützt | -| S4 | Runner-Reuse-Race | Isolation/Lifecycle | [plausibel] | hoch | – | – | – | R3a | OPEN | Race-Test | – | +| S4 | Runner-Reuse-Race | Isolation/Lifecycle | [Code] deterministisch reproduziert | hoch | – | – | – | R3a | DONE | Race-Test | Restfenster: Stop genau während `spawn` (ms); verwaiste Container fängt R4a | +| S4-CHILD | `ProcessController` leitet stdout/stderr/close/error eines ersetzten Kindprozesses an die Listener des nächsten Laufs weiter | Isolation/Lifecycle | [Code] reproduziert (neu bei R3a-Verifikation) | mittel | – | – | – | R3a | DONE | `process-controller-stale-child.test.ts` | – | | A1 | Überlappende Concurrency-Mechanismen, vermischte Statusmetriken | Concurrency | [Code] | mittel | – | – | – | R5b | OPEN | – | – | | A2 | Gatekeeper-TTL ohne Queue-Fortsetzung, tote Cache-Locks | Concurrency | [Code] | mittel | – | – | – | R5a | OPEN | – | – | | A3 | Fallback umgeht Lastgrenze, keine Worker-Recovery, unbegrenzte Queue | Concurrency | [Code] | mittel | – | – | – | R5b | OPEN | – | – | @@ -81,4 +82,5 @@ Abweichungen werden unter „Reihenfolge-Änderungen“ begründet. | #160 | R1: Include-Grenze für den REST-Compiler | `1c2995d9` | PR-CI 5/5 grün; Post-Merge-CI von #159 grün | | #161 | R2: Per-Identity-Isolation von Code-Fallback und Sketch-CRUD | `18c7d89b` | PR-CI 5/5 grün; Post-Merge-CI von #160 grün | | #162 | R6: Einheitlicher Compile-Hash, kein Binary in der REST-Antwort | `9716910d` | PR-CI 5/5 grün; Post-Merge-CI von #161 grün | -| R7 | Globaler API-Limiter nach Gateway-Subject | – | – | +| #163 | R7: Globaler API-Limiter nach Gateway-Subject | `7a39d554` | PR-CI 5/5 grün; Post-Merge-CI von #162 grün | +| R3a | Lauf-Generation, abbrechbares Start-Slot-Warten, Kindprozess-Guard | – | – | diff --git a/server/services/process-controller.ts b/server/services/process-controller.ts index 61144f294..b7c377a81 100644 --- a/server/services/process-controller.ts +++ b/server/services/process-controller.ts @@ -96,8 +96,12 @@ export class ProcessController implements IProcessController { /* ignore */ } - // attach existing listeners (guard for nullability) + // attach existing listeners (guard for nullability). Every forwarder checks + // that its child is still the current one: a later spawn belongs to another + // run, and the previous child's late output or close must not reach it. + const child = this.proc; this.proc?.stdout?.on("data", (d: Buffer) => { + if (this.proc !== child) return; this.stdoutListeners.forEach((cb) => cb(d)); }); @@ -110,8 +114,10 @@ export class ProcessController implements IProcessController { private _setupStderrHandling(createInterface: (options: any) => import("node:readline").Interface): void { if (!this.proc?.stderr) return; + const child = this.proc; this.proc.stderr.on("data", (d: Buffer) => { + if (this.proc !== child) return; if (process.env.NODE_ENV === "test") { // convert low-level wrapper events into buffered debug logs try { @@ -135,6 +141,7 @@ export class ProcessController implements IProcessController { crlfDelay: Infinity, }); this.stderrReadline.on("line", (line: string) => { + if (this.proc !== child) return; this.stderrLineListeners.forEach((cb) => cb(line)); }); } @@ -142,11 +149,16 @@ export class ProcessController implements IProcessController { private _setupProcessEventListeners(): void { if (!this.proc) return; + const child = this.proc; this.proc.on("close", (code: number | null) => { + if (this.proc !== child) return; this.closeListeners.forEach((cb) => cb(code)); }); - this.proc.on("error", (err: Error) => this.errorListeners.forEach((cb) => cb(err))); + this.proc.on("error", (err: Error) => { + if (this.proc !== child) return; + this.errorListeners.forEach((cb) => cb(err)); + }); } onStdout(cb: StdDataCb) { diff --git a/server/services/sandbox-runner.ts b/server/services/sandbox-runner.ts index 1bc437d22..524d1aab6 100644 --- a/server/services/sandbox-runner.ts +++ b/server/services/sandbox-runner.ts @@ -398,6 +398,8 @@ export class SandboxRunner { async stop(): Promise { 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.state = SimulationState.STOPPED; s.processKilled = true; diff --git a/server/services/sandbox/docker-compile-semaphore.ts b/server/services/sandbox/docker-compile-semaphore.ts index 9be7f69cf..b44b50069 100644 --- a/server/services/sandbox/docker-compile-semaphore.ts +++ b/server/services/sandbox/docker-compile-semaphore.ts @@ -26,25 +26,40 @@ export class SandboxStartSemaphore { * * @param onQueued Optional callback invoked exactly once when this caller is * placed in the queue (i.e. no slot is immediately available). + * @param signal Optional cancellation: an aborted waiter leaves the queue + * and never takes a slot. * @returns A release function. Must be called exactly once. */ - acquire(onQueued?: () => void, timeoutMs = 60_000): Promise<() => void> { + acquire(onQueued?: () => void, timeoutMs = 60_000, signal?: AbortSignal): Promise<() => void> { return new Promise<() => void>((resolve, reject) => { + if (signal?.aborted) { + reject(new Error("Sandbox start slot acquire cancelled")); + return; + } let settled = false; let attempt: () => void; - const timer = setTimeout(() => { + const leaveQueue = (error: Error) => { if (settled) return; const index = this.queue.findIndex((entry) => entry.attempt === attempt); if (index !== -1) this.queue.splice(index, 1); settled = true; - reject(new Error(`Sandbox start slot timeout after ${timeoutMs}ms`)); - }, timeoutMs); + clearTimeout(timer); + signal?.removeEventListener("abort", onAbort); + reject(error); + }; + const onAbort = () => leaveQueue(new Error("Sandbox start slot acquire cancelled")); + const timer = setTimeout( + () => leaveQueue(new Error(`Sandbox start slot timeout after ${timeoutMs}ms`)), + timeoutMs, + ); + signal?.addEventListener("abort", onAbort, { once: true }); attempt = () => { if (settled) return; if (this._active < this.max) { settled = true; clearTimeout(timer); + signal?.removeEventListener("abort", onAbort); this._active++; resolve(this._makeRelease()); } else { diff --git a/server/services/sandbox/execution-manager.ts b/server/services/sandbox/execution-manager.ts index be48aacca..bb1f83590 100644 --- a/server/services/sandbox/execution-manager.ts +++ b/server/services/sandbox/execution-manager.ts @@ -133,6 +133,20 @@ export interface ExecutionState { dockerAvailable?: boolean; dockerImageBuilt?: boolean; outputCollector?: OutputCollector; + /** Increments with every run on this state; a run whose number is no longer current is stale. */ + runGeneration?: number; + /** Aborted when the current run is stopped or superseded. */ + runAbort?: AbortController | null; +} + +/** + * Identity of one runSketch call. The execution state is reused across runs + * (pooled runners), so every step after an await must check that its run is + * still the current one before it touches shared state. + */ +interface RunToken { + readonly signal: AbortSignal; + isStale(): boolean; } export class ExecutionManager { @@ -170,7 +184,7 @@ export class ExecutionManager { * Main execution entry point: orchestrates prepare → start → stream → timeout → cleanup */ async runSketch(options: RunSketchOptions, state: ExecutionState): Promise { - const { code, onOutput, onError, onExit, onCompileError, onPinState, timeoutSec, onIORegistry, onTelemetry, onPinStateBatch } = options; + const { code, onOutput, onError, onPinState, timeoutSec, onIORegistry, onTelemetry, onPinStateBatch } = options; // Transition to STARTING state const canStart = this.transitionTo(state, SimulationState.STARTING); @@ -179,6 +193,8 @@ export class ExecutionManager { return false; } + const run = this.beginRun(state); + // Clear pending cleanup for a fresh run state.pendingCleanup = false; @@ -252,10 +268,16 @@ export class ExecutionManager { state.serialOutputBatcher.start(); this.registryManager.setSerialOutputBatcher(state.serialOutputBatcher); + let files: { sketchDir: string; sketchFile: string; exeFile: string } | undefined; try { // Prepare environment - const files = await this.prepareEnvironment(code, state, options.headers, options.entryFile); - state.processKilled = false; + files = await this.prepareEnvironment(code, options.headers, options.entryFile); + if (run.isStale()) { + this.cleanupAbandonedSketchDir(files.sketchDir); + return false; + } + state.currentSketchDir = files.sketchDir; + state.processKilled = false; if (state.pendingCleanup || state.processKilled || state.state === SimulationState.STOPPED) { this.filesystemHelper.markTempDirForCleanup(this.extractFilesystemState(state)); @@ -269,27 +291,63 @@ export class ExecutionManager { }); // Setup and run simulation - await this.setupSimulationProcess(files, wrapped, options, state); + await this.setupSimulationProcess(files, wrapped, options, state, run); + if (run.isStale()) return false; return ( state.processController.hasProcess() && (state.state === SimulationState.RUNNING || state.state === SimulationState.PAUSED) && !state.processKilled ); } catch (err) { - const errorMessage = err instanceof Error ? err.message : String(err); - this.logger.error(`Kompilierfehler oder Timeout: ${errorMessage}`); - if (onCompileError) { - onCompileError(errorMessage); - } - if (onExit) { - onExit(-1); - } - state.processController.destroySockets(); - this.filesystemHelper.markTempDirForCleanup(this.extractFilesystemState(state)); + this.handleRunFailure(err, files, options, state, run); return false; } } + private handleRunFailure( + err: unknown, + files: { sketchDir: string } | undefined, + options: RunSketchOptions, + state: ExecutionState, + run: RunToken, + ): void { + if (run.isStale()) { + // A stopped or superseded run reports nothing and leaves the shared state alone. + if (files) this.cleanupAbandonedSketchDir(files.sketchDir); + return; + } + const errorMessage = err instanceof Error ? err.message : String(err); + this.logger.error(`Kompilierfehler oder Timeout: ${errorMessage}`); + options.onCompileError?.(errorMessage); + options.onExit?.(-1); + state.processController.destroySockets(); + this.filesystemHelper.markTempDirForCleanup(this.extractFilesystemState(state)); + } + + /** Starts a new run on the state and supersedes any run still in flight. */ + private beginRun(state: ExecutionState): RunToken { + state.runAbort?.abort(); + const abort = new AbortController(); + const generation = (state.runGeneration ?? 0) + 1; + state.runAbort = abort; + state.runGeneration = generation; + return { + signal: abort.signal, + isStale: () => abort.signal.aborted || state.runGeneration !== generation, + }; + } + + /** Removes the sketch directory of a stale run without touching the shared state. */ + private cleanupAbandonedSketchDir(sketchDir: string): void { + this.filesystemHelper.markTempDirForCleanup({ + currentSketchDir: sketchDir, + isCompiling: false, + pendingCleanup: false, + cleanupRetries: new Map(), + currentRegistryFile: null, + }); + } + /** * Initialize run state for a new execution */ @@ -332,14 +390,11 @@ export class ExecutionManager { */ private async prepareEnvironment( code: string, - state: ExecutionState, headers: Array<{ name: string; content: string }> = [], entryFile?: string, ): Promise<{ sketchDir: string; sketchFile: string; exeFile: string }> { const sketchId = randomUUID(); - const files = await this.fileBuilder.build(code, sketchId, headers, entryFile); - state.currentSketchDir = files.sketchDir; - return files; + return this.fileBuilder.build(code, sketchId, headers, entryFile); } /** @@ -368,9 +423,10 @@ export class ExecutionManager { callbacks: ExecutionCallbacks, opts: RunSketchOptions, state: ExecutionState, + run: RunToken, ): Promise { if (config.serverMode === "local") { - await this.runLocal(files, callbacks, opts, state); + await this.runLocal(files, callbacks, opts, state, run); return; } @@ -378,7 +434,7 @@ export class ExecutionManager { throw new Error("Docker sandbox is unavailable"); } - await this.runDocker(files, callbacks, opts, state); + await this.runDocker(files, callbacks, opts, state, run); } /** @@ -389,10 +445,10 @@ export class ExecutionManager { callbacks: ExecutionCallbacks, opts: RunSketchOptions, state: ExecutionState, + run: RunToken, ): Promise { const executionTimeout = normalizeSimulationTimeout(opts.timeoutSec); const containerName = `unosim-sandbox-${randomUUID()}`; - state.currentContainerName = containerName; const { onCompileError, onCompileSuccess, onExit } = opts; @@ -404,14 +460,15 @@ export class ExecutionManager { const queueStartTime = Date.now(); const releaseSemaphore = await getSandboxStartSemaphore().acquire(() => { opts.onCompileQueued?.(); - }, config.capacity.sandboxStartSlotTimeoutMs); + }, config.capacity.sandboxStartSlotTimeoutMs, run.signal); const queueWaitTimeMs = Date.now() - queueStartTime; - // Guard: abort if the simulation was stopped while we were waiting - if (state.processKilled || state.pendingCleanup || state.state === SimulationState.STOPPED) { + // Guard: abort if the simulation was stopped or the runner reused while we were waiting + if (run.isStale() || state.processKilled || state.pendingCleanup || state.state === SimulationState.STOPPED) { releaseSemaphore(); return; } + state.currentContainerName = containerName; // Release wrapper – idempotent, called from startup callbacks or onClose let semaphoreReleased = false; @@ -451,7 +508,7 @@ export class ExecutionManager { resolveRuntimeStart = resolve; }); const onRuntimeStart = () => { - if (runtimeStarted || state.processKilled || state.pendingCleanup || state.state === SimulationState.STOPPED) return; + if (runtimeStarted || run.isStale() || state.processKilled || state.pendingCleanup || state.state === SimulationState.STOPPED) return; runtimeStarted = true; this.registryManager.enableWaitMode(config.timeouts.registryWaitModeAfterStartMs); this.transitionTo(state, SimulationState.RUNNING); @@ -506,6 +563,8 @@ export class ExecutionManager { state.processController.onClose((_code) => { if (!runtimeStarted) resolveRuntimeStart(false); releaseOnce(); + // stop() already removed this run's container; the state may belong to the next run. + if (run.isStale()) return; this.transitionTo(state, SimulationState.STOPPED); if (state.flushTimer) { clearTimeout(state.flushTimer); @@ -536,9 +595,10 @@ export class ExecutionManager { ); await runtimeStartPromise; } catch (err) { + releaseOnce(); + if (run.isStale()) throw err; const isTimeout = err instanceof Error && err.message.includes("timeout"); compileMetricsTracker.recordCompileComplete(compileStartTime, queueWaitTimeMs, false, isTimeout); - releaseOnce(); this.logger.error(`Docker process spawn failed: ${err instanceof Error ? err.message : String(err)}`); this.transitionTo(state, SimulationState.STOPPED); state.processController.destroySockets(); @@ -555,6 +615,7 @@ export class ExecutionManager { callbacks: ExecutionCallbacks, opts: RunSketchOptions, state: ExecutionState, + run: RunToken, ): Promise { const executionTimeout = normalizeSimulationTimeout(opts.timeoutSec); const { onCompileError, onExit } = opts; @@ -568,6 +629,10 @@ export class ExecutionManager { state.isCompiling = true; await this.performCompilation(files.sketchFile, files.exeFile, opts, state); state.isCompiling = false; + if (run.isStale()) { + this.cleanupAbandonedSketchDir(files.sketchDir); + return; + } // Compile success compileMetricsTracker.recordCompileComplete(compileStartTime, queueWaitTimeMs, true, compileTimedOut); @@ -590,6 +655,7 @@ export class ExecutionManager { // Stream-Phase: Event-Handler registrieren this.setupLocalHandlers(callbacks, onExit, executionTimeout, state); } catch (err) { + if (run.isStale()) throw err; state.isCompiling = false; const isTimeout = err instanceof Error && err.message.includes("timeout"); compileMetricsTracker.recordCompileComplete(compileStartTime, queueWaitTimeMs, false, isTimeout); diff --git a/tests/server/services/process-controller-stale-child.test.ts b/tests/server/services/process-controller-stale-child.test.ts new file mode 100644 index 000000000..95bbf4379 --- /dev/null +++ b/tests/server/services/process-controller-stale-child.test.ts @@ -0,0 +1,25 @@ +import { describe, expect, it } from "vitest"; +import { ProcessController } from "../../../server/services/process-controller"; + +describe("ProcessController with a replaced child", () => { + it("does not forward output or close of a previous child to the next run's listeners", async () => { + const controller = new ProcessController(); + await controller.spawn("node", ["-e", String.raw`setTimeout(() => { process.stdout.write('OLD-CHILD\n'); process.exit(0); }, 150);`]); + + // The next run replaces the child and its listeners before the old child speaks. + controller.clearListeners(); + let output = ""; + const closes: Array = []; + controller.onStdout((data) => { output += data.toString(); }); + controller.onClose((code) => closes.push(code)); + await controller.spawn("node", ["-e", String.raw`setTimeout(() => { process.stdout.write('NEW-CHILD\n'); process.exit(3); }, 400);`]); + + await new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("new child did not close")), 3_000); + controller.onClose(() => { clearTimeout(timer); resolve(); }); + }); + + expect(output).toBe("NEW-CHILD\n"); + expect(closes).toEqual([3]); + }); +}); diff --git a/tests/server/services/sandbox/docker-compile-semaphore-abort.test.ts b/tests/server/services/sandbox/docker-compile-semaphore-abort.test.ts new file mode 100644 index 000000000..e93773263 --- /dev/null +++ b/tests/server/services/sandbox/docker-compile-semaphore-abort.test.ts @@ -0,0 +1,33 @@ +import { describe, expect, it } from "vitest"; +import { SandboxStartSemaphore } from "../../../../server/services/sandbox/docker-compile-semaphore"; + +describe("SandboxStartSemaphore cancellation", () => { + it("removes an aborted waiter so the slot goes to the next one", async () => { + const semaphore = new SandboxStartSemaphore(1); + const release = await semaphore.acquire(); + const abort = new AbortController(); + + const cancelled = semaphore.acquire(undefined, 60_000, abort.signal); + const next = semaphore.acquire(undefined, 60_000); + expect(semaphore.queueLength).toBe(2); + + abort.abort(); + await expect(cancelled).rejects.toThrow(/cancelled/); + expect(semaphore.queueLength).toBe(1); + + release(); + const releaseNext = await next; + expect(semaphore.activeCount).toBe(1); + releaseNext(); + expect(semaphore.activeCount).toBe(0); + }); + + it("rejects immediately for an already aborted signal without taking a slot", async () => { + const semaphore = new SandboxStartSemaphore(1); + const abort = new AbortController(); + abort.abort(); + + await expect(semaphore.acquire(undefined, 60_000, abort.signal)).rejects.toThrow(/cancelled/); + expect(semaphore.activeCount).toBe(0); + }); +}); diff --git a/tests/server/services/sandbox/runner-reuse-race.test.ts b/tests/server/services/sandbox/runner-reuse-race.test.ts new file mode 100644 index 000000000..f840a4432 --- /dev/null +++ b/tests/server/services/sandbox/runner-reuse-race.test.ts @@ -0,0 +1,168 @@ +/** + * Runner reuse while a start is still waiting for a sandbox-start slot. + * + * A waits for the only slot -> A stops -> the pool hands the same runner to B. + * A must neither start a sandbox nor reach B's output or container lifecycle. + * Real pool, runner, execution manager and semaphore; only the child process + * and the Docker CLI are faked. + */ +import { readFile } from "node:fs/promises"; +import { join } from "node:path"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; + +const fakes = vi.hoisted(() => { + type Listener = (value: T) => void; + class FakeChild { + readonly stdout: Listener[] = []; + readonly stderrLine: Listener[] = []; + readonly close: Listener[] = []; + constructor(readonly command: string, readonly args: string[]) {} + } + const spawned: FakeChild[] = []; + const executed: string[][] = []; + + class FakeProcessController { + private child: FakeChild | null = null; + private stdoutListeners: Listener[] = []; + private stderrLineListeners: Listener[] = []; + private closeListeners: Listener[] = []; + async spawn(command: string, args: string[] = []) { + const child = new FakeChild(command, args); + // Mirrors ProcessController: the child forwards to whatever listeners the controller holds. + child.stdout.push((data) => this.stdoutListeners.forEach((cb) => cb(data))); + child.stderrLine.push((line) => this.stderrLineListeners.forEach((cb) => cb(line))); + child.close.push((code) => this.closeListeners.forEach((cb) => cb(code))); + this.child = child; + spawned.push(child); + return null; + } + onStdout(cb: Listener) { this.stdoutListeners.push(cb); } + onStderr() {} + onStderrLine(cb: Listener) { this.stderrLineListeners.push(cb); } + supportsStderrLineStreaming() { return true; } + onClose(cb: Listener) { this.closeListeners.push(cb); } + onError() {} + writeStdin() { return true; } + kill() {} + destroySockets() {} + hasProcess() { return this.child !== null; } + clearListeners() { this.stdoutListeners = []; this.stderrLineListeners = []; this.closeListeners = []; } + getPid() { return null; } + } + + class FakeProcessExecutor { + isBusy = false; + async execute(command: string, args: string[]) { + executed.push([command, ...args]); + return { code: 0, stdout: "Docker version 27.0.0", stderr: "", error: null }; + } + kill() {} + } + + return { spawned, executed, FakeProcessController, FakeProcessExecutor }; +}); + +vi.mock("../../../../server/services/process-controller", () => ({ ProcessController: fakes.FakeProcessController })); +vi.mock("../../../../server/services/process-executor", () => ({ ProcessExecutor: fakes.FakeProcessExecutor })); + +async function waitFor(predicate: () => boolean, label: string): Promise { + const deadline = Date.now() + 3_000; + while (!predicate()) { + if (Date.now() >= deadline) throw new Error(`Timed out waiting for ${label}`); + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +function mountedSketchDir(args: string[]): string { + const mount = args[args.indexOf("-v") + 1]; + return mount.slice(0, mount.indexOf(":/sandbox")); +} + +const ENV = { + UNOSIM_SERVER_MODE: "docker", + UNOSIM_DOCKER_TEST_BYPASS_GATEWAY: "1", + SANDBOX_START_MAX_CONCURRENT: "1", + DOCKER_HOST: "tcp://unosim-test-docker:2375", +}; +const previousEnv: Record = {}; + +describe("runner reuse while a start waits for a sandbox slot", () => { + beforeAll(() => { + for (const [key, value] of Object.entries(ENV)) { + previousEnv[key] = process.env[key]; + process.env[key] = value; + } + vi.resetModules(); + }); + + afterAll(() => { + for (const [key, value] of Object.entries(previousEnv)) { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; + } + vi.resetModules(); + }); + + it("A waits -> A stops -> B takes over: A neither starts nor touches B", async () => { + const { SandboxRunnerPool } = await import("../../../../server/services/sandbox-runner-pool"); + const { getSandboxStartSemaphore } = await import("../../../../server/services/sandbox/docker-compile-semaphore"); + const semaphore = getSandboxStartSemaphore(1); + const releaseHeldSlot = await semaphore.acquire(); + + const pool = new SandboxRunnerPool({ minRunners: 1, maxRunners: 1, acquireTimeoutMs: 1_000, resetTimeoutMs: 1_000 }); + await pool.initialize(); + + const outputA: string[] = []; + const exitA = vi.fn(); + const runnerA = await pool.acquireRunner(); + const runA = runnerA.runSketch({ + code: "void setup(){} void loop(){} // RUN_A", + onOutput: (line) => outputA.push(line), + onError: vi.fn(), + onExit: exitA, + timeoutSec: 30, + }); + await waitFor(() => semaphore.queueLength === 1, "A queued for a start slot"); + + await runnerA.stop(); + await pool.releaseRunner(runnerA); + + const outputB: string[] = []; + const runnerB = await pool.acquireRunner(); + expect(runnerB).toBe(runnerA); + const runB = runnerB.runSketch({ + code: "void setup(){} void loop(){} // RUN_B", + onOutput: (line) => outputB.push(line), + onError: vi.fn(), + onExit: vi.fn(), + timeoutSec: 30, + }); + await waitFor(() => semaphore.queueLength >= 1, "B queued for a start slot"); + + releaseHeldSlot(); + await waitFor(() => fakes.spawned.length >= 1, "a sandbox spawn"); + await new Promise((resolve) => setTimeout(resolve, 50)); + + expect(fakes.spawned).toHaveLength(1); + const [sandbox] = fakes.spawned; + // A stale start mounts A's directory, which A's stop already cleaned up. + const sketch = await readFile(join(mountedSketchDir(sandbox.args), "sketch.cpp"), "utf8").catch(() => ""); + expect(sketch).toContain("RUN_B"); + expect(sketch).not.toContain("RUN_A"); + + const serialEvent = `[[SERIAL_EVENT:0:${Buffer.from("hello from B\n").toString("base64")}]]`; + sandbox.stdout.forEach((cb) => cb(Buffer.from(`[[RUNTIME_START]]\n${serialEvent}\n`))); + // The sandbox reports its (empty) I/O registry, which ends the registry wait. + for (const line of ["[[IO_REGISTRY_START]]", "[[IO_REGISTRY_END]]"]) sandbox.stderrLine.forEach((cb) => cb(line)); + await expect(runB).resolves.toBe(true); + await waitFor(() => outputB.some((line) => line.includes("hello from B")), "B output"); + expect(outputA.join("")).not.toContain("hello from B"); + await expect(runA).resolves.toBe(false); + + const removed = fakes.executed.filter(([command, verb]) => command === "docker" && verb === "rm"); + expect(removed).toEqual([]); + + await runnerB.stop(); + await pool.shutdown(); + }); +});