diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 7fd229c9d..1d70022c3 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,12 +1,20 @@ # Changelog ## [Unreleased] + ### Added - Added a `telemetry` option to `createAgentSession` for passing OpenTelemetry configuration through to the underlying Agent +### Changed + +- Changed Python session pooling to remove the previous 4-session retention cap and 5-minute idle-session eviction, so kernels now stay alive for a session until explicitly disposed via `disposeKernelSessionsByOwner` or `disposeAllKernelSessions` +- Changed kernel cleanup behavior to avoid automatic eviction by idle timeout and capacity pressure, so additional Python sessions are not queued behind retained-session shutdown retries + ### Fixed +- Fixed Python execution cancellation and timeouts by escalating to kernel shutdown if `SIGINT` did not terminate a running cell within 2 seconds, preventing indefinite hangs in queued or stuck sessions +- Fixed cleanup blocking during long-running executions by forcing a kernel shutdown path when interrupt-based cancellation is ignored - Fixed bash output emitting a spurious `[… 0 lines elided (NB) …]` marker (and reordering the artifact link before the command output) when the shell minimizer rewrote a small command's output. After `OutputSink.replace()` swapped the minimized text into the buffer, the subsequent `sink.push("[raw output: artifact://N]\n")` chunk was funneled back into the (now empty) head-retention window while the pre-replace `#totalBytes` still tracked the original raw stream — so `dump()` composed ` + + ` instead of ` + `. `replace()` now realigns `#totalBytes`/`#totalLines`/`#sawData`/`#truncated` to the authoritative buffer and disables head retention for the lifetime of the sink, so further pushes append to the tail buffer in order. The bash executor also drops the leading `\n` on the artifact-link push when the minimized text already ends with one so the separator stays single-newline. ### Fixed diff --git a/packages/coding-agent/src/eval/py/executor.ts b/packages/coding-agent/src/eval/py/executor.ts index 2011349b1..1030ae82b 100644 --- a/packages/coding-agent/src/eval/py/executor.ts +++ b/packages/coding-agent/src/eval/py/executor.ts @@ -13,11 +13,6 @@ import { } from "./kernel"; import { ensurePyToolBridge, registerPyToolBridge } from "./tool-bridge"; -const IDLE_TIMEOUT_MS = 5 * 60 * 1000; // 5 minutes -const MAX_KERNEL_SESSIONS = 4; -const CLEANUP_INTERVAL_MS = 30 * 1000; // 30 seconds -const OWNER_CLEANUP_KERNEL_SHUTDOWN_TIMEOUT_MS = 2_000; - export type PythonKernelMode = "session" | "per-call"; export interface PythonExecutorOptions { @@ -94,42 +89,27 @@ export interface PythonResult { stdinRequested: boolean; } -interface KernelSession { - id: string; +// --------------------------------------------------------------------------- +// Session bookkeeping +// +// One PythonKernel subprocess per session id. Sessions are reused until they +// die or are explicitly disposed. Multiple agent owners can register against +// the same session id; the kernel stays alive until the last owner detaches. +// --------------------------------------------------------------------------- + +interface PythonSession { + sessionId: string; kernel: PythonKernel; - queue: Promise; - restartCount: number; - dead: boolean; - needsRestart: boolean; - disposing: boolean; - disposeCapacityPromise?: Promise; - resolveDisposeCapacity?: () => void; - disposeAttemptPromise?: Promise; - resolveDisposeAttempt?: () => void; - disposeResultPromise?: Promise; - disposeResultTimeoutMs?: number; - nextDisposalRetryAt?: number; - lastUsedAt: number; ownerIds: Set; hasFallbackOwner: boolean; - heartbeatTimer?: NodeJS.Timeout; + queue: Promise; } -const kernelSessions = new Map(); -const disposingKernelSessions = new Set(); -let cleanupTimer: NodeJS.Timeout | null = null; +const sessions = new Map(); -interface KernelSessionExecutionOptions { - sessionFile?: string; - artifactsDir?: string; - signal?: AbortSignal; - deadlineMs?: number; - kernelOwnerId?: string; - /** Bridge session identifier exported into the kernel env as PI_TOOL_BRIDGE_SESSION. */ - bridgeSessionId?: string; - /** Cached bridge connection info. When present, env vars for tool.() get injected. */ - bridge?: { url: string; token: string }; -} +// --------------------------------------------------------------------------- +// Cancellation plumbing +// --------------------------------------------------------------------------- class PythonExecutionCancelledError extends Error { readonly timedOut: boolean; @@ -147,29 +127,6 @@ function getExecutionDeadlineMs(options?: Pick | undefined { - const env: Record = {}; - if (options.sessionFile) env.PI_SESSION_FILE = options.sessionFile; - if (options.artifactsDir) env.PI_ARTIFACTS_DIR = options.artifactsDir; - if (options.bridge && options.bridgeSessionId) { - env.PI_TOOL_BRIDGE_URL = options.bridge.url; - env.PI_TOOL_BRIDGE_TOKEN = options.bridge.token; - env.PI_TOOL_BRIDGE_SESSION = options.bridgeSessionId; - } - return Object.keys(env).length > 0 ? env : undefined; -} - function getRemainingTimeoutMs(deadlineMs?: number): number | undefined { if (deadlineMs === undefined) return undefined; return deadlineMs - Date.now(); @@ -203,62 +160,48 @@ function isTimedOutCancellation(error: unknown, signal?: AbortSignal): boolean { async function waitForPromiseWithCancellation( promise: Promise, - options: Pick, + options: Pick, ): Promise { if (options.signal?.aborted) { throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); } - const remainingMs = getRemainingTimeoutMs(options.deadlineMs); if (remainingMs !== undefined && remainingMs <= 0) { throw new PythonExecutionCancelledError(true); } - if (!options.signal && remainingMs === undefined) { return await promise; } - return await new Promise((resolve, reject) => { - const cleanups: Array<() => void> = []; - const finish = (callback: () => void) => { - while (cleanups.length > 0) { - cleanups.pop()?.(); - } - callback(); - }; - - const onAbort = () => { + const { promise: resultPromise, resolve, reject } = Promise.withResolvers(); + const cleanups: Array<() => void> = []; + const finish = (cb: () => void): void => { + while (cleanups.length > 0) cleanups.pop()?.(); + cb(); + }; + if (options.signal) { + const onAbort = (): void => finish(() => reject(new PythonExecutionCancelledError(isTimedOutCancellation(options.signal?.reason, options.signal))), ); - }; - - if (options.signal) { - options.signal.addEventListener("abort", onAbort, { once: true }); - cleanups.push(() => options.signal?.removeEventListener("abort", onAbort)); - } - - if (remainingMs !== undefined) { - const timeout = setTimeout(() => { - finish(() => reject(new PythonExecutionCancelledError(true))); - }, remainingMs); - timeout.unref(); - cleanups.push(() => clearTimeout(timeout)); - } - - promise.then( - value => finish(() => resolve(value)), - error => finish(() => reject(error)), - ); - }); + options.signal.addEventListener("abort", onAbort, { once: true }); + cleanups.push(() => options.signal?.removeEventListener("abort", onAbort)); + } + if (remainingMs !== undefined) { + const timer = setTimeout(() => finish(() => reject(new PythonExecutionCancelledError(true))), remainingMs); + timer.unref(); + cleanups.push(() => clearTimeout(timer)); + } + promise.then( + value => finish(() => resolve(value)), + err => finish(() => reject(err)), + ); + return await resultPromise; } -async function waitForQueueTurn( - queue: Promise, - options: Pick, -): Promise { - await waitForPromiseWithCancellation(queue, options); -} +// --------------------------------------------------------------------------- +// Result formatting +// --------------------------------------------------------------------------- function formatTimeoutAnnotation(timeoutMs?: number): string | undefined { if (timeoutMs === undefined) return "Command timed out"; @@ -284,534 +227,139 @@ function createCancelledPythonResult(timedOut: boolean, timeoutMs?: number): Pyt }; } -function buildKernelStartOptions( - cwd: string, - env: Record | undefined, - options: KernelSessionExecutionOptions, -) { - return { +// --------------------------------------------------------------------------- +// Kernel start helpers +// --------------------------------------------------------------------------- + +function buildKernelEnv(options: { + sessionFile?: string; + artifactsDir?: string; + bridgeSessionId?: string; + bridge?: { url: string; token: string }; +}): Record | undefined { + const env: Record = {}; + if (options.sessionFile) env.PI_SESSION_FILE = options.sessionFile; + if (options.artifactsDir) env.PI_ARTIFACTS_DIR = options.artifactsDir; + if (options.bridge && options.bridgeSessionId) { + env.PI_TOOL_BRIDGE_URL = options.bridge.url; + env.PI_TOOL_BRIDGE_TOKEN = options.bridge.token; + env.PI_TOOL_BRIDGE_SESSION = options.bridgeSessionId; + } + return Object.keys(env).length > 0 ? env : undefined; +} + +async function startKernel(cwd: string, options: PythonExecutorOptions): Promise { + requireRemainingTimeoutMs(options.deadlineMs); + return await PythonKernel.start({ cwd, - env, + env: buildKernelEnv(options), signal: options.signal, deadlineMs: options.deadlineMs, - }; + }); } -function startCleanupTimer(): void { - if (cleanupTimer) return; - cleanupTimer = setInterval(() => { - void cleanupIdleSessions(); - }, CLEANUP_INTERVAL_MS); - cleanupTimer.unref(); -} - -function stopCleanupTimer(): void { - if (cleanupTimer) { - clearInterval(cleanupTimer); - cleanupTimer = null; - } -} - -function attachKernelOwner(sessionId: string, ownerId?: string): boolean { - const session = kernelSessions.get(sessionId); - if (!session || session.disposing) return false; +function attachOwner(session: PythonSession, sessionId: string, ownerId: string | undefined): void { if (ownerId !== undefined) { if (session.hasFallbackOwner) { session.ownerIds.delete(sessionId); session.hasFallbackOwner = false; } session.ownerIds.add(ownerId); - } else if (session.hasFallbackOwner || session.ownerIds.size === 0) { + return; + } + if (session.hasFallbackOwner || session.ownerIds.size === 0) { session.ownerIds.add(sessionId); session.hasFallbackOwner = true; } - session.lastUsedAt = Date.now(); - return true; } -function getRetainedKernelSessionCount(): number { - return kernelSessions.size + disposingKernelSessions.size; -} - -function syncCleanupTimer(): void { - if (kernelSessions.size === 0 && disposingKernelSessions.size === 0) { - stopCleanupTimer(); - return; +async function acquireSession(sessionId: string, cwd: string, options: PythonExecutorOptions): Promise { + const existing = sessions.get(sessionId); + if (existing) { + attachOwner(existing, sessionId, options.kernelOwnerId); + return existing; } - startCleanupTimer(); -} - -function retryPendingKernelSessionDisposals(now: number = Date.now()): void { - for (const session of disposingKernelSessions.values()) { - if (session.disposeResultPromise) continue; - if (session.nextDisposalRetryAt !== undefined && session.nextDisposalRetryAt > now) continue; - session.nextDisposalRetryAt = undefined; - void disposeKernelSession(session); - } -} - -function beginDisposingKernelSession(session: KernelSession): boolean { - if (session.disposing) return false; - session.disposing = true; - disposingKernelSessions.add(session); - if (kernelSessions.get(session.id) === session) { - kernelSessions.delete(session.id); - } - if (session.heartbeatTimer) { - clearInterval(session.heartbeatTimer); - session.heartbeatTimer = undefined; - } - syncCleanupTimer(); - return true; -} - -function finishDisposingKernelSession(session: KernelSession): void { - disposingKernelSessions.delete(session); - session.resolveDisposeCapacity?.(); - session.resolveDisposeCapacity = undefined; - session.disposeCapacityPromise = undefined; - session.resolveDisposeAttempt = undefined; - session.disposeAttemptPromise = undefined; - session.disposeResultPromise = undefined; - session.disposeResultTimeoutMs = undefined; - session.nextDisposalRetryAt = undefined; - syncCleanupTimer(); -} - -async function waitForDisposalCapacity( - options: Pick, -): Promise { - retryPendingKernelSessionDisposals(); - - const disposalPromises: Promise[] = []; - let nextRetryAt: number | undefined; - for (const session of disposingKernelSessions.values()) { - if (session.disposeCapacityPromise) { - disposalPromises.push( - session.disposeCapacityPromise.then( - () => undefined, - () => undefined, - ), - ); - } - if (session.disposeAttemptPromise) { - disposalPromises.push( - session.disposeAttemptPromise.then( - () => undefined, - () => undefined, - ), - ); - } - if (session.nextDisposalRetryAt !== undefined) { - nextRetryAt = - nextRetryAt === undefined - ? session.nextDisposalRetryAt - : Math.min(nextRetryAt, session.nextDisposalRetryAt); - } - } - if (disposalPromises.length > 0) { - await waitForPromiseWithCancellation(Promise.race(disposalPromises), options); - return; - } - if (nextRetryAt === undefined) return; - await waitForPromiseWithCancellation( - Bun.sleep(Math.max(0, nextRetryAt - Date.now())).then(() => undefined), - options, - ); -} - -async function ensureKernelSessionCapacity( - options: Pick, -): Promise { - while (getRetainedKernelSessionCount() >= MAX_KERNEL_SESSIONS) { - if (disposingKernelSessions.size > 0) { - await waitForDisposalCapacity(options); - continue; - } - if (kernelSessions.size === 0) { - await waitForDisposalCapacity(options); - continue; - } - await evictOldestSession(); - } -} - -async function cleanupIdleSessions(): Promise { - const now = Date.now(); - const toDispose: KernelSession[] = []; - - for (const session of kernelSessions.values()) { - if (session.dead || now - session.lastUsedAt > IDLE_TIMEOUT_MS) { - toDispose.push(session); - } - } - - if (toDispose.length > 0) { - logger.debug("Cleaning up idle kernel sessions", { count: toDispose.length }); - await Promise.allSettled(toDispose.map(session => disposeKernelSession(session))); - } - - retryPendingKernelSessionDisposals(now); - syncCleanupTimer(); -} - -async function evictOldestSession(): Promise { - let oldest: KernelSession | null = null; - for (const session of kernelSessions.values()) { - if (!oldest || session.lastUsedAt < oldest.lastUsedAt) { - oldest = session; - } - } - if (oldest) { - logger.debug("Evicting oldest kernel session", { id: oldest.id }); - await disposeKernelSession(oldest); - } -} - -export async function disposeAllKernelSessions(): Promise { - stopCleanupTimer(); - const sessions = Array.from(new Set([...kernelSessions.values(), ...disposingKernelSessions.values()])); - await Promise.allSettled(sessions.map(session => disposeKernelSession(session))); -} - -export async function disposeKernelSessionsByOwner(ownerId: string): Promise { - const sessionsToDispose: KernelSession[] = []; - for (const session of new Set([...kernelSessions.values(), ...disposingKernelSessions.values()])) { - if (!session.ownerIds.delete(ownerId)) continue; - if (session.ownerIds.size === 0) { - sessionsToDispose.push(session); - } - } - await Promise.allSettled( - sessionsToDispose.map(session => disposeKernelSession(session, OWNER_CLEANUP_KERNEL_SHUTDOWN_TIMEOUT_MS)), - ); - syncCleanupTimer(); -} - -async function ensureKernelAvailable( - cwd: string, - options: Pick = {}, -): Promise { - const availability = await waitForPromiseWithCancellation(checkPythonKernelAvailability(cwd), options); - if (!availability.ok) { - throw new Error(availability.reason ?? "Python kernel unavailable"); - } -} - -function ensureKernelHeartbeat(session: KernelSession): void { - if (session.heartbeatTimer) return; - session.heartbeatTimer = setInterval(() => { - if (session.dead || session.needsRestart) return; - if (!session.kernel.isAlive()) { - session.dead = true; - } - }, 5000); - session.heartbeatTimer.unref(); -} - -async function createKernelSession( - sessionId: string, - cwd: string, - options: KernelSessionExecutionOptions = {}, -): Promise { - requireRemainingTimeoutMs(options.deadlineMs); - const env = buildKernelEnv(options); - const startOptions = buildKernelStartOptions(cwd, env, options); - - const kernel = await logger.time("createKernelSession:PythonKernel.start", PythonKernel.start, startOptions); - - const hasFallbackOwner = options.kernelOwnerId === undefined; - const initialOwnerId = options.kernelOwnerId ?? sessionId; - const session: KernelSession = { - id: sessionId, + const kernel = await startKernel(cwd, options); + const session: PythonSession = { + sessionId, kernel, + ownerIds: new Set(), + hasFallbackOwner: false, queue: Promise.resolve(), - restartCount: 0, - dead: false, - needsRestart: false, - disposing: false, - disposeResultPromise: undefined, - nextDisposalRetryAt: undefined, - lastUsedAt: Date.now(), - ownerIds: new Set([initialOwnerId]), - hasFallbackOwner, }; - - ensureKernelHeartbeat(session); - + attachOwner(session, sessionId, options.kernelOwnerId); + sessions.set(sessionId, session); return session; } -async function restartKernelSession( - session: KernelSession, +async function replaceSessionKernel( + session: PythonSession, cwd: string, - options: KernelSessionExecutionOptions = {}, + options: PythonExecutorOptions, ): Promise { - session.restartCount += 1; - if (session.restartCount > 1) { - throw new Error("Python kernel restarted too many times in this session"); - } + const old = session.kernel; + const remaining = getRemainingTimeoutMs(options.deadlineMs); + await old + .shutdown(remaining !== undefined ? { timeoutMs: Math.max(0, remaining) } : undefined) + .catch(() => undefined); requireRemainingTimeoutMs(options.deadlineMs); - try { - const deadKernel = session.dead || !session.kernel.isAlive(); - const shutdownTimeoutMs = requireRemainingTimeoutMs(options.deadlineMs); - const shutdownResult = await session.kernel.shutdown({ signal: options.signal, timeoutMs: shutdownTimeoutMs }); - if (!shutdownResult.confirmed && !deadKernel) { - throw new Error("Failed to confirm crashed kernel shutdown before restart"); - } - if (!shutdownResult.confirmed) { - logger.warn("Proceeding with retained kernel restart after unconfirmed dead-kernel shutdown", { - sessionId: session.id, - }); - } - const env = buildKernelEnv(options); - const startOptions = buildKernelStartOptions(cwd, env, options); - const kernel = await PythonKernel.start(startOptions); - session.kernel = kernel; - session.dead = false; - session.needsRestart = false; - session.lastUsedAt = Date.now(); - ensureKernelHeartbeat(session); - } catch (err) { - session.restartCount = 0; - logger.warn("Failed to restart kernel", { error: err instanceof Error ? err.message : String(err) }); - throw err; - } + session.kernel = await startKernel(cwd, options); } -type KernelDisposalResult = { status: "confirmed" } | { status: "unconfirmed" } | { status: "failed"; err: unknown }; -type KernelDisposalWaitResult = KernelDisposalResult | { status: "timedOut" }; - -function createKernelDisposalResultPromise(session: KernelSession, timeoutMs?: number): Promise { - return Promise.resolve() - .then(() => session.kernel.shutdown(timeoutMs === undefined ? undefined : { timeoutMs })) - .then( - result => (result.confirmed ? { status: "confirmed" as const } : { status: "unconfirmed" as const }), - (err: unknown) => ({ status: "failed" as const, err }), - ); +async function resetSession(sessionId: string): Promise { + const existing = sessions.get(sessionId); + if (!existing) return; + sessions.delete(sessionId); + await existing.kernel.shutdown().catch(() => undefined); } -function getOrStartKernelDisposalResultPromise( - session: KernelSession, - timeoutMs?: number, -): Promise { - if (!session.disposeResultPromise) { - session.disposeResultTimeoutMs = timeoutMs; - const releaseDisposalAttempt = Promise.withResolvers(); - session.disposeAttemptPromise = releaseDisposalAttempt.promise; - session.resolveDisposeAttempt = releaseDisposalAttempt.resolve; - const disposeResultPromise = createKernelDisposalResultPromise(session, timeoutMs); - void disposeResultPromise.then(result => { - if (result.status === "confirmed") { - finishDisposingKernelSession(session); - return; - } - if (session.disposing) { - session.nextDisposalRetryAt = Date.now() + CLEANUP_INTERVAL_MS; - syncCleanupTimer(); - } - }); - const disposalAttemptPromise = disposeResultPromise.finally(() => { - releaseDisposalAttempt.resolve(); - if (session.disposeResultPromise === disposalAttemptPromise) { - session.disposeResultPromise = undefined; - session.disposeResultTimeoutMs = undefined; - } - if (session.disposeAttemptPromise === releaseDisposalAttempt.promise) { - session.disposeAttemptPromise = undefined; - session.resolveDisposeAttempt = undefined; - } - }); - session.disposeResultPromise = disposalAttemptPromise; - } - return session.disposeResultPromise; -} - -async function waitForKernelSessionDisposal( - session: KernelSession, - timeoutMs?: number, -): Promise { - const disposeResultPromise = session.disposeResultPromise; - if (!disposeResultPromise) { - return undefined; - } - if (timeoutMs === undefined) { - return await disposeResultPromise; - } - - let timeoutId: NodeJS.Timeout | undefined; - const result = await Promise.race([ - disposeResultPromise, - new Promise<{ status: "timedOut" }>(resolve => { - timeoutId = setTimeout(() => resolve({ status: "timedOut" }), timeoutMs); - timeoutId.unref(); - }), - ]); - - if (timeoutId) { - clearTimeout(timeoutId); - } - return result; -} - -function retryKernelSessionDisposalInBackground(session: KernelSession): void { - session.nextDisposalRetryAt = undefined; - void disposeKernelSession(session); -} - -async function disposeKernelSession(session: KernelSession, shutdownTimeoutMs?: number): Promise { - if (!session.disposing) { - if (!beginDisposingKernelSession(session)) return; - const releaseDisposalCapacity = Promise.withResolvers(); - session.disposeCapacityPromise = releaseDisposalCapacity.promise; - session.resolveDisposeCapacity = releaseDisposalCapacity.resolve; - } - - if ( - shutdownTimeoutMs === undefined && - session.disposeResultPromise && - session.disposeResultTimeoutMs !== undefined - ) { - const inheritedResult = await session.disposeResultPromise; - if (inheritedResult.status === "confirmed") { - finishDisposingKernelSession(session); - return; - } - session.disposeResultPromise = undefined; - session.disposeResultTimeoutMs = undefined; - logger.warn("Retained kernel shutdown was not confirmed during owner cleanup; retrying without timeout", { - sessionId: session.id, - }); - } - - const inheritedBackgroundRetryTimeoutMs = - shutdownTimeoutMs === undefined && session.disposeResultPromise && session.disposeResultTimeoutMs === undefined - ? OWNER_CLEANUP_KERNEL_SHUTDOWN_TIMEOUT_MS - : shutdownTimeoutMs; - - getOrStartKernelDisposalResultPromise(session, shutdownTimeoutMs); - const result = await waitForKernelSessionDisposal(session, inheritedBackgroundRetryTimeoutMs); - if (!result) { - return; - } - if (result.status === "timedOut") { - logger.warn( - shutdownTimeoutMs === undefined - ? "Timed out waiting for retained kernel shutdown during global cleanup; retained capacity remains reserved" - : "Timed out shutting down retained kernel during owner cleanup", - { - sessionId: session.id, - timeoutMs: inheritedBackgroundRetryTimeoutMs, - }, - ); - if (shutdownTimeoutMs !== undefined) { - retryKernelSessionDisposalInBackground(session); - } - return; - } - if (result.status === "confirmed") { - finishDisposingKernelSession(session); - return; - } - if (result.status === "unconfirmed") { - logger.warn( - shutdownTimeoutMs === undefined - ? "Kernel shutdown was not confirmed; retained capacity remains reserved" - : "Retained kernel shutdown was not confirmed during owner cleanup", - { sessionId: session.id }, - ); - if (shutdownTimeoutMs !== undefined) { - retryKernelSessionDisposalInBackground(session); - } - return; - } - logger.warn( - shutdownTimeoutMs === undefined - ? "Failed to shutdown kernel" - : "Failed to shutdown retained kernel during owner cleanup", - { - sessionId: session.id, - error: result.err instanceof Error ? result.err.message : String(result.err), - }, - ); - if (shutdownTimeoutMs !== undefined) { - retryKernelSessionDisposalInBackground(session); - } - return; -} - -async function withKernelSession( - sessionId: string, - cwd: string, - handler: (kernel: PythonKernel) => Promise, - options: KernelSessionExecutionOptions = {}, +async function runQueued( + session: PythonSession, + options: Pick, + work: () => Promise, ): Promise { - let session = kernelSessions.get(sessionId); - if (session?.disposing) { - session = undefined; - } - if (!session) { - await ensureKernelSessionCapacity(options); - requireRemainingTimeoutMs(options.deadlineMs); - if (options.signal?.aborted) { - throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); - } - session = await logger.time("kernel:createKernelSession", createKernelSession, sessionId, cwd, options); - kernelSessions.set(sessionId, session); - startCleanupTimer(); - } - attachKernelOwner(sessionId, options.kernelOwnerId); - - if (session.disposing) { - return await withKernelSession(sessionId, cwd, handler, options); - } - - const run = async (): Promise => { - session!.lastUsedAt = Date.now(); - if (session!.dead || session!.needsRestart || !session!.kernel.isAlive()) { - await logger.time("kernel:restartKernelSession", restartKernelSession, session!, cwd, options); - } - try { - const result = await logger.time("kernel:withSession:handler", handler, session!.kernel); - session!.restartCount = 0; - return result; - } catch (err) { - if (!session!.dead && !session!.needsRestart && session!.kernel.isAlive()) { - throw err; - } - await logger.time("kernel:restartKernelSession", restartKernelSession, session!, cwd, options); - const result = await logger.time("kernel:postRestart:handler", handler, session!.kernel); - session!.restartCount = 0; - return result; - } - }; - - const queue = session.queue; - let releaseTurn: (() => void) | undefined; - const turn = new Promise(resolve => { - releaseTurn = resolve; - }); - session.queue = queue - .then( - () => turn, - () => turn, - ) - .then( - () => undefined, - () => undefined, - ); - + const previous = session.queue; + const { promise: ourSlot, resolve: releaseSlot } = Promise.withResolvers(); + // Keep the queue chained even if WE bail out: future runs must still wait + // for `previous` to finish before they touch the kernel. + session.queue = previous.catch(() => undefined).then(() => ourSlot); try { - await waitForQueueTurn(queue, options); - if (session.disposing) { - return await withKernelSession(sessionId, cwd, handler, options); - } - return await run(); + await waitForPromiseWithCancellation( + previous.catch(() => undefined), + options, + ); + return await work(); } finally { - releaseTurn?.(); + releaseSlot(); } } +// --------------------------------------------------------------------------- +// Public dispose entry points +// --------------------------------------------------------------------------- + +export async function disposeAllKernelSessions(): Promise { + const all = [...sessions.values()]; + sessions.clear(); + await Promise.allSettled(all.map(session => session.kernel.shutdown().catch(() => undefined))); +} + +export async function disposeKernelSessionsByOwner(ownerId: string): Promise { + const toShutdown: PythonSession[] = []; + for (const session of [...sessions.values()]) { + if (!session.ownerIds.delete(ownerId)) continue; + if (session.ownerIds.size === 0) toShutdown.push(session); + } + for (const session of toShutdown) { + if (sessions.get(session.sessionId) === session) sessions.delete(session.sessionId); + } + await Promise.allSettled(toShutdown.map(session => session.kernel.shutdown().catch(() => undefined))); +} + +// --------------------------------------------------------------------------- +// Execution +// --------------------------------------------------------------------------- + async function executeWithKernel( kernel: PythonKernelExecutor, code: string, @@ -900,6 +448,64 @@ async function executeWithKernel( } } +async function ensureKernelAvailable(cwd: string, options: PythonExecutorOptions): Promise { + const availability = await waitForPromiseWithCancellation(checkPythonKernelAvailability(cwd), options); + if (!availability.ok) { + throw new Error(availability.reason ?? "Python kernel unavailable"); + } +} + +async function ensureToolBridge(options: PythonExecutorOptions): Promise { + if (!options.toolSession || options.bridge) return; + try { + options.bridge = await ensurePyToolBridge(); + } catch (err) { + logger.warn("Failed to start Python tool bridge", { + error: err instanceof Error ? err.message : String(err), + }); + } +} + +async function executePerCall(code: string, cwd: string, options: PythonExecutorOptions): Promise { + if (options.bridge && !options.bridgeSessionId) { + options.bridgeSessionId = `py-bridge:${crypto.randomUUID()}`; + } + const kernel = await startKernel(cwd, options); + try { + return await executeWithKernel(kernel, code, options); + } finally { + await kernel.shutdown().catch(() => undefined); + } +} + +async function executeOnSession(code: string, cwd: string, options: PythonExecutorOptions): Promise { + const sessionId = options.sessionId ?? `session:${cwd}`; + if (options.bridge && !options.bridgeSessionId) { + options.bridgeSessionId = sessionId; + } + if (options.reset) { + await resetSession(sessionId); + } + const session = await acquireSession(sessionId, cwd, options); + return await runQueued(session, options, async () => { + if (options.signal?.aborted) { + throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); + } + if (!session.kernel.isAlive()) { + await replaceSessionKernel(session, cwd, options); + } + try { + return await executeWithKernel(session.kernel, code, options); + } catch (err) { + if (isCancellationError(err) || options.signal?.aborted) throw err; + if (session.kernel.isAlive()) throw err; + // Kernel died during execute. Replace it and retry once on a fresh one. + await replaceSessionKernel(session, cwd, options); + return await executeWithKernel(session.kernel, code, options); + } + }); +} + export async function executePythonWithKernel( kernel: PythonKernelExecutor, code: string, @@ -923,55 +529,14 @@ export async function executePython(code: string, options?: PythonExecutorOption isTimedOutCancellation(executionOptions.signal.reason, executionOptions.signal), ); } - - await ensureKernelAvailable(cwd); + await ensureKernelAvailable(cwd, executionOptions); + await ensureToolBridge(executionOptions); const kernelMode = executionOptions.kernelMode ?? "session"; - - if (executionOptions.toolSession && !executionOptions.bridge) { - try { - executionOptions.bridge = await ensurePyToolBridge(); - } catch (err) { - logger.warn("Failed to start Python tool bridge", { - error: err instanceof Error ? err.message : String(err), - }); - } - } - if (kernelMode === "per-call") { - if (executionOptions.bridge && !executionOptions.bridgeSessionId) { - executionOptions.bridgeSessionId = `py-bridge:${crypto.randomUUID()}`; - } - const env = buildKernelEnv(executionOptions); - requireRemainingTimeoutMs(deadlineMs); - const startOptions = buildKernelStartOptions(cwd, env, executionOptions); - const kernel = await PythonKernel.start(startOptions); - try { - return await executeWithKernel(kernel, code, executionOptions); - } finally { - await kernel.shutdown(); - } + return await executePerCall(code, cwd, executionOptions); } - - const sessionId = executionOptions.sessionId ?? `session:${cwd}`; - if (executionOptions.bridge && !executionOptions.bridgeSessionId) { - executionOptions.bridgeSessionId = sessionId; - } - if (executionOptions.reset) { - const existing = kernelSessions.get(sessionId); - if (existing) { - await disposeKernelSession(existing); - if (existing.disposing && existing.nextDisposalRetryAt !== undefined) { - retryKernelSessionDisposalInBackground(existing); - } - } - } - return await withKernelSession( - sessionId, - cwd, - async kernel => executeWithKernel(kernel, code, executionOptions), - executionOptions, - ); + return await executeOnSession(code, cwd, executionOptions); } catch (err) { if (isCancellationError(err) || executionOptions.signal?.aborted) { return createCancelledPythonResult(isTimedOutCancellation(err, executionOptions.signal)); diff --git a/packages/coding-agent/src/eval/py/kernel.ts b/packages/coding-agent/src/eval/py/kernel.ts index f8758c8ae..dfd43bbf3 100644 --- a/packages/coding-agent/src/eval/py/kernel.ts +++ b/packages/coding-agent/src/eval/py/kernel.ts @@ -43,6 +43,11 @@ async function ensureRunnerScript(): Promise { const SHUTDOWN_GRACE_MS = 1_000; const STARTUP_TIMEOUT_MS = 10_000; +// How long to wait after SIGINT for the runner to emit `done`. If the cell is +// stuck in code that ignores Python signals (e.g. a C extension holding the +// GIL), we escalate to a full subprocess shutdown so the host queue unblocks +// instead of hanging the session forever. +const INTERRUPT_ESCALATION_MS = 2_000; export interface KernelExecuteOptions { signal?: AbortSignal; @@ -158,6 +163,7 @@ interface PendingExecution { timedOut: boolean; stdinRequested: boolean; settled: boolean; + escalationTimer?: NodeJS.Timeout; } export class PythonKernel { @@ -273,22 +279,41 @@ export class PythonKernel { }); }; + const requestCancel = () => { + if (pending.settled || pending.escalationTimer) return; + void this.interrupt(); + const escalation = setTimeout(() => { + if (pending.settled) return; + logger.warn("Python runner did not respond to SIGINT; terminating subprocess", { + kernelId: this.id, + }); + // `shutdown()` aborts pending executions immediately and escalates to + // SIGTERM/SIGKILL, so the host queue unblocks even if the runner is + // stuck in a non-interruptible state. + void this.shutdown(); + }, INTERRUPT_ESCALATION_MS); + escalation.unref?.(); + pending.escalationTimer = escalation; + }; + const onAbort = () => { pending.cancelled = true; pending.timedOut = pending.timedOut || isTimeoutReason(options?.signal?.reason); - void this.interrupt(); + requestCancel(); }; const timeoutId = typeof options?.timeoutMs === "number" && options.timeoutMs > 0 ? setTimeout(() => { pending.timedOut = true; pending.cancelled = true; - void this.interrupt(); + requestCancel(); }, options.timeoutMs) : undefined; const cleanup = () => { if (timeoutId) clearTimeout(timeoutId); + if (pending.escalationTimer) clearTimeout(pending.escalationTimer); + pending.escalationTimer = undefined; options?.signal?.removeEventListener("abort", onAbort); }; diff --git a/packages/coding-agent/test/core/python-executor-owner-cleanup.test.ts b/packages/coding-agent/test/core/python-executor-owner-cleanup.test.ts index 1c0a45efd..545142f5f 100644 --- a/packages/coding-agent/test/core/python-executor-owner-cleanup.test.ts +++ b/packages/coding-agent/test/core/python-executor-owner-cleanup.test.ts @@ -192,101 +192,6 @@ describe("python executor owner cleanup", () => { expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); }); - it("keeps tracked disposals counted against retained kernel capacity until shutdown settles", async () => { - const retainedKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel(), new FakeKernel()]; - const replacementKernel = new FakeKernel(); - const shutdownDeferreds = retainedKernels.map(() => Promise.withResolvers()); - for (const [index, kernel] of retainedKernels.entries()) { - kernel.shutdown = vi.fn(() => shutdownDeferreds[index]!.promise); - } - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start"); - for (const kernel of [...retainedKernels, replacementKernel]) { - startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - } - - for (const [index] of retainedKernels.entries()) { - await executePython(`print(${index})`, { - cwd: `/tmp/capacity-tracking-${index}`, - sessionId: `capacity-session-${index}`, - kernelMode: "session", - }); - } - expect(startSpy).toHaveBeenCalledTimes(4); - - const globalDisposal = disposeAllKernelSessions(); - await Promise.resolve(); - - const fifthExecution = executePython("print('replacement')", { - cwd: "/tmp/capacity-tracking-replacement", - sessionId: "capacity-session-replacement", - kernelMode: "session", - }); - await Promise.resolve(); - - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - - shutdownDeferreds[0]!.resolve({ confirmed: true }); - await fifthExecution; - expect(startSpy).toHaveBeenCalledTimes(5); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - - for (const deferred of shutdownDeferreds.slice(1)) { - deferred.resolve({ confirmed: true }); - } - await globalDisposal; - await disposeAllKernelSessions(); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("waits with the owner-cleanup timeout when the last owner is removed from an already-disposing session", async () => { - vi.useFakeTimers(); - try { - const kernel = new FakeKernel(); - const shutdownConfirmation = Promise.withResolvers(); - kernel.shutdown = vi.fn(() => shutdownConfirmation.promise); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - await executePython("print('owner-a')", { - cwd: "/tmp/disposing-owner-cleanup-session", - sessionId: "disposing-owner-cleanup-session", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - - let globalCleanupResolved = false; - const globalCleanup = disposeAllKernelSessions().finally(() => { - globalCleanupResolved = true; - }); - await flushMicrotasks(); - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - expect(globalCleanupResolved).toBe(false); - - let ownerCleanupResolved = false; - const ownerCleanup = disposeKernelSessionsByOwner("owner-a").finally(() => { - ownerCleanupResolved = true; - }); - await flushMicrotasks(); - expect(ownerCleanupResolved).toBe(false); - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - - vi.advanceTimersByTime(2_000); - await ownerCleanup; - expect(ownerCleanupResolved).toBe(true); - expect(globalCleanupResolved).toBe(false); - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - - shutdownConfirmation.resolve({ confirmed: true }); - await globalCleanup; - expect(globalCleanupResolved).toBe(true); - expect(startSpy).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }); - it("returns a cancelled result when a dead session restart shutdown times out", async () => { const kernel = new FakeKernel(); kernel.alive = false; @@ -296,13 +201,13 @@ describe("python executor owner cleanup", () => { if (shutdownCallCount > 1) { return { confirmed: true }; } - return await new Promise((_, reject) => { - const timer = setTimeout( - () => reject(new DOMException("Python kernel shutdown timed out", "TimeoutError")), - options?.timeoutMs ?? 0, - ); - timer.unref?.(); - }); + const { promise, reject } = Promise.withResolvers(); + const timer = setTimeout( + () => reject(new DOMException("Python kernel shutdown timed out", "TimeoutError")), + options?.timeoutMs ?? 0, + ); + timer.unref?.(); + return await promise; }); vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); @@ -319,353 +224,21 @@ describe("python executor owner cleanup", () => { expect(kernel.shutdown).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: expect.any(Number) })); expect(startSpy).toHaveBeenCalledTimes(1); }); - it("returns owner cleanup promptly but keeps retained capacity reserved until shutdown is confirmed", async () => { - vi.useFakeTimers(); - try { - const retainedKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel(), new FakeKernel()]; - const replacementKernel = new FakeKernel(); - const shutdownDeferreds = retainedKernels.map(() => Promise.withResolvers()); - for (const [index, kernel] of retainedKernels.entries()) { - kernel.shutdown = vi.fn(() => shutdownDeferreds[index]!.promise); - } - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start"); - for (const kernel of [...retainedKernels, replacementKernel]) { - startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - } - - for (const [index] of retainedKernels.entries()) { - await executePython(`print(${index})`, { - cwd: `/tmp/owner-timeout-kernel-${index}`, - sessionId: `timeout-session-${index}`, - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - } - - let ownerCleanupResolved = false; - const ownerCleanup = disposeKernelSessionsByOwner("owner-a").finally(() => { - ownerCleanupResolved = true; - }); - await flushMicrotasks(); - - for (const kernel of retainedKernels) { - expect(kernel.shutdown).toHaveBeenCalledWith({ timeoutMs: 2_000 }); - } - expect(ownerCleanupResolved).toBe(false); - - vi.advanceTimersByTime(2_000); - await ownerCleanup; - expect(ownerCleanupResolved).toBe(true); - - const blockedExecution = executePython("print('replacement')", { - cwd: "/tmp/owner-timeout-kernel-replacement", - sessionId: "timeout-session-replacement", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - await Promise.resolve(); - - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - - shutdownDeferreds[0]!.resolve({ confirmed: true }); - await blockedExecution; - expect(startSpy).toHaveBeenCalledTimes(5); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - - for (const deferred of shutdownDeferreds.slice(1)) { - deferred.resolve({ confirmed: true }); - } - await Promise.resolve(); - await disposeAllKernelSessions(); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }); - - it("keeps owner-cleanup retries on the timer path without evicting unrelated live sessions", async () => { - vi.useFakeTimers(); - try { - const ownerKernel = new FakeKernel(); - const unrelatedKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel()]; - const replacementKernel = new FakeKernel(); - const retryConfirmation = Promise.withResolvers(); - let shutdownCallCount = 0; - ownerKernel.shutdown = vi.fn(async (options?: FakeKernelShutdownOptions): Promise => { - shutdownCallCount += 1; - if (shutdownCallCount === 1) { - expect(options).toEqual({ timeoutMs: 2_000 }); - return { confirmed: false }; - } - if (shutdownCallCount === 2) { - expect(options).toBeUndefined(); - return { confirmed: false }; - } - return await retryConfirmation.promise; - }); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start"); - for (const kernel of [ownerKernel, ...unrelatedKernels, replacementKernel]) { - startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - } - - await executePython("print('owner-a')", { - cwd: "/tmp/timer-retry-owner-a", - sessionId: "timer-retry-owner-a", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - for (const [index, ownerId] of ["owner-b", "owner-c", "owner-d"].entries()) { - await executePython(`print(${index})`, { - cwd: `/tmp/timer-retry-${ownerId}`, - sessionId: `timer-retry-${ownerId}`, - kernelMode: "session", - kernelOwnerId: ownerId, - }); - } - - await disposeKernelSessionsByOwner("owner-a"); - expect(ownerKernel.shutdown).toHaveBeenCalledTimes(2); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - const blockedExecution = executePython("print('replacement')", { - cwd: "/tmp/timer-retry-replacement", - sessionId: "timer-retry-replacement", - kernelMode: "session", - kernelOwnerId: "owner-e", - }); - await Promise.resolve(); - await Promise.resolve(); - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - vi.advanceTimersByTime(29_999); - await Promise.resolve(); - await Promise.resolve(); - expect(ownerKernel.shutdown).toHaveBeenCalledTimes(2); - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - vi.advanceTimersByTime(1); - await Promise.resolve(); - await Promise.resolve(); - expect(ownerKernel.shutdown).toHaveBeenCalledTimes(3); - expect(ownerKernel.shutdown).toHaveBeenNthCalledWith(3, undefined); - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - retryConfirmation.resolve({ confirmed: true }); - await blockedExecution; - expect(startSpy).toHaveBeenCalledTimes(5); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - await disposeAllKernelSessions(); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }); - - it("waits for confirmed disposal capacity before evicting unrelated retained sessions", async () => { - const ownerKernel = new FakeKernel(); - const unrelatedKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel()]; - const replacementKernel = new FakeKernel(); - const retryConfirmation = Promise.withResolvers(); - let shutdownCallCount = 0; - ownerKernel.shutdown = vi.fn(async (_options?: FakeKernelShutdownOptions): Promise => { - shutdownCallCount += 1; - if (shutdownCallCount === 1) { - return { confirmed: false }; - } - return await retryConfirmation.promise; - }); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start"); - for (const kernel of [ownerKernel, ...unrelatedKernels, replacementKernel]) { - startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - } - - await executePython("print('owner-a')", { - cwd: "/tmp/unconfirmed-owner-cleanup-a", - sessionId: "unconfirmed-owner-cleanup-a", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - for (const [index, ownerId] of ["owner-b", "owner-c", "owner-d"].entries()) { - await executePython(`print(${index})`, { - cwd: `/tmp/unconfirmed-owner-cleanup-${ownerId}`, - sessionId: `unconfirmed-owner-cleanup-${ownerId}`, - kernelMode: "session", - kernelOwnerId: ownerId, - }); - } - - await disposeKernelSessionsByOwner("owner-a"); - expect(ownerKernel.shutdown).toHaveBeenCalledTimes(2); - expect(ownerKernel.shutdown).toHaveBeenNthCalledWith(1, { timeoutMs: 2_000 }); - expect(ownerKernel.shutdown).toHaveBeenNthCalledWith(2, undefined); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - const blockedExecution = executePython("print('replacement')", { - cwd: "/tmp/unconfirmed-owner-cleanup-replacement", - sessionId: "unconfirmed-owner-cleanup-replacement", - kernelMode: "session", - kernelOwnerId: "owner-e", - }); - await Promise.resolve(); - await Promise.resolve(); - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - for (const [index, ownerId] of ["owner-b", "owner-c", "owner-d"].entries()) { - await executePython(`print('reuse-${ownerId}')`, { - cwd: `/tmp/unconfirmed-owner-cleanup-${ownerId}`, - sessionId: `unconfirmed-owner-cleanup-${ownerId}`, - kernelMode: "session", - kernelOwnerId: ownerId, - }); - expect(unrelatedKernels[index]!.execute).toHaveBeenCalledTimes(2); - expect(unrelatedKernels[index]!.shutdown).not.toHaveBeenCalled(); - } - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - - retryConfirmation.resolve({ confirmed: true }); - await blockedExecution; - expect(startSpy).toHaveBeenCalledTimes(5); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - for (const kernel of unrelatedKernels) { - expect(kernel.shutdown).not.toHaveBeenCalled(); - } - - await disposeAllKernelSessions(); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("owner cleanup retries every shutdown in background and frees retained capacity one confirmation at a time", async () => { - const retainedKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel(), new FakeKernel()]; - const replacementKernel = new FakeKernel(); - const laterKernel = new FakeKernel(); - const retryConfirmations = retainedKernels.map(() => Promise.withResolvers()); - for (const [index, kernel] of retainedKernels.entries()) { - let shutdownCallCount = 0; - kernel.shutdown = vi.fn(async (_options?: FakeKernelShutdownOptions): Promise => { - shutdownCallCount += 1; - if (shutdownCallCount === 1) { - return { confirmed: false }; - } - return await retryConfirmations[index]!.promise; - }); - } - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start"); - for (const kernel of [...retainedKernels, replacementKernel, laterKernel]) { - startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - } - - for (const [index] of retainedKernels.entries()) { - await executePython(`print(${index})`, { - cwd: `/tmp/unconfirmed-capacity-${index}`, - sessionId: `unconfirmed-capacity-session-${index}`, - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - } - - await disposeKernelSessionsByOwner("owner-a"); - for (const kernel of retainedKernels) { - expect(kernel.shutdown).toHaveBeenCalledTimes(2); - expect(kernel.shutdown).toHaveBeenNthCalledWith(1, { timeoutMs: 2_000 }); - expect(kernel.shutdown).toHaveBeenNthCalledWith(2, undefined); - } - - let globalCleanupResolved = false; - const globalCleanup = disposeAllKernelSessions().then(() => { - globalCleanupResolved = true; - }); - await Promise.resolve(); - await Promise.resolve(); - - for (const kernel of retainedKernels) { - expect(kernel.shutdown).toHaveBeenCalledTimes(2); - } - expect(globalCleanupResolved).toBe(false); - - const blockedExecution = executePython("print('replacement')", { - cwd: "/tmp/unconfirmed-capacity-replacement", - sessionId: "unconfirmed-capacity-replacement", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - await Promise.resolve(); - await Promise.resolve(); - expect(startSpy).toHaveBeenCalledTimes(4); - expect(replacementKernel.execute).not.toHaveBeenCalled(); - - retryConfirmations[0]!.resolve({ confirmed: true }); - await blockedExecution; - expect(globalCleanupResolved).toBe(false); - expect(startSpy).toHaveBeenCalledTimes(5); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - - const secondBlockedExecution = executePython("print('later')", { - cwd: "/tmp/unconfirmed-capacity-later", - sessionId: "unconfirmed-capacity-later", - kernelMode: "session", - kernelOwnerId: "owner-c", - }); - await Promise.resolve(); - await Promise.resolve(); - expect(globalCleanupResolved).toBe(false); - expect(startSpy).toHaveBeenCalledTimes(5); - expect(laterKernel.execute).not.toHaveBeenCalled(); - - retryConfirmations[1]!.resolve({ confirmed: true }); - await secondBlockedExecution; - expect(globalCleanupResolved).toBe(false); - expect(startSpy).toHaveBeenCalledTimes(6); - expect(laterKernel.execute).toHaveBeenCalledTimes(1); - - for (const confirmation of retryConfirmations.slice(2)) { - confirmation.resolve({ confirmed: true }); - } - await globalCleanup; - expect(globalCleanupResolved).toBe(true); - - await disposeAllKernelSessions(); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - expect(laterKernel.shutdown).toHaveBeenCalledTimes(1); - }); it("does not let stuck retained executions block owner or global cleanup", async () => { const ownerKernel = new FakeKernel(); const globalKernel = new FakeKernel(); const ownerExecutionStarted = Promise.withResolvers(); const globalExecutionStarted = Promise.withResolvers(); + const ownerExecutionHang = Promise.withResolvers(); + const globalExecutionHang = Promise.withResolvers(); ownerKernel.execute = vi.fn(async () => { ownerExecutionStarted.resolve(); - return await new Promise(() => {}); + return await ownerExecutionHang.promise; }); globalKernel.execute = vi.fn(async () => { globalExecutionStarted.resolve(); - return await new Promise(() => {}); + return await globalExecutionHang.promise; }); vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); vi.spyOn(PythonKernel, "start") @@ -697,7 +270,11 @@ describe("python executor owner cleanup", () => { await flushMicrotasks(); expect(globalKernel.shutdown).toHaveBeenCalledTimes(1); await globalCleanup; + + ownerExecutionHang.resolve(OK_RESULT); + globalExecutionHang.resolve(OK_RESULT); }); + it("leaves per-call kernels out of owner-scoped retained cleanup and keeps global cleanup intact", async () => { const perCallKernel = new FakeKernel(); const retainedKernel = new FakeKernel();