Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 38 additions & 6 deletions server/routes/simulation.ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
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<void> {
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);
Expand All @@ -96,13 +120,21 @@ function handlePauseSimulation(
/**
* Handle "resume_simulation" WebSocket message
*/
function handleResumeSimulation(
async function handleResumeSimulation(
_ws: WebSocket,
clientState: ClientState,
sessionManager: WsSessionManager,
): void {
): Promise<void> {
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;
Expand Down
65 changes: 47 additions & 18 deletions server/services/sandbox-runner-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ interface PooledRunner {
runner: SandboxRunner;
inUse: boolean;
resetting: boolean;
quarantined?: boolean;
lastReleasedTime: number;
idleTimer: ReturnType<typeof setTimeout> | null;
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -179,15 +181,20 @@ export class SandboxRunnerPool {
return;
}

if (!pooledRunner.inUse) {
if (!pooledRunner.inUse && !pooledRunner.quarantined) {
this.logger.warn(
"[SandboxRunnerPool] Attempt to release already-released runner (ignored)",
);
return;
}

// 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();

Expand All @@ -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);
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -317,6 +345,7 @@ export class SandboxRunnerPool {
}

async shutdown(): Promise<void> {
this.shuttingDown = true;
this.logger.info("[SandboxRunnerPool] Shutting down...");

for (const entry of this.queue) {
Expand All @@ -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) {
Expand Down
Loading
Loading