diff --git a/docs/tools/python.md b/docs/tools/python.md index 177bc39e4dd..318487f95d6 100644 --- a/docs/tools/python.md +++ b/docs/tools/python.md @@ -40,6 +40,10 @@ The REPL runs in the session cwd and retains variables, imports, and loaded data The kernel owner id is `python:`, deliberately distinct from the `eval` owner so the two never alias and cleanup remains scoped to the owning session. +An invocation captures its cwd, session file, session id, and settings before preflight and remains tracked through its transcript append. Clearing detaches the captured generation and joins its pending operations and kernel shutdown. Its session cleanup callback remains registered until cleanup succeeds, so session teardown also joins an older generation still clearing while a successor is active. + +Owner ids remain string labels, not private identity authority. A later cleanup using the same label can still select a successor kernel; this lifecycle improvement does not fix that legacy ABA boundary. Captured invocation metadata grants no transcript or audit filesystem append authority. + The session's kernel is disposed on: - the `clear` action, diff --git a/packages/coding-agent/changelog.d/python-physical-preflight-lifetime.md b/packages/coding-agent/changelog.d/python-physical-preflight-lifetime.md new file mode 100644 index 00000000000..a9c649eee4d --- /dev/null +++ b/packages/coding-agent/changelog.d/python-physical-preflight-lifetime.md @@ -0,0 +1,8 @@ +### Fixed + +- Register Python operations before availability and initialization, and join their captured physical work during owner cleanup. +- Keep standalone Python invocations tracked through transcript append, capture their execution context before preflight, and retain cleanup joins for earlier generations while a successor runs. + +### Security + +- Python owner ids remain legacy string labels; these lifecycle changes do not establish private owner authority or transcript/audit filesystem append permission. diff --git a/packages/coding-agent/src/eval/py/executor.ts b/packages/coding-agent/src/eval/py/executor.ts index d0c2672ecf7..ac62ba67427 100644 --- a/packages/coding-agent/src/eval/py/executor.ts +++ b/packages/coding-agent/src/eval/py/executor.ts @@ -11,6 +11,7 @@ import { checkPythonKernelAvailability, type KernelExecuteOptions, type KernelExecuteResult, + type KernelShutdownResult, PythonKernel, } from "./kernel"; import type { PythonRuntimeOptions } from "./runtime"; @@ -116,12 +117,26 @@ interface PythonSession { ownerIds: Set; hasFallbackOwner: boolean; queue: Promise; + cleanupPromise?: Promise; + cleanupStart?: () => void; } interface InitializingPythonSession { - sessionId: string; promise: Promise; + ownerIds: Set; + retirementOwnerIds: Set; + hasFallbackOwner: boolean; cancelled?: PythonExecutionCancelledError; + kernel?: PythonKernel; + cleanupError?: unknown; + cleanupFailed?: boolean; +} + +interface ActivePythonRequest { + readonly ownerId?: string; + readonly controller: AbortController; + readonly completion: Promise; + readonly finish: () => void; } let pythonResourceCleanupRegistered = false; @@ -134,6 +149,11 @@ function ensurePythonResourceCleanup(): void { const sessions = new Map(); const retiringKernels = new Set(); const settingsScopes = new WeakMap(); +const activeRequests = new Set(); +const retiringSessions = new Set(); +const initializingSessions = new Set(); +const retiringKernelShutdowns = new Map>(); +const retiringKernelOwners = new Map>(); /** * Fingerprint the settings values that actually shape a spawned kernel process. @@ -169,11 +189,101 @@ function scopedSessionId(sessionId: string, activeSettings: Settings | undefined return `${sessionId}:settings-${settingsScope(activeSettings)}`; } -async function shutdownOrRetainKernel(kernel: PythonKernel): Promise { - const result = await kernel.shutdown().catch(() => undefined); - if (result?.confirmed) return; +function retainKernelForOwners(kernel: PythonKernel, ownerIds: Set): void { retiringKernels.add(kernel); + const owners = retiringKernelOwners.get(kernel) ?? new Set(); + for (const ownerId of ownerIds) owners.add(ownerId); + retiringKernelOwners.set(kernel, owners); +} + +async function shutdownOrRetainKernel(kernel: PythonKernel, ownerIds: Set): Promise { + let result: KernelShutdownResult; + try { + result = await callKernelShutdown(kernel); + } catch (error) { + retainKernelForOwners(kernel, ownerIds); + logger.warn("Python kernel shutdown not confirmed", { kernelId: kernel.id, reason: error }); + throw error; + } + if (result.confirmed) return; + retainKernelForOwners(kernel, ownerIds); logger.warn("Python kernel shutdown not confirmed", { kernelId: kernel.id }); + throw unconfirmedShutdownError(kernel); +} + +function prepareSessionShutdown(session: PythonSession): Promise { + if (session.cleanupPromise) return session.cleanupPromise; + const { promise: cleanup, resolve, reject } = Promise.withResolvers(); + retiringSessions.add(session); + session.cleanupPromise = cleanup; + session.cleanupStart = () => void confirmedKernelShutdown(session.kernel).then(resolve, reject); + void cleanup.then( + result => { + if (result.confirmed) retiringSessions.delete(session); + if (session.cleanupPromise === cleanup) { + session.cleanupPromise = undefined; + session.cleanupStart = undefined; + } + }, + () => { + retiringSessions.add(session); + if (session.cleanupPromise === cleanup) { + session.cleanupPromise = undefined; + session.cleanupStart = undefined; + } + }, + ); + return cleanup; +} + +function startSessionShutdown(session: PythonSession): void { + const start = session.cleanupStart; + if (!start) return; + session.cleanupStart = undefined; + start(); +} + +function callKernelShutdown(kernel: PythonKernel): Promise { + try { + return kernel.shutdown(); + } catch (error) { + return Promise.reject(error); + } +} + +function confirmedKernelShutdown(kernel: PythonKernel): Promise { + return callKernelShutdown(kernel).then(result => { + if (!result.confirmed) throw unconfirmedShutdownError(kernel); + return result; + }); +} + +function prepareRetiringKernelShutdown(kernel: PythonKernel): [Promise, () => void] { + const existing = retiringKernelShutdowns.get(kernel); + if (existing) return [existing, () => {}]; + const { promise, resolve, reject } = Promise.withResolvers(); + retiringKernelShutdowns.set(kernel, promise); + let started = false; + void promise.then( + result => { + if (result.confirmed) { + retiringKernels.delete(kernel); + retiringKernelOwners.delete(kernel); + } + if (retiringKernelShutdowns.get(kernel) === promise) retiringKernelShutdowns.delete(kernel); + }, + () => { + if (retiringKernelShutdowns.get(kernel) === promise) retiringKernelShutdowns.delete(kernel); + }, + ); + return [ + promise, + () => { + if (started) return; + started = true; + void confirmedKernelShutdown(kernel).then(resolve, reject); + }, + ]; } function isInitializingSession( @@ -202,6 +312,59 @@ function getExecutionDeadlineMs(options?: Pick(); + let deadlineTimer: NodeJS.Timeout | undefined; + const request: ActivePythonRequest = { + ownerId, + controller, + completion, + finish: () => { + activeRequests.delete(request); + if (deadlineTimer) clearTimeout(deadlineTimer); + resolve(); + }, + }; + activeRequests.add(request); + + const signals = [controller.signal]; + if (capturedOptions.signal) signals.push(capturedOptions.signal); + const signal = AbortSignal.any(signals); + const remainingMs = getRemainingTimeoutMs(deadlineMs); + deadlineTimer = + remainingMs !== undefined + ? setTimeout(() => controller.abort(new PythonExecutionCancelledError(true)), Math.max(0, remainingMs)) + : undefined; + deadlineTimer?.unref?.(); + + return { + options: { ...capturedOptions, signal, deadlineMs }, + request, + cwd, + }; +} + function getRemainingTimeoutMs(deadlineMs?: number): number | undefined { if (deadlineMs === undefined) return undefined; return deadlineMs - Date.now(); @@ -233,6 +396,13 @@ function isTimedOutCancellation(error: unknown, signal?: AbortSignal): boolean { return reason instanceof Error ? reason.name === "TimeoutError" : false; } +function throwIfExecutionCancelled(options: Pick): void { + if (options.signal?.aborted) { + throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); + } + requireRemainingTimeoutMs(options.deadlineMs); +} + async function waitForPromiseWithCancellation( promise: Promise, options: Pick, @@ -333,7 +503,7 @@ function buildKernelEnv(options: { } async function startKernel(cwd: string, options: PythonExecutorOptions): Promise { - requireRemainingTimeoutMs(options.deadlineMs); + throwIfExecutionCancelled(options); return await PythonKernel.start({ cwd, settings: options.settings, @@ -344,7 +514,11 @@ async function startKernel(cwd: string, options: PythonExecutorOptions): Promise }); } -function attachOwner(session: PythonSession, sessionId: string, ownerId: string | undefined): void { +function attachOwner( + session: PythonSession | InitializingPythonSession, + sessionId: string, + ownerId: string | undefined, +): void { if (ownerId !== undefined) { if (session.hasFallbackOwner) { session.ownerIds.delete(sessionId); @@ -362,6 +536,7 @@ function attachOwner(session: PythonSession, sessionId: string, ownerId: string async function acquireSession(sessionId: string, cwd: string, options: PythonExecutorOptions): Promise { const existing = sessions.get(sessionId); if (existing) { + attachOwner(existing, sessionId, options.kernelOwnerId); let session: PythonSession; if (isInitializingSession(existing)) { try { @@ -379,22 +554,37 @@ async function acquireSession(sessionId: string, cwd: string, options: PythonExe } else { session = existing; } - attachOwner(session, sessionId, options.kernelOwnerId); - ensurePythonResourceCleanup(); return session; } const initializing: InitializingPythonSession = { - sessionId, + ownerIds: new Set(), + retirementOwnerIds: new Set(), + hasFallbackOwner: false, promise: Promise.resolve().then(async () => { + throwIfExecutionCancelled(options); const kernel = await startKernel(cwd, options); - if (initializing.cancelled) { - await shutdownOrRetainKernel(kernel); - throw initializing.cancelled; + initializing.kernel = kernel; + if (initializing.cancelled || options.signal?.aborted) { + try { + await shutdownOrRetainKernel(kernel, initializing.ownerIds); + } catch (error) { + initializing.cleanupError = error; + initializing.cleanupFailed = true; + } + throw ( + initializing.cancelled ?? + new PythonExecutionCancelledError(isTimedOutCancellation(options.signal?.reason, options.signal)) + ); } const current = sessions.get(sessionId); if (current !== initializing) { - await shutdownOrRetainKernel(kernel); + try { + await shutdownOrRetainKernel(kernel, initializing.ownerIds); + } catch (error) { + initializing.cleanupError = error; + initializing.cleanupFailed = true; + } const winner = current ? isInitializingSession(current) ? await waitForPromiseWithCancellation(current.promise, options) @@ -408,16 +598,17 @@ async function acquireSession(sessionId: string, cwd: string, options: PythonExe kernel, kernelInstanceId: crypto.randomUUID(), bridgeCapability: options.bridge?.capability, - ownerIds: new Set(), - hasFallbackOwner: false, + ownerIds: new Set(initializing.ownerIds), + hasFallbackOwner: initializing.hasFallbackOwner, queue: Promise.resolve(), }; sessions.set(sessionId, session); return session; }), }; + attachOwner(initializing, sessionId, options.kernelOwnerId); + initializingSessions.add(initializing); sessions.set(sessionId, initializing); - ensurePythonResourceCleanup(); let cancellationTimer: NodeJS.Timeout | undefined; const retireCancelledInitialization = (timedOut: boolean): void => { if (initializing.cancelled) return; @@ -444,14 +635,13 @@ async function acquireSession(sessionId: string, cwd: string, options: PythonExe } void initializing.promise .finally(() => { + initializingSessions.delete(initializing); options.signal?.removeEventListener("abort", onAbort); if (cancellationTimer) clearTimeout(cancellationTimer); }) .catch(() => undefined); try { const session = await waitForPromiseWithCancellation(initializing.promise, options); - attachOwner(session, sessionId, options.kernelOwnerId); - ensurePythonResourceCleanup(); return session; } catch (err) { if (isCancellationError(err)) { @@ -469,22 +659,26 @@ async function replaceSessionKernel( cwd: string, options: PythonExecutorOptions, ): Promise { + throwIfExecutionCancelled(options); const old = session.kernel; const remaining = getRemainingTimeoutMs(options.deadlineMs); await old .shutdown(remaining !== undefined ? { timeoutMs: Math.max(0, remaining) } : undefined) .catch(() => undefined); + throwIfExecutionCancelled(options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } - requireRemainingTimeoutMs(options.deadlineMs); + throwIfExecutionCancelled(options); const bridge = options.bridge; const previousCapability = bridge?.capability; const nextCapability = bridge ? crypto.randomUUID() : undefined; if (bridge && nextCapability) bridge.capability = nextCapability; let next: PythonKernel | undefined; try { + throwIfExecutionCancelled(options); next = await startKernel(cwd, options); + throwIfExecutionCancelled(options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } @@ -503,9 +697,14 @@ async function replaceSessionKernel( async function resetSession(sessionId: string): Promise { const existing = sessions.get(sessionId); if (!existing) return; - sessions.delete(sessionId); - const session = isInitializingSession(existing) ? await existing.promise.catch(() => undefined) : existing; - await session?.kernel.shutdown().catch(() => undefined); + if (sessions.get(sessionId) === existing) sessions.delete(sessionId); + if (isInitializingSession(existing)) { + await existing.promise.catch(() => undefined); + return; + } + const shutdown = prepareSessionShutdown(existing); + startSessionShutdown(existing); + await shutdown.catch(() => undefined); } async function runQueued( @@ -533,68 +732,125 @@ async function runQueued( // Public dispose entry points // --------------------------------------------------------------------------- -export async function disposeAllKernelSessions(): Promise { - const all = [...sessions.entries()]; - const retiring = [...retiringKernels]; - retiringKernels.clear(); - for (const [id, session] of all) { - if (sessions.get(id) === session) sessions.delete(id); - } - const resolved = await Promise.all( - all.map( - async ([id, entry]) => - [id, isInitializingSession(entry) ? await entry.promise.catch(() => undefined) : entry] as const, +export function disposeAllKernelSessions(): Promise { + const entries = [...sessions.entries()]; + const requests = [...activeRequests]; + const initializing = [...initializingSessions]; + const targets = new Set([ + ...[...entries.map(([, entry]) => entry), ...retiringSessions].filter( + (entry): entry is PythonSession => !isInitializingSession(entry), ), - ); - const shutdownTargets = resolved.filter( - (entry): entry is readonly [string, PythonSession] => entry[1] !== undefined, - ); - const results = await Promise.allSettled(shutdownTargets.map(([, session]) => session.kernel.shutdown())); - for (let i = 0; i < shutdownTargets.length; i += 1) { - const [id, session] = shutdownTargets[i]; - const result = results[i]; - if (result.status === "fulfilled" && result.value.confirmed) continue; - const reason = result.status === "rejected" ? result.reason : "not confirmed"; - logger.warn("Python kernel shutdown not confirmed", { sessionId: id, reason }); - if (!sessions.has(id)) sessions.set(id, session); - } - const retiringResults = await Promise.allSettled(retiring.map(kernel => kernel.shutdown())); - for (let i = 0; i < retiring.length; i += 1) { - const result = retiringResults[i]; - if (result.status === "fulfilled" && result.value.confirmed) continue; - retiringKernels.add(retiring[i]); - logger.warn("Python kernel shutdown not confirmed", { - kernelId: retiring[i].id, - reason: result.status === "rejected" ? result.reason : "not confirmed", - }); + ]); + const shutdowns = [...targets].map(session => [session, prepareSessionShutdown(session)] as const); + const retiring = [...retiringKernels].map(kernel => [kernel, ...prepareRetiringKernelShutdown(kernel)] as const); + for (const [id, entry] of entries) { + if (sessions.get(id) === entry) sessions.delete(id); + } + // Publish every exact physical cleanup future before abort listeners or shutdown callbacks can reenter disposal. + for (const request of requests) { + if (!request.controller.signal.aborted) request.controller.abort(new PythonExecutionCancelledError(false)); } + for (const [session] of shutdowns) startSessionShutdown(session); + for (const [, , start] of retiring) start(); + return (async () => { + await Promise.all([ + ...requests.map(request => request.completion), + ...initializing.map(entry => entry.promise.catch(() => undefined)), + ]); + const failures: unknown[] = initializing.filter(entry => entry.cleanupFailed).map(entry => entry.cleanupError); + const results = await Promise.allSettled(shutdowns.map(([, cleanup]) => cleanup)); + for (let index = 0; index < shutdowns.length; index += 1) { + const [session] = shutdowns[index]; + const result = results[index]; + if (result.status === "fulfilled" && result.value.confirmed) continue; + if (!sessions.has(session.sessionId)) sessions.set(session.sessionId, session); + if (result.status === "rejected") failures.push(result.reason); + else failures.push(unconfirmedShutdownError(session.kernel)); + } + const retiringResults = await Promise.allSettled(retiring.map(([, cleanup]) => cleanup)); + for (let index = 0; index < retiring.length; index += 1) { + const [kernel] = retiring[index]; + const result = retiringResults[index]; + if (result.status === "fulfilled" && result.value.confirmed) continue; + retiringKernels.add(kernel); + failures.push(result.status === "rejected" ? result.reason : unconfirmedShutdownError(kernel)); + } + if (failures.length) throw failures[0]; + })(); } -export async function disposeKernelSessionsByOwner(ownerId: string): Promise { +export function disposeKernelSessionsByOwner(ownerId: string): Promise { + const requests = [...activeRequests].filter(request => request.ownerId === ownerId); + const entries = [...sessions.entries()].filter(([, entry]) => entry.ownerIds.has(ownerId)); + const initializing = [...initializingSessions].filter( + entry => entry.ownerIds.has(ownerId) || entry.retirementOwnerIds.has(ownerId), + ); + const sessionsToRetire = new Set([ + ...[...entries.map(([, entry]) => entry), ...retiringSessions].filter( + (entry): entry is PythonSession => !isInitializingSession(entry) && entry.ownerIds.has(ownerId), + ), + ]); const toShutdown: PythonSession[] = []; - for (const session of [...sessions.values()]) { - if (isInitializingSession(session) || !session.ownerIds.has(ownerId)) continue; + for (const entry of initializing) entry.retirementOwnerIds.add(ownerId); + for (const [id, entry] of entries) + if (isInitializingSession(entry) && entry.ownerIds.size === 1) sessions.delete(id); + for (const entry of initializing) entry.ownerIds.delete(ownerId); + for (const session of sessionsToRetire) { if (session.ownerIds.size === 1) { toShutdown.push(session); - continue; + if (sessions.get(session.sessionId) === session) sessions.delete(session.sessionId); + } else { + session.ownerIds.delete(ownerId); } - session.ownerIds.delete(ownerId); } - for (const session of toShutdown) { - if (sessions.get(session.sessionId) === session) sessions.delete(session.sessionId); + const shutdowns = toShutdown.map(session => [session, prepareSessionShutdown(session)] as const); + const retiring = [...retiringKernels] + .filter(kernel => retiringKernelOwners.get(kernel)?.has(ownerId)) + .map(kernel => [kernel, ...prepareRetiringKernelShutdown(kernel)] as const); + // All cleanup futures are visible before an abort listener can reenter this API. + for (const request of requests) { + if (!request.controller.signal.aborted) request.controller.abort(new PythonExecutionCancelledError(false)); } - const results = await Promise.allSettled(toShutdown.map(session => session.kernel.shutdown())); - for (let i = 0; i < toShutdown.length; i += 1) { - const session = toShutdown[i]; - const result = results[i]; - if (result.status === "fulfilled" && result.value.confirmed) { - session.ownerIds.delete(ownerId); - continue; + for (const [session] of shutdowns) startSessionShutdown(session); + for (const [, , start] of retiring) start(); + return (async () => { + await Promise.all([ + ...requests.map(request => request.completion), + ...initializing.map(entry => entry.promise.catch(() => undefined)), + ]); + const failures: unknown[] = []; + for (const entry of initializing) + if (entry.cleanupFailed) { + if (entry.kernel) retainKernelForOwners(entry.kernel, new Set([ownerId])); + failures.push(entry.cleanupError); + } + const results = await Promise.allSettled(shutdowns.map(([, shutdown]) => shutdown)); + for (let index = 0; index < shutdowns.length; index += 1) { + const [session] = shutdowns[index]; + const result = results[index]; + if (result.status === "fulfilled" && result.value.confirmed) { + session.ownerIds.delete(ownerId); + continue; + } + if (!sessions.has(session.sessionId)) sessions.set(session.sessionId, session); + if (result.status === "rejected") failures.push(result.reason); + else failures.push(unconfirmedShutdownError(session.kernel)); } - const reason = result.status === "rejected" ? result.reason : "not confirmed"; - logger.warn("Python kernel shutdown not confirmed", { sessionId: session.sessionId, reason }); - if (!sessions.has(session.sessionId)) sessions.set(session.sessionId, session); - } + const retiringResults = await Promise.allSettled(retiring.map(([, cleanup]) => cleanup)); + for (let index = 0; index < retiring.length; index += 1) { + const [kernel] = retiring[index]; + const result = retiringResults[index]; + if (result.status === "fulfilled" && result.value.confirmed) continue; + failures.push(result.status === "rejected" ? result.reason : unconfirmedShutdownError(kernel)); + } + if (failures.length) throw failures[0]; + })(); +} + +function unconfirmedShutdownError(kernel: PythonKernel): Error { + const error = new Error(`Python kernel shutdown not confirmed: ${kernel.id}`); + error.name = "PythonKernelShutdownUnconfirmedError"; + return error; } // --------------------------------------------------------------------------- @@ -607,6 +863,7 @@ async function executeWithKernel( options: PythonExecutorOptions | undefined, ): Promise { const settings = options?.settings ?? (await Settings.init()); + if (options) throwIfExecutionCancelled(options); const sink = new OutputSink({ onChunk: options?.onChunk, artifactPath: options?.artifactPath, @@ -623,16 +880,19 @@ async function executeWithKernel( ((event: JsStatusEvent) => { displayOutputs.push({ type: "status", event }); }); - const unregisterBridge = - options?.toolSession && options?.bridgeSessionId && options.bridge - ? registerPyToolBridge(options.bridgeSessionId, options.bridge.capability, { - toolSession: options.toolSession, - signal: options.signal, - emitStatus, - }) - : null; + let unregisterBridge: (() => void) | null = null; try { + if (options) throwIfExecutionCancelled(options); + unregisterBridge = + options?.toolSession && options?.bridgeSessionId && options.bridge + ? registerPyToolBridge(options.bridgeSessionId, options.bridge.capability, { + toolSession: options.toolSession, + signal: options.signal, + emitStatus, + }) + : null; + if (options) throwIfExecutionCancelled(options); executionTimeoutMs = requireRemainingTimeoutMs(deadlineMs); const result = await kernel.execute(code, { signal: options?.signal, @@ -640,6 +900,7 @@ async function executeWithKernel( onChunk: text => sink.push(text), onDisplay: output => void displayOutputs.push(output), }); + if (options) throwIfExecutionCancelled(options); if (result.cancelled) { // Prefer the caller-configured timeout for the user-facing annotation. @@ -714,6 +975,7 @@ async function executeWithKernel( } async function ensureKernelAvailable(cwd: string, options: PythonExecutorOptions): Promise { + throwIfExecutionCancelled(options); const availability = await checkPythonKernelAvailability( cwd, options.runtimeOptions, @@ -723,6 +985,7 @@ async function ensureKernelAvailable(cwd: string, options: PythonExecutorOptions }, options.settings, ); + throwIfExecutionCancelled(options); if (!availability.ok) { throw new Error(availability.reason ?? "Python kernel unavailable"); } @@ -730,10 +993,13 @@ async function ensureKernelAvailable(cwd: string, options: PythonExecutorOptions async function ensureToolBridge(options: PythonExecutorOptions): Promise { if (!options.toolSession || options.bridge) return; + throwIfExecutionCancelled(options); try { const bridge = await ensurePyToolBridge(); + throwIfExecutionCancelled(options); options.bridge = { ...bridge, capability: crypto.randomUUID() }; } catch (err) { + if (isCancellationError(err) || options.signal?.aborted) throw err; logger.warn("Failed to start Python tool bridge", { error: err instanceof Error ? err.message : String(err), }); @@ -741,11 +1007,13 @@ async function ensureToolBridge(options: PythonExecutorOptions): Promise { } async function executePerCall(code: string, cwd: string, options: PythonExecutorOptions): Promise { + throwIfExecutionCancelled(options); if (options.bridge && !options.bridgeSessionId) { options.bridgeSessionId = `py-bridge:${crypto.randomUUID()}`; } const kernel = await startKernel(cwd, options); try { + throwIfExecutionCancelled(options); return await executeWithKernel(kernel, code, options); } finally { await kernel.shutdown().catch(() => undefined); @@ -753,6 +1021,7 @@ async function executePerCall(code: string, cwd: string, options: PythonExecutor } async function executeOnSession(code: string, cwd: string, options: PythonExecutorOptions): Promise { + throwIfExecutionCancelled(options); const sessionId = scopedSessionId( options.sessionId ?? `session:${cwd}`, options.settings, @@ -762,28 +1031,32 @@ async function executeOnSession(code: string, cwd: string, options: PythonExecut options.bridgeSessionId = sessionId; } if (options.reset) { + throwIfExecutionCancelled(options); await resetSession(sessionId); + throwIfExecutionCancelled(options); } const session = await acquireSession(sessionId, cwd, options); + throwIfExecutionCancelled(options); if (options.bridge && session.bridgeCapability) { options.bridge.capability = session.bridgeCapability; } options.onKernelStart?.(session.kernelInstanceId); return await runQueued(session, options, async () => { - if (options.signal?.aborted) { - throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); - } + throwIfExecutionCancelled(options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } if (!session.kernel.isAlive()) { + throwIfExecutionCancelled(options); await replaceSessionKernel(session, cwd, options); + throwIfExecutionCancelled(options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } options.onKernelStart?.(session.kernelInstanceId); } try { + throwIfExecutionCancelled(options); return await executeWithKernel(session.kernel, code, options); } catch (err) { if (isCancellationError(err) || options.signal?.aborted) throw err; @@ -792,7 +1065,9 @@ async function executeOnSession(code: string, cwd: string, options: PythonExecut throw new PythonExecutionCancelledError(false); } // Kernel died during execute. Replace it and retry once on a fresh one. + throwIfExecutionCancelled(options); await replaceSessionKernel(session, cwd, options); + throwIfExecutionCancelled(options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } @@ -807,43 +1082,32 @@ export async function executePythonWithKernel( code: string, options?: PythonExecutorOptions, ): Promise { - return await executeWithKernel(kernel, code, options); + const tracked = beginPythonRequest(options); + try { + return await executeWithKernel(kernel, code, tracked.options); + } catch (err) { + if (isCancellationError(err) || tracked.options.signal?.aborted) { + return createCancelledPythonResult( + isTimedOutCancellation(err, tracked.options.signal), + tracked.options.timeoutMs, + ); + } + throw err; + } finally { + tracked.request.finish(); + } } export async function executePython(code: string, options?: PythonExecutorOptions): Promise { - const cwd = options?.cwd ?? getProjectDir(); - const deadlineMs = getExecutionDeadlineMs(options); - const deadlineController = deadlineMs === undefined ? undefined : new AbortController(); - const combinedController = deadlineController && options?.signal ? new AbortController() : undefined; - const forwardAbort = (): void => combinedController?.abort(options?.signal?.reason); - if (combinedController) { - if (options?.signal?.aborted) forwardAbort(); - else options?.signal?.addEventListener("abort", forwardAbort, { once: true }); - } - const forwardDeadlineAbort = (): void => combinedController?.abort(deadlineController?.signal.reason); - if (combinedController && deadlineController) - deadlineController.signal.addEventListener("abort", forwardDeadlineAbort, { once: true }); - const signal = combinedController?.signal ?? deadlineController?.signal ?? options?.signal; - const remainingMs = getRemainingTimeoutMs(deadlineMs); - const deadlineTimer = - deadlineController && remainingMs !== undefined - ? setTimeout(() => deadlineController.abort(new PythonExecutionCancelledError(true)), Math.max(0, remainingMs)) - : undefined; - deadlineTimer?.unref(); - const executionOptions: PythonExecutorOptions = { - ...(options ?? {}), - signal, - deadlineMs, - }; + const tracked = beginPythonRequest(options, true); + const executionOptions = tracked.options; + const cwd = tracked.cwd!; try { - requireRemainingTimeoutMs(deadlineMs); - if (executionOptions.signal?.aborted) { - throw new PythonExecutionCancelledError( - isTimedOutCancellation(executionOptions.signal.reason, executionOptions.signal), - ); - } + throwIfExecutionCancelled(executionOptions); await ensureKernelAvailable(cwd, executionOptions); + throwIfExecutionCancelled(executionOptions); await ensureToolBridge(executionOptions); + throwIfExecutionCancelled(executionOptions); const kernelMode = executionOptions.kernelMode ?? "session"; if (kernelMode === "per-call") { @@ -856,10 +1120,6 @@ export async function executePython(code: string, options?: PythonExecutorOption } throw err; } finally { - if (deadlineTimer) clearTimeout(deadlineTimer); - if (combinedController) { - options?.signal?.removeEventListener("abort", forwardAbort); - deadlineController?.signal.removeEventListener("abort", forwardDeadlineAbort); - } + tracked.request.finish(); } } diff --git a/packages/coding-agent/src/sdk/session.ts b/packages/coding-agent/src/sdk/session.ts index ebe56be8150..e853d02c6f4 100644 --- a/packages/coding-agent/src/sdk/session.ts +++ b/packages/coding-agent/src/sdk/session.ts @@ -85,6 +85,8 @@ import { Settings, type SkillsSettings } from "../config/settings"; import { resolveEagerTaskDelegation } from "../config/task-delegation"; import { CursorExecHandlers } from "../cursor"; import { EditTool } from "../edit"; +import { disposeVmContextsByOwner } from "../eval/js/context-manager"; +import { disposeKernelSessionsByOwner } from "../eval/py/executor"; import type { MasterModeContext } from "../master-mode/context"; import { describeFoldReceipt } from "../session/fold-coordinator"; import type { BashRestrictionProfile } from "../tools/bash-allowed-prefixes"; @@ -3213,9 +3215,18 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} getSessionFile: () => sessionManager.getSessionFile() ?? null, getSessionAgentDir: () => session?.getSessionAgentDir() ?? options.agentDir ?? settings.getAgentDir(), getEvalKernelOwnerId: () => evalKernelOwnerId, - assertEvalExecutionAllowed: () => session?.assertEvalExecutionAllowed(), - trackEvalExecution: (execution, abortController) => - session ? session.trackEvalExecution(execution, abortController) : execution, + assertEvalExecutionAllowed: () => { + if (!session) throw new Error("Eval execution is unavailable until session initialization completes"); + session.assertEvalExecutionAllowed(); + }, + trackEvalExecution: (execution, abortController) => { + if (!session) { + abortController.abort(new Error("Eval execution is unavailable until session initialization completes")); + void execution.catch(() => {}); + throw new Error("Eval execution is unavailable until session initialization completes"); + } + return session.trackEvalExecution(execution, abortController); + }, getAsyncJobManager: () => asyncJobManager, waitForUserSteering: signal => { if (agent) return agent.waitForSteeringArrival(signal); @@ -3309,7 +3320,10 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} adoptArtifactManager: manager => sessionManager.adoptArtifactManager(manager), releaseArtifactManager: manager => sessionManager.releaseArtifactManager(manager), ensureArtifactManager: () => sessionManager.ensureArtifactManager(), - registerSessionCleanup: cleanup => session?.registerToolSessionTransitionCleanup(cleanup) ?? (() => {}), + registerSessionCleanup: cleanup => { + if (!session) throw new Error("Cannot register session cleanup before session initialization completes"); + return session.registerToolSessionTransitionCleanup(cleanup); + }, mcpConfigPath: explicitMcpConfigPath, settings, authStorage, @@ -6055,6 +6069,9 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} gjcRuntimeSnapshot: gjcRuntimeStore, }; } catch (error) { + // Capture pending owner work before asynchronous startup teardown can yield. + const startupPythonCleanup = disposeKernelSessionsByOwner(evalKernelOwnerId); + void startupPythonCleanup.catch(() => {}); let cleanupDiagnostic: unknown; let ownedMcpCleanupFailed = false; let ownedMcpCleanupError: unknown; @@ -6092,6 +6109,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} await session.awaitDisposeCompletion(); } }); + await attemptCleanup(() => startupPythonCleanup); } else { if (hasRegistered) await attemptCleanup(() => { @@ -6116,18 +6134,8 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} ownedMcpCleanupFailed = true; ownedMcpCleanupError = cleanupError; }); - const evalCleanup = Promise.all([import("../eval/py/executor"), import("../eval/js/context-manager")]); - let evalCleanupModules: Awaited | undefined; - try { - evalCleanupModules = await evalCleanup; - } catch (cleanupError) { - recordCleanupFailure(cleanupError); - } - if (evalCleanupModules) { - const [kernelExecutor, contextManager] = evalCleanupModules; - await attemptCleanup(() => kernelExecutor.disposeKernelSessionsByOwner(evalKernelOwnerId)); - await attemptCleanup(() => contextManager.disposeVmContextsByOwner(evalKernelOwnerId)); - } + await attemptCleanup(() => startupPythonCleanup); + await attemptCleanup(() => disposeVmContextsByOwner(evalKernelOwnerId)); await attemptCleanup(closeOwnedSettings); } if (processCwdClaimed) diff --git a/packages/coding-agent/src/tools/descriptors.ts b/packages/coding-agent/src/tools/descriptors.ts index b206272af6b..eaeff506202 100644 --- a/packages/coding-agent/src/tools/descriptors.ts +++ b/packages/coding-agent/src/tools/descriptors.ts @@ -351,10 +351,14 @@ const loaders: Record = { cwd: session.cwd, settings: session.settings, getCwd: () => session.cwd, + getSessionFile: () => session.getSessionFile(), getSessionId: () => session.getSessionId?.() ?? null, registerSessionCleanup: (cleanup: () => Promise | void) => { - session.registerSessionCleanup?.(cleanup); + return session.registerSessionCleanup?.(cleanup); }, + assertEvalExecutionAllowed: () => session.assertEvalExecutionAllowed?.(), + trackEvalExecution: (execution: Promise, abortController: AbortController) => + session.trackEvalExecution?.(execution, abortController) ?? execution, }), ), job: session => cached("job", () => import("./job")).then(module => module.JobTool.createIf(session)), diff --git a/packages/coding-agent/src/tools/python.ts b/packages/coding-agent/src/tools/python.ts index 42b25f73ca7..79cd977f649 100644 --- a/packages/coding-agent/src/tools/python.ts +++ b/packages/coding-agent/src/tools/python.ts @@ -22,12 +22,20 @@ export function pythonKernelOwnerId(sessionId: string): string { export interface SessionPythonToolInput { /** Working directory for kernel execution (session cwd). */ cwd: string; + /** Resolve the current working directory for each admitted invocation. */ + getCwd?: () => string; /** Session settings used for Python runtime policy. */ settings?: SettingsType; + /** Resolve the session file associated with each admitted invocation. */ + getSessionFile?: () => string | null; /** Resolve the GJC session id used for the kernel owner and transcript paths. */ getSessionId: () => string | null; /** Register cleanup with the current logical session lifecycle. */ - registerSessionCleanup: (cleanup: () => Promise | void) => void; + registerSessionCleanup: (cleanup: () => Promise | void) => (() => void) | void; + /** Reject execution after the owning session has begun disposal. */ + assertEvalExecutionAllowed?: () => void; + /** Track this whole invocation through its transcript append. */ + trackEvalExecution?: (execution: Promise, abortController: AbortController) => Promise; } const paramsSchema = z.object({ @@ -52,6 +60,25 @@ interface TranscriptExecutionResult { truncated: boolean; } +interface PythonGeneration { + readonly sessionId: string; + readonly ownerId: string; + readonly abortControllers: Set; + readonly completions: Set>; + readonly transcripts: Map; + unregisterCleanup?: () => void; + cleanupPromise?: Promise; +} + +interface PythonInvocationContext { + readonly cwd: string; + readonly sessionFile: string | null; + readonly sessionId: string; + readonly ownerId: string; + readonly settings: SettingsType; + readonly generation: PythonGeneration; +} + function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } @@ -71,33 +98,87 @@ function appendFailureTrailer(output: string, appendFailure: string | undefined) } export function createSessionPythonTool(input: SessionPythonToolInput): AgentTool { - let armedForSession: string | null = null; - const seenOwnerIds = new Set(); - let currentTranscript: PythonKernelTranscript | null = null; - - const armCleanupForSession = (sessionId: string): void => { - if (armedForSession === sessionId) return; - input.registerSessionCleanup(async () => { - await Promise.all([...seenOwnerIds].map(ownerId => disposeKernelSessionsByOwner(ownerId))); - seenOwnerIds.clear(); - currentTranscript = null; - armedForSession = null; + const activeGenerations = new Map(); + + const transcriptFor = ( + generation: PythonGeneration, + context: Pick, + kernelInstanceId: string, + ): PythonKernelTranscript => { + const key = JSON.stringify([context.cwd, context.sessionId, kernelInstanceId]); + let transcript = generation.transcripts.get(key); + if (!transcript) { + transcript = openPythonKernelTranscript({ + cwd: context.cwd, + sessionId: context.sessionId, + kernelInstanceId, + }); + generation.transcripts.set(key, transcript); + } + return transcript; + }; + + const retireGeneration = (generation: PythonGeneration): Promise => { + if (generation.cleanupPromise) return generation.cleanupPromise; + const cleanup = Promise.withResolvers(); + generation.cleanupPromise = cleanup.promise; + const completions = [...generation.completions]; + const controllers = [...generation.abortControllers]; + if (activeGenerations.get(generation.sessionId) === generation) { + activeGenerations.delete(generation.sessionId); + } + + let ownerCleanup: Promise; + try { + // Core registers pending operations by this existing owner label before + // availability work begins, so this synchronous call also captures work + // that has not acquired a kernel yet. + ownerCleanup = disposeKernelSessionsByOwner(generation.ownerId); + } catch (error) { + ownerCleanup = Promise.reject(error); + } + for (const controller of controllers) controller.abort(); + void (async () => { + const results = await Promise.allSettled([ownerCleanup, ...completions]); + const coreResult = results[0]!; + if (coreResult.status === "rejected") throw coreResult.reason; + generation.unregisterCleanup?.(); + generation.unregisterCleanup = undefined; + generation.transcripts.clear(); + })().then(cleanup.resolve, error => { + if (generation.cleanupPromise === cleanup.promise) generation.cleanupPromise = undefined; + cleanup.reject(error); }); - armedForSession = sessionId; + return cleanup.promise; + }; + + const generationFor = (sessionId: string): PythonGeneration => { + const current = activeGenerations.get(sessionId); + if (current) return current; + const generation: PythonGeneration = { + sessionId, + ownerId: pythonKernelOwnerId(sessionId), + abortControllers: new Set(), + completions: new Set(), + transcripts: new Map(), + }; + activeGenerations.set(sessionId, generation); + try { + const unregister = input.registerSessionCleanup(() => retireGeneration(generation)); + if (typeof unregister === "function") generation.unregisterCleanup = unregister; + } catch (error) { + activeGenerations.delete(sessionId); + throw error; + } + return generation; }; const appendTranscript = async ( - sessionId: string, + context: PythonInvocationContext, + transcript: PythonKernelTranscript | undefined, code: string, result: TranscriptExecutionResult, ): Promise => { - if (currentTranscript === null) { - currentTranscript = openPythonKernelTranscript({ - cwd: input.cwd, - sessionId, - kernelInstanceId: crypto.randomUUID(), - }); - } const record: PythonTranscriptRecord = { timestamp: new Date().toISOString(), code, @@ -107,11 +188,12 @@ export function createSessionPythonTool(input: SessionPythonToolInput): AgentToo truncated: result.truncated, }; try { - await currentTranscript.append(record); + const target = transcript ?? transcriptFor(context.generation, context, crypto.randomUUID()); + await target.append(record); return undefined; } catch (error) { const message = errorMessage(error); - logger.warn("Python transcript append failed", { sessionId, error: message }); + logger.warn("Python transcript append failed", { sessionId: context.sessionId, error: message }); return message; } }; @@ -135,11 +217,9 @@ export function createSessionPythonTool(input: SessionPythonToolInput): AgentToo isError: true, }; } - armCleanupForSession(sessionId); - const ownerId = pythonKernelOwnerId(sessionId); if (params.action === "clear") { - await disposeKernelSessionsByOwner(ownerId); - currentTranscript = null; + const generation = generationFor(sessionId); + await retireGeneration(generation); return { content: [{ type: "text", text: "Python kernel cleared; the next execute starts a fresh kernel." }], }; @@ -152,48 +232,89 @@ export function createSessionPythonTool(input: SessionPythonToolInput): AgentToo }; } - seenOwnerIds.add(ownerId); - try { - const activeSettings = input.settings ?? Settings.instance; - const result = await executePython(code, { - cwd: input.cwd, - settings: activeSettings, - kernelMode: "session", - sessionId: ownerId, - kernelOwnerId: ownerId, - artifactsDir: sessionIpykernelsArtifactsDir(input.cwd, sessionId), - signal, - onKernelStart: kernelInstanceId => { - if (currentTranscript?.kernelInstanceId !== kernelInstanceId) { - currentTranscript = openPythonKernelTranscript({ - cwd: input.cwd, - sessionId, - kernelInstanceId, - }); - } - }, - }); - const appendFailure = await appendTranscript(sessionId, code, { - output: result.output, - exitCode: result.exitCode ?? null, - cancelled: result.cancelled, - truncated: result.truncated, + const cwd = input.getCwd?.() ?? input.cwd; + const sessionFile = input.getSessionFile?.() ?? null; + const settings = input.settings ?? Settings.instance; + input.assertEvalExecutionAllowed?.(); + const contextGeneration = generationFor(sessionId); + const context: PythonInvocationContext = { + cwd, + sessionFile, + sessionId, + ownerId: contextGeneration.ownerId, + settings, + generation: contextGeneration, + }; + const abortController = new AbortController(); + const abortFromCaller = (): void => abortController.abort(signal?.reason); + if (signal?.aborted) abortFromCaller(); + else signal?.addEventListener("abort", abortFromCaller, { once: true }); + contextGeneration.abortControllers.add(abortController); + + let trackingAccepted = false; + const execution = Promise.resolve() + .then(async (): Promise => { + if (!trackingAccepted) throw new Error("Python execution was not admitted by its owning session."); + input.assertEvalExecutionAllowed?.(); + let transcript: PythonKernelTranscript | undefined; + try { + const result = await executePython(code, { + cwd: context.cwd, + settings: context.settings, + sessionFile: context.sessionFile ?? undefined, + kernelMode: "session", + sessionId: context.ownerId, + kernelOwnerId: context.ownerId, + artifactsDir: sessionIpykernelsArtifactsDir(context.cwd, context.sessionId), + signal: abortController.signal, + onKernelStart: kernelInstanceId => { + transcript = transcriptFor(context.generation, context, kernelInstanceId); + }, + }); + const appendFailure = await appendTranscript(context, transcript, code, { + output: result.output, + exitCode: result.exitCode ?? null, + cancelled: result.cancelled, + truncated: result.truncated, + }); + const output = result.output.length > 0 ? result.output : "(no output)"; + return { content: [{ type: "text", text: appendFailureTrailer(output, appendFailure) }] }; + } catch (error) { + const output = errorMessage(error); + const appendFailure = await appendTranscript(context, transcript, code, { + output, + exitCode: null, + cancelled: isCancellationError(error, abortController.signal), + truncated: false, + }); + return { + content: [{ type: "text", text: appendFailureTrailer(output, appendFailure) }], + isError: true, + }; + } + }) + .finally(() => { + signal?.removeEventListener("abort", abortFromCaller); + contextGeneration.abortControllers.delete(abortController); }); - const output = result.output.length > 0 ? result.output : "(no output)"; - return { content: [{ type: "text", text: appendFailureTrailer(output, appendFailure) }] }; + let completion: Promise; + try { + completion = input.trackEvalExecution ? input.trackEvalExecution(execution, abortController) : execution; + trackingAccepted = true; } catch (error) { - const output = errorMessage(error); - const appendFailure = await appendTranscript(sessionId, code, { - output, - exitCode: null, - cancelled: isCancellationError(error, signal), - truncated: false, - }); - return { - content: [{ type: "text", text: appendFailureTrailer(output, appendFailure) }], - isError: true, - }; + abortController.abort(error); + signal?.removeEventListener("abort", abortFromCaller); + contextGeneration.abortControllers.delete(abortController); + void execution.catch(() => {}); + throw error; } + const settled = completion.then( + () => {}, + () => {}, + ); + contextGeneration.completions.add(settled); + void settled.then(() => contextGeneration.completions.delete(settled)); + return await completion; }, }; const agentTool = { diff --git a/packages/coding-agent/test/agent-session-python-cleanup.test.ts b/packages/coding-agent/test/agent-session-python-cleanup.test.ts index 8054c1ff741..bbce51ffd44 100644 --- a/packages/coding-agent/test/agent-session-python-cleanup.test.ts +++ b/packages/coding-agent/test/agent-session-python-cleanup.test.ts @@ -2,16 +2,20 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; +import type { AgentToolResult } from "@gajae-code/agent-core"; import { getBundledModel } from "@gajae-code/ai"; import { Settings } from "@gajae-code/coding-agent/config/settings"; import { pythonBackend } from "@gajae-code/coding-agent/eval"; -import * as pythonExecutor from "@gajae-code/coding-agent/eval/py/executor"; -import type { PythonKernel as PythonKernelInstance } from "@gajae-code/coding-agent/eval/py/kernel"; -import * as pythonKernel from "@gajae-code/coding-agent/eval/py/kernel"; import { AgentRegistry } from "@gajae-code/coding-agent/registry/agent-registry"; import { createAgentSession, type ExtensionFactory, type WorkspaceTree } from "@gajae-code/coding-agent/sdk"; +import { isSessionDisposalIncompleteError } from "@gajae-code/coding-agent/session/agent-session"; import { SessionManager } from "@gajae-code/coding-agent/session/session-manager"; import { Snowflake } from "@gajae-code/utils"; +import * as pythonExecutor from "../src/eval/py/executor"; +import type { PythonKernel as PythonKernelInstance } from "../src/eval/py/kernel"; +import * as pythonKernel from "../src/eval/py/kernel"; +import { sessionIpykernelsDir } from "../src/gjc-runtime/session-layout"; +import { PYTHON_TOOL_NAME } from "../src/tools/python"; const OK_EXECUTION = { status: "ok", cancelled: false, timedOut: false, stdinRequested: false } as const; @@ -99,7 +103,7 @@ const mockLongPythonDisposeSleepsImmediate = () => { const createSession = async ( tempDir: string, cwd: string, - options: { extensions?: ExtensionFactory[]; sessionManager?: SessionManager } = {}, + options: { extensions?: ExtensionFactory[]; sessionManager?: SessionManager; toolNames?: string[] } = {}, ) => ( await createAgentSession({ @@ -117,7 +121,7 @@ const createSession = async ( slashCommands: [], enableMCP: false, enableLsp: false, - toolNames: ["eval"], + toolNames: options.toolNames ?? ["eval"], }) ).session; @@ -137,6 +141,72 @@ const createMockKernel = () => { }; }; +function toolText(result: AgentToolResult): string { + return result.content.map(block => (block.type === "text" ? block.text : "")).join("\n"); +} + +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +async function waitForProcessFile(filePath: string, timeoutMs = 10_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const file = Bun.file(filePath); + if (await file.exists()) { + const pid = Number((await file.text()).trim()); + if (Number.isSafeInteger(pid) && pid > 0 && isProcessAlive(pid)) return pid; + } + await Bun.sleep(10); + } + throw new Error(`Timed out waiting for a live Python process in ${filePath}`); +} + +async function waitForProcessGone(pid: number, timeoutMs = 10_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (!isProcessAlive(pid)) return true; + await Bun.sleep(25); + } + return !isProcessAlive(pid); +} + +async function transcriptDirectories(cwd: string, sessionId: string): Promise { + const root = sessionIpykernelsDir(cwd, sessionId); + const directories = new Set(); + try { + for await (const file of new Bun.Glob("*/transcript.jsonl").scan({ cwd: root })) { + directories.add(path.dirname(file)); + } + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + return [...directories].sort(); +} + +async function transcriptRecords( + cwd: string, + sessionId: string, + directory: string, +): Promise>> { + const raw = await Bun.file(path.join(sessionIpykernelsDir(cwd, sessionId), directory, "transcript.jsonl")).text(); + return raw + .split(/\r?\n/) + .filter(Boolean) + .map(line => JSON.parse(line) as Record); +} + +async function transcriptBytes(cwd: string, sessionId: string, directory: string): Promise { + return new Uint8Array( + await Bun.file(path.join(sessionIpykernelsDir(cwd, sessionId), directory, "transcript.jsonl")).arrayBuffer(), + ); +} + describe("AgentSession python cleanup", () => { const tempDirs: string[] = []; let originalNullPrompt: string | undefined; @@ -226,6 +296,166 @@ describe("AgentSession python cleanup", () => { expect(unrelatedKernel.execute).toHaveBeenCalledTimes(2); }); + it("joins held generation A cleanup and generation B process shutdown during SDK disposal", async () => { + const { tempDir, cwd } = createTempProject(); + tempDirs.push(tempDir); + const session = await createSession(tempDir, cwd, { toolNames: [PYTHON_TOOL_NAME] }); + const pythonTool = session.getToolByName(PYTHON_TOOL_NAME); + expect(pythonTool).toBeDefined(); + if (!pythonTool) throw new Error("Expected the SDK Python tool"); + + const sessionId = session.sessionManager.getSessionId(); + const pidFileA = path.join(tempDir, "python-generation-a.pid"); + const pidFileB = path.join(tempDir, "python-generation-b.pid"); + const codeA = `import os, time\nwith open(${JSON.stringify(pidFileA)}, "w") as pid_file:\n pid_file.write(str(os.getpid()))\nprint("generation-a-held", flush=True)\ntime.sleep(30)`; + const codeB = `import os\nwith open(${JSON.stringify(pidFileB)}, "w") as pid_file:\n pid_file.write(str(os.getpid()))\nprint("generation-b-ready", flush=True)`; + const realAvailability = pythonKernel.checkPythonKernelAvailability.bind(pythonKernel); + vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockImplementation((...args) => + realAvailability(...args), + ); + const realExecutePython = pythonExecutor.executePython.bind(pythonExecutor); + let executorCalls = 0; + vi.spyOn(pythonExecutor, "executePython").mockImplementation((...args) => { + executorCalls += 1; + return realExecutePython(...args); + }); + const realStart = pythonKernel.PythonKernel.start.bind(pythonKernel.PythonKernel); + const shutdownAStarted = Promise.withResolvers(); + const shutdownBStarted = Promise.withResolvers(); + const releaseShutdownA = Promise.withResolvers(); + let firstKernel: pythonKernel.PythonKernel | undefined; + let kernelStarts = 0; + vi.spyOn(pythonKernel.PythonKernel, "start").mockImplementation(async options => { + const kernel = await realStart(options); + kernelStarts += 1; + if (kernelStarts === 1) { + firstKernel = kernel; + const originalShutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async shutdownOptions => { + shutdownAStarted.resolve(); + await releaseShutdownA.promise; + return await originalShutdown(shutdownOptions); + }; + } else { + const originalShutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async shutdownOptions => { + shutdownBStarted.resolve(); + return await originalShutdown(shutdownOptions); + }; + } + return kernel; + }); + const executionA = pythonTool.execute("python-generation-a", { code: codeA }); + let executionB: Promise | undefined; + let clearSettled = false; + let clearA: Promise | undefined; + let disposeSettled = false; + let disposePromise: Promise | undefined; + let pidA: number | undefined; + let pidB: number | undefined; + let cleanupResults: PromiseSettledResult[] = []; + try { + pidA = await waitForProcessFile(pidFileA); + expect(isProcessAlive(pidA)).toBe(true); + clearA = pythonTool.execute("python-clear-generation-a", { action: "clear" }).then(result => { + clearSettled = true; + return result; + }); + await shutdownAStarted.promise; + expect(firstKernel).toBeDefined(); + expect(isProcessAlive(pidA)).toBe(true); + await Bun.sleep(0); + expect(clearSettled).toBe(false); + + executionB = pythonTool.execute("python-generation-b", { code: codeB }); + const executionBStarted = executionB; + const resultA = await executionA; + expect(resultA.isError).toBeUndefined(); + const resultB = await executionBStarted; + expect(resultB.isError).toBeUndefined(); + pidB = await waitForProcessFile(pidFileB); + expect(isProcessAlive(pidB)).toBe(true); + expect(session.sessionManager.getSessionId()).toBe(sessionId); + const directories = await transcriptDirectories(cwd, sessionId); + expect(directories).toHaveLength(2); + const transcriptSnapshots = await Promise.all( + directories.map(async directory => ({ + directory, + bytes: await transcriptBytes(cwd, sessionId, directory), + records: await transcriptRecords(cwd, sessionId, directory), + })), + ); + expect(transcriptSnapshots.flatMap(transcript => transcript.records)).toEqual( + expect.arrayContaining([ + expect.objectContaining({ code: codeA, cancelled: true }), + expect.objectContaining({ code: codeB, cancelled: false }), + ]), + ); + + disposePromise = session.dispose().then(() => { + disposeSettled = true; + }); + await shutdownBStarted.promise; + expect(isProcessAlive(pidA)).toBe(true); + expect(disposeSettled).toBe(false); + expect(clearSettled).toBe(false); + expect(kernelStarts).toBe(2); + expect(executorCalls).toBe(2); + + releaseShutdownA.resolve(); + await disposePromise; + expect(disposeSettled).toBe(true); + expect(clearSettled).toBe(true); + expect(await waitForProcessGone(pidA)).toBe(true); + expect(await waitForProcessGone(pidB)).toBe(true); + expect(toolText(resultA)).toContain("generation-a-held"); + expect(toolText(resultB)).toContain("generation-b-ready"); + await session.awaitDisposeCompletion(); + for (const transcript of transcriptSnapshots) { + expect(await transcriptBytes(cwd, sessionId, transcript.directory)).toEqual(transcript.bytes); + } + await clearA; + } finally { + releaseShutdownA.resolve(); + if (!disposePromise) disposePromise = session.dispose().then(() => undefined); + const disposeCleanup = (async (): Promise => { + let callerFailure: unknown; + try { + await disposePromise; + } catch (error) { + if (!isSessionDisposalIncompleteError(error)) callerFailure = error; + } + let completionFailure: unknown; + try { + await session.awaitDisposeCompletion(); + } catch (error) { + completionFailure = error; + } + if (callerFailure !== undefined && completionFailure !== undefined) { + throw new AggregateError([callerFailure, completionFailure], "SDK Python disposal failed."); + } + if (callerFailure !== undefined) throw callerFailure; + if (completionFailure !== undefined) throw completionFailure; + })(); + const processWaits = [ + ...(pidA !== undefined ? [waitForProcessGone(pidA)] : []), + ...(pidB !== undefined ? [waitForProcessGone(pidB)] : []), + ]; + cleanupResults = await Promise.allSettled([ + disposeCleanup, + executionA, + ...(executionB ? [executionB] : []), + ...(clearA ? [clearA] : []), + ...processWaits, + ]); + } + const cleanupFailures = cleanupResults.flatMap(result => (result.status === "rejected" ? [result.reason] : [])); + if (cleanupFailures.length === 1) throw cleanupFailures[0]; + if (cleanupFailures.length > 1) throw new AggregateError(cleanupFailures, "SDK Python test cleanup failed."); + const processResults = cleanupResults.slice(4); + expect(processResults.every(result => result.status === "fulfilled" && result.value === true)).toBe(true); + }, 30_000); + it("does not dispose unrelated Python owners when createAgentSession fails after session construction", async () => { const { tempDir, cwd } = createTempProject(); tempDirs.push(tempDir); @@ -430,7 +660,7 @@ describe("AgentSession python cleanup", () => { ); }); - it("detaches retained kernel ownership even when dispose times out waiting for Python work", async () => { + it("retains kernel ownership cleanup until blocked Python work settles", async () => { const { tempDir, cwd } = createTempProject(); tempDirs.push(tempDir); const kernel = new FakeKernel(); @@ -444,50 +674,83 @@ describe("AgentSession python cleanup", () => { vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); const sleepSpy = mockLongPythonDisposeSleepsImmediate(); - const startSpy = vi - .spyOn(pythonKernel.PythonKernel, "start") - .mockResolvedValue(kernel as unknown as PythonKernelInstance); + let kernelStarts = 0; + vi.spyOn(pythonKernel.PythonKernel, "start").mockImplementation(async () => { + kernelStarts += 1; + return kernel as unknown as PythonKernelInstance; + }); const firstSession = await createSession(tempDir, cwd); const secondSession = await createSession(tempDir, cwd); - await secondSession.executePython("print('owner-b warmup')"); - const firstExecution = firstSession.executePython("print('blocked')"); - await blockedExecutionStarted.promise; - let firstExecutionSettled = false; - void firstExecution.finally(() => { - firstExecutionSettled = true; - }); - - let firstDisposed = false; - const disposeFirst = firstSession.dispose().then(() => { - firstDisposed = true; - }); - await disposeFirst; - expect(sleepSpy.mock.calls.some(([duration]) => isPythonDisposeWaitDuration(duration))).toBe(true); - - expect(firstDisposed).toBe(true); - expect(firstExecutionSettled).toBe(false); - expect(kernel.shutdownCalls).toBe(0); - expect(startSpy).toHaveBeenCalledTimes(1); - - blockedExecution.resolve(OK_EXECUTION); - await expect(firstExecution).resolves.toMatchObject({ - cancelled: false, - exitCode: 0, - stdinRequested: false, - }); - await secondSession.executePython("print('owner-b after detach')"); - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.executeCalls).toEqual([ - "print('owner-b warmup')", - "print('blocked')", - "print('owner-b after detach')", - ]); - await secondSession.dispose(); + let firstExecution: Promise | undefined; + let firstDisposeCaller: Promise | undefined; + let firstDisposeCompletion: Promise | undefined; + let secondDisposeCaller: Promise | undefined; + let secondDisposeCompletion: Promise | undefined; + let cleanupFailures: unknown[] = []; + try { + await secondSession.executePython("print('owner-b warmup')"); + firstExecution = firstSession.executePython("print('blocked')"); + await blockedExecutionStarted.promise; + let firstExecutionSettled = false; + void firstExecution.then( + () => { + firstExecutionSettled = true; + }, + () => { + firstExecutionSettled = true; + }, + ); + + let disposeRejectedAsIncomplete = false; + firstDisposeCaller = firstSession.dispose().catch(error => { + if (!isSessionDisposalIncompleteError(error)) throw error; + disposeRejectedAsIncomplete = true; + }); + await firstDisposeCaller; + expect(sleepSpy.mock.calls.some(([duration]) => isPythonDisposeWaitDuration(duration))).toBe(true); + expect(disposeRejectedAsIncomplete).toBe(true); + expect(firstExecutionSettled).toBe(false); + expect(kernel.shutdownCalls).toBe(0); + expect(kernelStarts).toBe(1); - expect(kernel.shutdownCalls).toBe(1); - }, 10000); + blockedExecution.resolve(OK_EXECUTION); + await expect(firstExecution).resolves.toMatchObject({ + cancelled: true, + stdinRequested: false, + }); + firstDisposeCompletion = firstSession.awaitDisposeCompletion(); + await firstDisposeCompletion; + expect(firstExecutionSettled).toBe(true); + expect(kernel.shutdownCalls).toBe(0); + expect(kernelStarts).toBe(1); + await secondSession.executePython("print('owner-b after detach')"); + expect(kernelStarts).toBe(1); + expect(kernel.executeCalls).toEqual([ + "print('owner-b warmup')", + "print('blocked')", + "print('owner-b after detach')", + ]); + secondDisposeCaller = secondSession.dispose(); + secondDisposeCompletion = secondSession.awaitDisposeCompletion(); + await Promise.all([firstDisposeCompletion, secondDisposeCaller, secondDisposeCompletion]); + expect(kernel.shutdownCalls).toBe(1); + } finally { + blockedExecution.resolve(OK_EXECUTION); + const completions = [ + ...(firstExecution ? [firstExecution] : []), + ...(firstDisposeCaller ? [firstDisposeCaller] : []), + ...(secondDisposeCaller ? [secondDisposeCaller] : []), + firstSession.awaitDisposeCompletion(), + secondSession.awaitDisposeCompletion(), + ]; + const results = await Promise.allSettled(completions); + cleanupFailures = results.flatMap(result => (result.status === "rejected" ? [result.reason] : [])); + } + if (cleanupFailures.length === 1) throw cleanupFailures[0]; + if (cleanupFailures.length > 1) throw new AggregateError(cleanupFailures, "Python owner cleanup test failed."); + }, 30_000); it("rejects direct session Python starts once dispose begins", async () => { const { tempDir, cwd } = createTempProject(); 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 47aa6a0389c..60d750d2b97 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 @@ -1,413 +1,1247 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; +import * as path from "node:path"; +import { fileURLToPath } from "node:url"; import { disposeAllKernelSessions, disposeKernelSessionsByOwner, executePython, + executePythonWithKernel, + type PythonResult, } from "@gajae-code/coding-agent/eval/py/executor"; -import type { - KernelExecuteResult, - KernelShutdownResult, - PythonKernel as PythonKernelInstance, -} from "@gajae-code/coding-agent/eval/py/kernel"; import * as pythonKernel from "@gajae-code/coding-agent/eval/py/kernel"; -import { PythonKernel } from "@gajae-code/coding-agent/eval/py/kernel"; - -const OK_RESULT: KernelExecuteResult = { - status: "ok", - cancelled: false, - timedOut: false, - stdinRequested: false, -}; - -type FakeKernelShutdownOptions = { timeoutMs?: number }; - -class FakeKernel { - execute = vi.fn(async () => OK_RESULT); - shutdown = vi.fn( - async (_options?: FakeKernelShutdownOptions): Promise => ({ confirmed: true }), - ); - ping = vi.fn(async () => true); - alive = true; - - isAlive(): boolean { - return this.alive; +import { type KernelShutdownResult, PythonKernel } from "@gajae-code/coding-agent/eval/py/kernel"; +import { TempDir } from "@gajae-code/utils"; + +const originalStart = PythonKernel.start; +const originalAvailability = pythonKernel.checkPythonKernelAvailability; + +function trackStartedKernels(): PythonKernel[] { + const kernels: PythonKernel[] = []; + vi.spyOn(PythonKernel, "start").mockImplementation(async options => { + const kernel = await originalStart(options); + kernels.push(kernel); + return kernel; + }); + return kernels; +} + +function holdAvailability(cwd: string, release: Promise, entered: () => void): void { + let held = false; + vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockImplementation(async (...args) => { + if (args[0] === cwd && !held) { + held = true; + entered(); + await release; + } + return await originalAvailability(...args); + }); +} + +async function waitForFile(path: string, timeoutMs = 5_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await Bun.file(path).exists()) return; + await Bun.sleep(10); } + throw new Error(`Timed out waiting for Python marker file: ${path}`); } -async function flushMicrotasks(turns = 6): Promise { - for (let turn = 0; turn < turns; turn += 1) { - await Promise.resolve(); +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; } } +async function waitForProcessFile(filePath: string, timeoutMs = 5_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await Bun.file(filePath).exists()) { + const pid = Number((await Bun.file(filePath).text()).trim()); + if (Number.isSafeInteger(pid) && pid > 0 && isProcessAlive(pid)) return pid; + } + await Bun.sleep(10); + } + throw new Error(`Timed out waiting for a live Python process in ${filePath}`); +} + +async function waitForProcessGone(pid: number, timeoutMs = 5_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (!isProcessAlive(pid)) return; + await Bun.sleep(25); + } + expect(isProcessAlive(pid)).toBe(false); +} + +async function flushMicrotasks(turns = 6): Promise { + for (let turn = 0; turn < turns; turn += 1) await Promise.resolve(); +} + afterEach(async () => { await disposeAllKernelSessions(); + PythonKernel.start = originalStart; vi.restoreAllMocks(); }); describe("python executor owner cleanup", () => { - it("keeps shared retained kernels alive until the last owner is disposed", async () => { - const kernel = new FakeKernel(); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - await executePython("1 + 1", { - cwd: "/tmp/shared-owner-kernel", - sessionId: "shared-session", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - await executePython("2 + 2", { - cwd: "/tmp/shared-owner-kernel", - sessionId: "shared-session", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.execute).toHaveBeenCalledTimes(2); - - await disposeKernelSessionsByOwner("owner-a"); - - expect(kernel.shutdown).not.toHaveBeenCalled(); - - await executePython("3 + 3", { - cwd: "/tmp/shared-owner-kernel", - sessionId: "shared-session", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.execute).toHaveBeenCalledTimes(3); - - await disposeKernelSessionsByOwner("owner-b"); - - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("disposes every retained kernel owned by one owner across session ids and cwd values", async () => { - const kernelOne = new FakeKernel(); - const kernelTwo = new FakeKernel(); - const unrelatedKernel = new FakeKernel(); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi - .spyOn(PythonKernel, "start") - .mockResolvedValueOnce(kernelOne as unknown as PythonKernelInstance) - .mockResolvedValueOnce(kernelTwo as unknown as PythonKernelInstance) - .mockResolvedValueOnce(unrelatedKernel as unknown as PythonKernelInstance); - - await executePython("print('one')", { - cwd: "/tmp/owner-a-one", - sessionId: "session-one", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - await executePython("print('two')", { - cwd: "/tmp/owner-a-two", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - await executePython("print('other')", { - cwd: "/tmp/owner-b-one", - sessionId: "session-other", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - - expect(startSpy).toHaveBeenCalledTimes(3); - - await disposeKernelSessionsByOwner("owner-a"); - - expect(kernelOne.shutdown).toHaveBeenCalledTimes(1); - expect(kernelTwo.shutdown).toHaveBeenCalledTimes(1); - expect(unrelatedKernel.shutdown).not.toHaveBeenCalled(); - - await executePython("print('still alive')", { - cwd: "/tmp/owner-b-one", - sessionId: "session-other", - kernelMode: "session", - kernelOwnerId: "owner-b", - }); - - expect(startSpy).toHaveBeenCalledTimes(3); - expect(unrelatedKernel.execute).toHaveBeenCalledTimes(2); - }); - - it("falls back to the retained session id when no explicit owner id is provided during execution", async () => { - const kernel = new FakeKernel(); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - await executePython("1 + 1", { - cwd: "/tmp/fallback-owner-session", - sessionId: "fallback-session", - kernelMode: "session", - }); - - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.execute).toHaveBeenCalledTimes(1); - - await disposeKernelSessionsByOwner("fallback-session"); - - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("does not reattach a kernel after owner disposal has already claimed it", async () => { - const disposingKernel = new FakeKernel(); - const replacementKernel = new FakeKernel(); - const shutdownDeferred = Promise.withResolvers(); - disposingKernel.shutdown = vi.fn(() => shutdownDeferred.promise); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi - .spyOn(PythonKernel, "start") - .mockResolvedValueOnce(disposingKernel as unknown as PythonKernelInstance) - .mockResolvedValueOnce(replacementKernel as unknown as PythonKernelInstance); - - await executePython("1 + 1", { - cwd: "/tmp/disposal-race-kernel", - sessionId: "race-session", - kernelMode: "session", - kernelOwnerId: "owner-a", + it("registers cold-process cleanup synchronously for a borrowed direct request", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-registration-"); + const resultPath = path.join(tempDir.path(), "cold-registration-result.json"); + const pidPath = path.join(tempDir.path(), "cold-registration.pid"); + const readyPath = path.join(tempDir.path(), "cold-registration-ready"); + const pythonCode = `from pathlib import Path\nimport os, time\nPath(${JSON.stringify(pidPath)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyPath)}).touch()\ntime.sleep(60)`; + const probe = ` +import * as lifecycle from ${JSON.stringify(new URL("../../src/runtime/process-lifecycle.ts", import.meta.url).href)}; +import { vi, test } from "bun:test"; +import { PythonKernel } from ${JSON.stringify(new URL("../../src/eval/py/kernel.ts", import.meta.url).href)}; +import * as executor from ${JSON.stringify(new URL("../../src/eval/py/executor.ts", import.meta.url).href)}; +test("cold Python resource registration", async () => { +const originalRegister = lifecycle.registerResourceOwner; +const registrationSpy = vi.spyOn(lifecycle, "registerResourceOwner").mockImplementation((...args) => + Reflect.apply(originalRegister, lifecycle, args), +); +const kernel = await PythonKernel.start({ cwd: ${JSON.stringify(tempDir.path())} }); +let execution; +try { +const registrationsBefore = registrationSpy.mock.calls.filter(([name]) => name === "python-kernel-sessions").length; +let executionSettled = false; +execution = executor.executePythonWithKernel(kernel, ${JSON.stringify(pythonCode)}) + .finally(() => { executionSettled = true; }); +const registeredSynchronously = registrationSpy.mock.calls + .filter(([name]) => name === "python-kernel-sessions").length === registrationsBefore + 1; +if (!registeredSynchronously || executionSettled) throw new Error("Cold Python process cleanup was not registered synchronously"); +const readyPath = ${JSON.stringify(readyPath)}; +const pidPath = ${JSON.stringify(pidPath)}; +const deadline = Date.now() + 5_000; +while (!(await Bun.file(readyPath).exists()) && Date.now() < deadline) await Bun.sleep(10); +if (!(await Bun.file(readyPath).exists())) throw new Error("Borrowed Python execution did not start"); +const pid = Number((await Bun.file(pidPath).text()).trim()); +await executor.disposeAllKernelSessions(); +const result = await execution; +const borrowedKernelAlive = kernel.isAlive(); +const shutdown = await kernel.shutdown(); +await lifecycle.disposeAllResourceOwners(); +const resourceOwnersAfterCleanup = lifecycle.resourceOwnerCount(); +await Bun.write(${JSON.stringify(resultPath)}, JSON.stringify({ + pid, + registeredSynchronously, + resultCancelled: result.cancelled, + executionSettled, + borrowedKernelAlive, + shutdownConfirmed: shutdown.confirmed, + resourceOwnersAfterCleanup, +})); +} finally { + await executor?.disposeAllKernelSessions(); + await execution?.catch(() => undefined); + if (kernel.isAlive()) { + const shutdown = await kernel.shutdown(); + if (!shutdown.confirmed) throw new Error("Borrowed Python kernel shutdown was not confirmed"); + } + await lifecycle.disposeAllResourceOwners(); +} +}, 30_000); +`; + const probePath = path.join(tempDir.path(), "cold-registration.test.ts"); + await Bun.write(probePath, probe); + const child = Bun.spawn([process.execPath, "test", probePath], { + cwd: path.resolve(fileURLToPath(new URL("../../../../", import.meta.url))), + stdout: "pipe", + stderr: "pipe", }); - - const disposal = disposeKernelSessionsByOwner("owner-a"); - await executePython("2 + 2", { - cwd: "/tmp/disposal-race-kernel", - sessionId: "race-session", - kernelMode: "session", - kernelOwnerId: "owner-b", + const [exitCode, stdout, stderr] = await Promise.all([ + child.exited, + new Response(child.stdout).text(), + new Response(child.stderr).text(), + ]); + expect(exitCode, `${stdout}\n${stderr}`).toBe(0); + const result = JSON.parse(await Bun.file(resultPath).text()) as { + pid: number; + registeredSynchronously: boolean; + resultCancelled: boolean; + executionSettled: boolean; + borrowedKernelAlive: boolean; + shutdownConfirmed: boolean; + resourceOwnersAfterCleanup: number; + }; + expect(result.pid).toBeGreaterThan(0); + expect(result.registeredSynchronously).toBe(true); + expect(result.resultCancelled).toBe(true); + expect(result.executionSettled).toBe(true); + expect(result.borrowedKernelAlive).toBe(true); + expect(result.shutdownConfirmed).toBe(true); + expect(result.resourceOwnersAfterCleanup).toBe(0); + }, 30_000); + + it("shares a retained physical kernel across owner labels until the last owner detaches", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-shared-"); + const kernels = trackStartedKernels(); + try { + const options = { + cwd: tempDir.path(), + sessionId: "shared-session", + kernelMode: "session" as const, + }; + const first = await executePython("print('owner-a')", { ...options, kernelOwnerId: "owner-a" }); + const second = await executePython("print('owner-b')", { ...options, kernelOwnerId: "owner-b" }); + + expect(first.exitCode).toBe(0); + expect(second.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + + await disposeKernelSessionsByOwner("owner-a"); + expect(kernels[0].isAlive()).toBe(true); + + const stillShared = await executePython("print('still shared')", { ...options, kernelOwnerId: "owner-b" }); + expect(stillShared.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + + await disposeKernelSessionsByOwner("owner-b"); + expect(kernels[0].isAlive()).toBe(false); + } finally { + await disposeAllKernelSessions(); + } + }, 30_000); + + it("disposes all physical sessions for one owner without touching another owner's process", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-sessions-"); + const kernels = trackStartedKernels(); + try { + const first = await executePython("print('one')", { + cwd: tempDir.path(), + sessionId: "owner-a-one", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + const second = await executePython("print('two')", { + cwd: tempDir.path(), + sessionId: "owner-a-two", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + const unrelated = await executePython("print('other')", { + cwd: tempDir.path(), + sessionId: "owner-b-one", + kernelMode: "session", + kernelOwnerId: "owner-b", + }); + expect([first.exitCode, second.exitCode, unrelated.exitCode]).toEqual([0, 0, 0]); + expect(kernels).toHaveLength(3); + + await disposeKernelSessionsByOwner("owner-a"); + expect(kernels[0].isAlive()).toBe(false); + expect(kernels[1].isAlive()).toBe(false); + expect(kernels[2].isAlive()).toBe(true); + + const stillUnrelated = await executePython("print('still alive')", { + cwd: tempDir.path(), + sessionId: "owner-b-one", + kernelMode: "session", + kernelOwnerId: "owner-b", + }); + expect(stillUnrelated.exitCode).toBe(0); + expect(kernels).toHaveLength(3); + } finally { + await disposeAllKernelSessions(); + } + }, 30_000); + + it("uses the retained session id as the existing fallback owner label", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-fallback-"); + const kernels = trackStartedKernels(); + try { + const result = await executePython("print('fallback')", { + cwd: tempDir.path(), + sessionId: "fallback-session", + kernelMode: "session", + }); + expect(result.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + + await disposeKernelSessionsByOwner("fallback-session"); + expect(kernels[0].isAlive()).toBe(false); + } finally { + await disposeAllKernelSessions(); + } + }, 30_000); + + it("tracks a direct execution lifetime without claiming shutdown authority over its borrowed kernel", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-borrowed-"); + const kernel = await PythonKernel.start({ cwd: tempDir.path() }); + const readyFile = `${tempDir.path()}/borrowed-execution-ready`; + let execution: Promise | undefined; + try { + execution = executePythonWithKernel( + kernel, + `from pathlib import Path\nimport time\nPath(${JSON.stringify(readyFile)}).touch()\ntime.sleep(60)`, + { kernelOwnerId: "borrowed-owner", timeoutMs: 30_000 }, + ); + await waitForFile(readyFile); + await disposeKernelSessionsByOwner("borrowed-owner"); + const result = await execution; + expect(result.cancelled).toBe(true); + expect(kernel.isAlive()).toBe(true); + } finally { + await disposeKernelSessionsByOwner("borrowed-owner"); + await execution?.catch(() => undefined); + await kernel.shutdown(); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("finishes cleanup of the captured physical kernel without shutting down a same-name successor", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-replacement-"); + const kernels = trackStartedKernels(); + const shutdownStarted = Promise.withResolvers(); + const releaseShutdown = Promise.withResolvers(); + let cleanup: Promise | undefined; + let repeatedCleanup: Promise | undefined; + let shutdownCalls = 0; + try { + const options = { + cwd: tempDir.path(), + sessionId: "replacement-session", + kernelMode: "session" as const, + }; + const first = await executePython("print('A')", { ...options, kernelOwnerId: "owner-a" }); + expect(first.exitCode).toBe(0); + const kernelA = kernels[0]; + const shutdown = kernelA.shutdown.bind(kernelA); + kernelA.shutdown = async shutdownOptions => { + shutdownCalls += 1; + shutdownStarted.resolve(); + await releaseShutdown.promise; + return await shutdown(shutdownOptions); + }; + + cleanup = disposeKernelSessionsByOwner("owner-a"); + await shutdownStarted.promise; + repeatedCleanup = disposeKernelSessionsByOwner("owner-a"); + let firstSettled = false; + let repeatedSettled = false; + void cleanup.then(() => { + firstSettled = true; + }); + void repeatedCleanup.then(() => { + repeatedSettled = true; + }); + await flushMicrotasks(); + expect(firstSettled).toBe(false); + expect(repeatedSettled).toBe(false); + expect(shutdownCalls).toBe(1); + + const second = await executePython("print('B')", { ...options, kernelOwnerId: "owner-b" }); + expect(second.exitCode).toBe(0); + expect(kernels).toHaveLength(2); + expect(kernels[1].isAlive()).toBe(true); + + releaseShutdown.resolve(); + await Promise.all([cleanup, repeatedCleanup]); + expect(shutdownCalls).toBe(1); + expect(kernelA.isAlive()).toBe(false); + expect(kernels[1].isAlive()).toBe(true); + } finally { + releaseShutdown.resolve(); + await (cleanup ?? disposeKernelSessionsByOwner("owner-a")); + await repeatedCleanup; + await disposeAllKernelSessions(); + } + }, 30_000); + + it("captures a same-label successor during earlier held cleanup and joins both physical shutdowns", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-same-label-replacement-"); + const kernels = trackStartedKernels(); + const shutdownAStarted = Promise.withResolvers(); + const releaseShutdownA = Promise.withResolvers(); + const pidAFile = `${tempDir.path()}/kernel-a.pid`; + const pidBFile = `${tempDir.path()}/kernel-b.pid`; + let pidA: number | undefined; + let pidB: number | undefined; + let executionB: Promise | undefined; + let cleanupA: Promise | undefined; + let cleanupB: Promise | undefined; + let shutdownCallsA = 0; + let bodyFailed = false; + let bodyError: unknown; + let cleanupFailed = false; + let cleanupError: unknown; + try { + const options = { + cwd: tempDir.path(), + sessionId: "same-label-replacement-session", + kernelMode: "session" as const, + kernelOwnerId: "same-owner-label", + }; + const first = await executePython( + `from pathlib import Path\nimport os\nPath(${JSON.stringify(pidAFile)}).write_text(str(os.getpid()))`, + options, + ); + expect(first.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + pidA = await waitForProcessFile(pidAFile); + + const kernelA = kernels[0]; + const originalShutdownA = kernelA.shutdown.bind(kernelA); + kernelA.shutdown = async shutdownOptions => { + shutdownCallsA += 1; + shutdownAStarted.resolve(); + await releaseShutdownA.promise; + return await originalShutdownA(shutdownOptions); + }; + cleanupA = disposeKernelSessionsByOwner(options.kernelOwnerId); + await shutdownAStarted.promise; + expect(isProcessAlive(pidA)).toBe(true); + + const readyBFile = `${tempDir.path()}/kernel-b.ready`; + executionB = executePython( + `from pathlib import Path\nimport os, time\nPath(${JSON.stringify(pidBFile)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyBFile)}).touch()\ntime.sleep(60)`, + options, + ); + await waitForFile(readyBFile); + pidB = await waitForProcessFile(pidBFile); + expect(kernels).toHaveLength(2); + expect(kernels[0].isAlive()).toBe(true); + expect(kernels[1].isAlive()).toBe(true); + + cleanupB = disposeKernelSessionsByOwner(options.kernelOwnerId); + let cleanupASettled = false; + let cleanupBSettled = false; + void cleanupA.then(() => { + cleanupASettled = true; + }); + void cleanupB.then(() => { + cleanupBSettled = true; + }); + await waitForProcessGone(pidB); + await flushMicrotasks(); + expect(cleanupASettled).toBe(false); + expect(cleanupBSettled).toBe(false); + expect(isProcessAlive(pidA)).toBe(true); + expect(kernels[0].isAlive()).toBe(true); + expect(kernels[1].isAlive()).toBe(false); + expect(shutdownCallsA).toBe(1); + + releaseShutdownA.resolve(); + await Promise.all([cleanupA, cleanupB]); + await executionB; + await waitForProcessGone(pidA); + expect(cleanupASettled).toBe(true); + expect(cleanupBSettled).toBe(true); + expect(isProcessAlive(pidA)).toBe(false); + expect(isProcessAlive(pidB)).toBe(false); + expect(kernels[0].isAlive()).toBe(false); + expect(kernels[1].isAlive()).toBe(false); + expect(shutdownCallsA).toBe(1); + } catch (error) { + bodyFailed = true; + bodyError = error; + } finally { + releaseShutdownA.resolve(); + const finalCleanup = cleanupB ?? disposeKernelSessionsByOwner("same-owner-label"); + const settled = await Promise.allSettled([ + ...(cleanupA ? [cleanupA] : []), + finalCleanup, + ...(executionB ? [executionB] : []), + disposeAllKernelSessions(), + ]); + const failed = settled.find(result => result.status === "rejected"); + if (failed?.status === "rejected") { + cleanupFailed = true; + cleanupError = failed.reason; + } + } + if (bodyFailed) throw bodyError; + if (cleanupFailed) throw cleanupError; + }, 30_000); + + it("publishes every physical retirement before an abort listener reenters owner cleanup", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-reentrant-cleanup-"); + const kernels: PythonKernel[] = []; + const shutdownStarted = Promise.withResolvers(); + const shutdownBStarted = Promise.withResolvers(); + const releaseShutdown = Promise.withResolvers(); + const releaseShutdownB = Promise.withResolvers(); + const readyFile = `${tempDir.path()}/reentrant-running.ready`; + const pidFile = `${tempDir.path()}/reentrant-running.pid`; + const pidBFile = `${tempDir.path()}/reentrant-second.pid`; + const ownerId = "reentrant-owner"; + let pid: number | undefined; + let pidB: number | undefined; + let execution: Promise | undefined; + let outerCleanup: Promise | undefined; + let nestedCleanup: Promise | undefined; + let nestedSettled = false; + let shutdownCalls = 0; + let shutdownCallsB = 0; + let bodyFailed = false; + let bodyError: unknown; + let cleanupFailed = false; + let cleanupError: unknown; + vi.spyOn(PythonKernel, "start").mockImplementation(async options => { + const kernel = await originalStart(options); + kernels.push(kernel); + if (options.cwd === tempDir.path()) { + const originalExecute = kernel.execute.bind(kernel); + kernel.execute = async (code, executeOptions) => { + if (code.includes("reentrant-cleanup-session")) { + executeOptions?.signal?.addEventListener( + "abort", + () => { + nestedCleanup = disposeKernelSessionsByOwner(ownerId); + void nestedCleanup.then(() => { + nestedSettled = true; + }); + }, + { once: true }, + ); + } + return await originalExecute(code, executeOptions); + }; + } + return kernel; }); - - expect(startSpy).toHaveBeenCalledTimes(2); - expect(disposingKernel.execute).toHaveBeenCalledTimes(1); - expect(replacementKernel.execute).toHaveBeenCalledTimes(1); - expect(disposingKernel.shutdown).toHaveBeenCalledTimes(1); - expect(replacementKernel.shutdown).not.toHaveBeenCalled(); - - shutdownDeferred.resolve({ confirmed: true }); - await disposal; - - await disposeKernelSessionsByOwner("owner-b"); - expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("returns a cancelled result when a dead session restart shutdown times out", async () => { - const kernel = new FakeKernel(); - kernel.alive = false; - let shutdownCallCount = 0; - kernel.shutdown = vi.fn(async (options?: FakeKernelShutdownOptions): Promise => { - shutdownCallCount += 1; - if (shutdownCallCount > 1) { - return { confirmed: true }; + try { + execution = executePython( + `# reentrant-cleanup-session\nfrom pathlib import Path\nimport os, time\nPath(${JSON.stringify(pidFile)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyFile)}).touch()\ntime.sleep(60)`, + { + cwd: tempDir.path(), + sessionId: "reentrant-cleanup-session", + kernelMode: "session", + kernelOwnerId: ownerId, + }, + ); + await waitForFile(readyFile); + pid = await waitForProcessFile(pidFile); + const kernel = kernels[0]; + const second = await executePython( + `from pathlib import Path\nimport os\nPath(${JSON.stringify(pidBFile)}).write_text(str(os.getpid()))`, + { + cwd: tempDir.path(), + sessionId: "reentrant-second-session", + kernelMode: "session", + kernelOwnerId: ownerId, + }, + ); + expect(second.exitCode).toBe(0); + pidB = await waitForProcessFile(pidBFile); + const kernelB = kernels[1]; + const originalShutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async options => { + shutdownCalls += 1; + shutdownStarted.resolve(); + await releaseShutdown.promise; + return await originalShutdown(options); + }; + const originalShutdownB = kernelB.shutdown.bind(kernelB); + kernelB.shutdown = async options => { + shutdownCallsB += 1; + shutdownBStarted.resolve(); + await releaseShutdownB.promise; + return await originalShutdownB(options); + }; + + outerCleanup = disposeKernelSessionsByOwner(ownerId); + await Promise.all([shutdownStarted.promise, shutdownBStarted.promise]); + await flushMicrotasks(); + expect(nestedCleanup).toBeDefined(); + expect(nestedSettled).toBe(false); + expect(isProcessAlive(pid)).toBe(true); + expect(isProcessAlive(pidB)).toBe(true); + expect(shutdownCalls).toBe(1); + expect(shutdownCallsB).toBe(1); + + releaseShutdown.resolve(); + await waitForProcessGone(pid); + await flushMicrotasks(); + expect(isProcessAlive(pidB)).toBe(true); + expect(kernels[1].isAlive()).toBe(true); + expect(nestedSettled).toBe(false); + expect(shutdownCallsB).toBe(1); + + releaseShutdownB.resolve(); + await Promise.all([outerCleanup, nestedCleanup]); + const result = await execution; + expect(result.cancelled).toBe(true); + await waitForProcessGone(pidB); + expect(nestedSettled).toBe(true); + expect(shutdownCalls).toBe(1); + expect(shutdownCallsB).toBe(1); + } catch (error) { + bodyFailed = true; + bodyError = error; + } finally { + releaseShutdown.resolve(); + releaseShutdownB.resolve(); + const finalCleanup = outerCleanup ?? disposeKernelSessionsByOwner(ownerId); + const settled = await Promise.allSettled([ + finalCleanup, + ...(nestedCleanup ? [nestedCleanup] : []), + ...(execution ? [execution] : []), + disposeAllKernelSessions(), + ]); + const failed = settled.find(result => result.status === "rejected"); + if (failed?.status === "rejected") { + cleanupFailed = true; + cleanupError = failed.reason; } - const { promise, reject } = Promise.withResolvers(); - const timer = setTimeout( - () => reject(new DOMException("Python kernel shutdown timed out", "TimeoutError")), - options?.timeoutMs ?? 0, + } + if (bodyFailed) throw bodyError; + if (cleanupFailed) throw cleanupError; + }, 30_000); + + it("a second global cleanup captures a new same-name process while the first retirement is held", async () => { + using tempDir = TempDir.createSync("@gjc-python-global-recapture-"); + const kernels = trackStartedKernels(); + const releaseShutdownA = Promise.withResolvers(); + const shutdownAStarted = Promise.withResolvers(); + const readyA = `${tempDir.path()}/global-a.ready`; + const readyB = `${tempDir.path()}/global-b.ready`; + const pidAFile = `${tempDir.path()}/global-a.pid`; + const pidBFile = `${tempDir.path()}/global-b.pid`; + let pidA: number | undefined; + let pidB: number | undefined; + let executionA: Promise | undefined; + let executionB: Promise | undefined; + let cleanupA: Promise | undefined; + let cleanupB: Promise | undefined; + let shutdownCallsA = 0; + let bodyFailed = false; + let bodyError: unknown; + let cleanupFailed = false; + let cleanupError: unknown; + try { + const options = { + cwd: tempDir.path(), + sessionId: "global-recapture-session", + kernelMode: "session" as const, + }; + executionA = executePython( + `from pathlib import Path\nimport os, time\nPath(${JSON.stringify(pidAFile)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyA)}).touch()\ntime.sleep(60)`, + options, ); - timer.unref?.(); - return await promise; - }); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValueOnce(kernel as unknown as PythonKernelInstance); - - const result = await executePython("1 + 1", { - cwd: "/tmp/restart-timeout-session", - sessionId: "restart-timeout-session", - kernelMode: "session", - timeoutMs: 100, - }); - - expect(result.cancelled).toBe(true); - expect(result.exitCode).toBeUndefined(); - expect(kernel.shutdown).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: expect.any(Number) })); - expect(startSpy).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 ownerExecutionHang.promise; - }); - globalKernel.execute = vi.fn(async () => { - globalExecutionStarted.resolve(); - return await globalExecutionHang.promise; - }); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - vi.spyOn(PythonKernel, "start") - .mockResolvedValueOnce(ownerKernel as unknown as PythonKernelInstance) - .mockResolvedValueOnce(globalKernel as unknown as PythonKernelInstance); - - void executePython("print('owner hangs')", { - cwd: "/tmp/stuck-owner-cleanup", - sessionId: "stuck-owner-session", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - await ownerExecutionStarted.promise; - - void executePython("print('global hangs')", { - cwd: "/tmp/stuck-global-cleanup", - sessionId: "stuck-global-session", - kernelMode: "session", - }); - await globalExecutionStarted.promise; - - const ownerCleanup = disposeKernelSessionsByOwner("owner-a"); - await flushMicrotasks(); - expect(ownerKernel.shutdown).toHaveBeenCalledTimes(1); - expect(globalKernel.shutdown).not.toHaveBeenCalled(); - await ownerCleanup; - - const globalCleanup = disposeAllKernelSessions(); - 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(); - const unownedRetainedKernel = new FakeKernel(); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi - .spyOn(PythonKernel, "start") - .mockResolvedValueOnce(perCallKernel as unknown as PythonKernelInstance) - .mockResolvedValueOnce(retainedKernel as unknown as PythonKernelInstance) - .mockResolvedValueOnce(unownedRetainedKernel as unknown as PythonKernelInstance); - - await executePython("print('per-call')", { - cwd: "/tmp/per-call-owner", - kernelMode: "per-call", - kernelOwnerId: "owner-a", - }); - await executePython("print('retained')", { - cwd: "/tmp/retained-owner", - sessionId: "retained-session", - kernelMode: "session", - kernelOwnerId: "owner-a", - }); - await executePython("print('unowned')", { - cwd: "/tmp/unowned-retained", - sessionId: "unowned-session", - kernelMode: "session", - }); - - expect(startSpy).toHaveBeenCalledTimes(3); - expect(perCallKernel.shutdown).toHaveBeenCalledTimes(1); - - await disposeKernelSessionsByOwner("owner-a"); - - expect(perCallKernel.shutdown).toHaveBeenCalledTimes(1); - expect(retainedKernel.shutdown).toHaveBeenCalledTimes(1); - expect(unownedRetainedKernel.shutdown).not.toHaveBeenCalled(); - - await disposeAllKernelSessions(); - - expect(unownedRetainedKernel.shutdown).toHaveBeenCalledTimes(1); - }); - it("rejects a queued execute when its session is disposed before the slot runs", async () => { - const kernel = new FakeKernel(); - const executeHang = Promise.withResolvers(); - const executeStarted = Promise.withResolvers(); - kernel.execute = vi.fn(async () => { - executeStarted.resolve(); - return await executeHang.promise; - }); - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - const first = executePython("first", { - cwd: "/tmp/dispose-queue-race", - sessionId: "dispose-queue-session", - kernelMode: "session", - }); - await executeStarted.promise; - - const queued = executePython("queued", { - cwd: "/tmp/dispose-queue-race", - sessionId: "dispose-queue-session", - kernelMode: "session", + await waitForFile(readyA); + pidA = await waitForProcessFile(pidAFile); + const kernelA = kernels[0]; + const originalShutdownA = kernelA.shutdown.bind(kernelA); + kernelA.shutdown = async shutdownOptions => { + shutdownCallsA += 1; + shutdownAStarted.resolve(); + await releaseShutdownA.promise; + return await originalShutdownA(shutdownOptions); + }; + cleanupA = disposeAllKernelSessions(); + await shutdownAStarted.promise; + + executionB = executePython( + `from pathlib import Path\nimport os, time\nPath(${JSON.stringify(pidBFile)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyB)}).touch()\ntime.sleep(60)`, + options, + ); + await waitForFile(readyB); + pidB = await waitForProcessFile(pidBFile); + expect(kernels).toHaveLength(2); + cleanupB = disposeAllKernelSessions(); + let cleanupASettled = false; + let cleanupBSettled = false; + void cleanupA.then(() => { + cleanupASettled = true; + }); + void cleanupB.then(() => { + cleanupBSettled = true; + }); + await waitForProcessGone(pidB); + await flushMicrotasks(); + expect(cleanupASettled).toBe(false); + expect(cleanupBSettled).toBe(false); + expect(isProcessAlive(pidA)).toBe(true); + expect(kernels[0].isAlive()).toBe(true); + expect(kernels[1].isAlive()).toBe(false); + expect(shutdownCallsA).toBe(1); + + releaseShutdownA.resolve(); + await Promise.all([cleanupA, cleanupB]); + const [resultA, resultB] = await Promise.all([executionA, executionB]); + expect(resultA.cancelled).toBe(true); + expect(resultB.cancelled).toBe(true); + await waitForProcessGone(pidA); + expect(cleanupASettled).toBe(true); + expect(cleanupBSettled).toBe(true); + expect(kernels[0].isAlive()).toBe(false); + expect(kernels[1].isAlive()).toBe(false); + } catch (error) { + bodyFailed = true; + bodyError = error; + } finally { + releaseShutdownA.resolve(); + const finalCleanup = cleanupB ?? disposeAllKernelSessions(); + const settled = await Promise.allSettled([ + ...(cleanupA ? [cleanupA] : []), + finalCleanup, + ...(executionA ? [executionA] : []), + ...(executionB ? [executionB] : []), + disposeAllKernelSessions(), + ]); + const failed = settled.find(result => result.status === "rejected"); + if (failed?.status === "rejected") { + cleanupFailed = true; + cleanupError = failed.reason; + } + } + if (bodyFailed) throw bodyError; + if (cleanupFailed) throw cleanupError; + }, 30_000); + + it("rejects unsuccessful physical shutdown and retries the same retained kernel", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-shutdown-retry-"); + const kernels = trackStartedKernels(); + const pidFile = `${tempDir.path()}/retry-kernel.pid`; + let pid: number | undefined; + let shutdownCalls = 0; + try { + const result = await executePython( + `from pathlib import Path\nimport os\nPath(${JSON.stringify(pidFile)}).write_text(str(os.getpid()))`, + { + cwd: tempDir.path(), + sessionId: "shutdown-retry-session", + kernelMode: "session", + kernelOwnerId: "retry-owner", + }, + ); + expect(result.exitCode).toBe(0); + pid = await waitForProcessFile(pidFile); + const kernel = kernels[0]; + const originalShutdown = kernel.shutdown.bind(kernel); + const originalError = new Error("first real shutdown attempt rejected"); + kernel.shutdown = async options => { + shutdownCalls += 1; + if (shutdownCalls === 1) throw originalError; + return await originalShutdown(options); + }; + + await expect(disposeKernelSessionsByOwner("retry-owner")).rejects.toBe(originalError); + expect(isProcessAlive(pid)).toBe(true); + expect(kernel.isAlive()).toBe(true); + await disposeKernelSessionsByOwner("retry-owner"); + await waitForProcessGone(pid); + expect(shutdownCalls).toBe(2); + expect(kernel.isAlive()).toBe(false); + } finally { + await disposeAllKernelSessions(); + } + }, 30_000); + + it("propagates an actual unconfirmed shutdown result and retries that retained physical kernel", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-shutdown-unconfirmed-"); + const kernels = trackStartedKernels(); + const readyFile = `${tempDir.path()}/unconfirmed-shutdown.ready`; + const pidFile = `${tempDir.path()}/unconfirmed-shutdown.pid`; + const secondShutdownStarted = Promise.withResolvers(); + const releaseSecondShutdown = Promise.withResolvers(); + let pid: number | undefined; + let execution: Promise | undefined; + let firstCleanup: Promise | undefined; + let retryCleanup: Promise | undefined; + let firstShutdownResult: KernelShutdownResult | undefined; + let secondShutdownResult: KernelShutdownResult | undefined; + let secondShutdownReceiver: PythonKernel | undefined; + let shutdownCalls = 0; + let bodyFailed = false; + let bodyError: unknown; + let cleanupFailed = false; + let cleanupError: unknown; + try { + execution = executePython( + `from pathlib import Path\nimport os, signal, time\nsignal.signal(signal.SIGINT, signal.SIG_IGN)\nsignal.signal(signal.SIGTERM, signal.SIG_IGN)\nPath(${JSON.stringify(pidFile)}).write_text(str(os.getpid()))\nPath(${JSON.stringify(readyFile)}).touch()\ntime.sleep(60)`, + { + cwd: tempDir.path(), + sessionId: "shutdown-unconfirmed-session", + kernelMode: "session", + kernelOwnerId: "unconfirmed-owner", + }, + ); + await waitForFile(readyFile); + pid = await waitForProcessFile(pidFile); + const kernel = kernels[0]; + const originalShutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async function (this: PythonKernel, options?: Parameters[0]) { + const call = ++shutdownCalls; + if (call === 2) { + secondShutdownReceiver = this; + secondShutdownStarted.resolve(); + await releaseSecondShutdown.promise; + } + const result = await originalShutdown(call === 1 ? { ...options, timeoutMs: 0 } : options); + if (call === 1) firstShutdownResult = result; + else secondShutdownResult = result; + return result; + }; + + firstCleanup = disposeKernelSessionsByOwner("unconfirmed-owner"); + await expect(firstCleanup).rejects.toMatchObject({ + name: "PythonKernelShutdownUnconfirmedError", + }); + expect(kernels).toHaveLength(1); + expect(kernels[0]).toBe(kernel); + const executionResult = await execution; + expect(executionResult.cancelled).toBe(true); + expect(firstShutdownResult).toEqual({ confirmed: false }); + expect(shutdownCalls).toBe(1); + + await waitForProcessGone(pid); + retryCleanup = disposeKernelSessionsByOwner("unconfirmed-owner"); + const retryStartTimeout = Promise.withResolvers(); + const retryStartTimer = setTimeout(retryStartTimeout.resolve, 5_000); + try { + await Promise.race([ + secondShutdownStarted.promise, + retryStartTimeout.promise.then(() => { + throw new Error("Timed out waiting for retained Python kernel shutdown retry"); + }), + ]); + } finally { + clearTimeout(retryStartTimer); + } + let retrySettled = false; + void retryCleanup.then( + () => { + retrySettled = true; + }, + () => { + retrySettled = true; + }, + ); + await flushMicrotasks(); + expect(retrySettled).toBe(false); + expect(shutdownCalls).toBe(2); + expect(secondShutdownReceiver).toBe(kernel); + releaseSecondShutdown.resolve(); + await retryCleanup; + await waitForProcessGone(pid); + expect(secondShutdownResult).toEqual({ confirmed: true }); + expect(shutdownCalls).toBe(2); + expect(secondShutdownReceiver).toBe(kernel); + expect(kernels).toHaveLength(1); + expect(kernels[0]).toBe(kernel); + expect(kernel.isAlive()).toBe(false); + } catch (error) { + bodyFailed = true; + bodyError = error; + } finally { + releaseSecondShutdown.resolve(); + const settled = await Promise.allSettled([ + ...(retryCleanup ? [retryCleanup] : []), + disposeKernelSessionsByOwner("unconfirmed-owner"), + ...(execution ? [execution] : []), + ...(pid !== undefined ? [waitForProcessGone(pid)] : []), + disposeAllKernelSessions(), + ]); + const failed = settled.find(result => result.status === "rejected"); + if (failed?.status === "rejected") { + cleanupFailed = true; + cleanupError = failed.reason; + } + } + if (bodyFailed) throw bodyError; + if (cleanupFailed) throw cleanupError; + }, 30_000); + + it("joins owner cleanup across held real availability preflight and prevents a late kernel launch", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-preflight-"); + const cwd = tempDir.path(); + const entered = Promise.withResolvers(); + const release = Promise.withResolvers(); + const kernels = trackStartedKernels(); + let execution: Promise | undefined; + let cleanup: Promise | undefined; + holdAvailability(cwd, release.promise, entered.resolve); + try { + execution = executePython("print('must not launch')", { + cwd, + sessionId: "held-preflight", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + await entered.promise; + cleanup = disposeKernelSessionsByOwner("owner-a"); + let cleanupSettled = false; + void cleanup.then(() => { + cleanupSettled = true; + }); + await flushMicrotasks(); + expect(cleanupSettled).toBe(false); + expect(kernels).toHaveLength(0); + + release.resolve(); + const result = await execution; + await cleanup; + expect(result.cancelled).toBe(true); + expect(kernels).toHaveLength(0); + } finally { + release.resolve(); + await (cleanup ?? disposeKernelSessionsByOwner("owner-a")); + await execution?.catch(() => undefined); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("joins a held genuine initializer and shuts down its unpublished real kernel", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-initializer-"); + const cwd = tempDir.path(); + const initialized = Promise.withResolvers(); + const release = Promise.withResolvers(); + const kernels: PythonKernel[] = []; + vi.spyOn(PythonKernel, "start").mockImplementation(async options => { + const kernel = await originalStart(options); + kernels.push(kernel); + if (options.cwd === cwd) { + initialized.resolve(kernel); + await release.promise; + } + return kernel; }); - await flushMicrotasks(); - - await disposeAllKernelSessions(); - executeHang.resolve(OK_RESULT); - - const firstResult = await first; - const queuedResult = await queued; - - expect(firstResult.cancelled).toBe(false); - expect(queuedResult.cancelled).toBe(true); - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.execute).toHaveBeenCalledTimes(1); - expect(kernel.shutdown).toHaveBeenCalledTimes(1); - }); - - it("retains sessions whose kernel shutdown is not confirmed so a later dispose retries", async () => { - const kernel = new FakeKernel(); - const unconfirmedShutdown = vi.fn(async (): Promise => ({ confirmed: false })); - kernel.shutdown = unconfirmedShutdown; - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - await executePython("1", { - cwd: "/tmp/unconfirmed-shutdown", - sessionId: "unconfirmed-shutdown-session", - kernelMode: "session", + let execution: Promise | undefined; + let cleanup: Promise | undefined; + try { + execution = executePython("print('must not execute')", { + cwd, + sessionId: "held-initializer", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + const kernel = await initialized.promise; + expect(kernel.isAlive()).toBe(true); + + cleanup = disposeKernelSessionsByOwner("owner-a"); + let cleanupSettled = false; + void cleanup.then(() => { + cleanupSettled = true; + }); + await flushMicrotasks(); + expect(cleanupSettled).toBe(false); + expect(kernel.isAlive()).toBe(true); + + release.resolve(); + const result = await execution; + await cleanup; + expect(result.cancelled).toBe(true); + expect(result.output).not.toContain("must not execute"); + expect(kernel.isAlive()).toBe(false); + } finally { + release.resolve(); + await (cleanup ?? disposeKernelSessionsByOwner("owner-a")); + await execution?.catch(() => undefined); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("keeps an A waiter joined through synchronous cleanup reentry while B's shared initializer survives", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-shared-initializer-"); + const cwd = tempDir.path(); + const initialized = Promise.withResolvers(); + const release = Promise.withResolvers(); + const waiterAvailability = Promise.withResolvers(); + const readyPath = `${cwd}/initializer-b.ready`; + const pidPath = `${cwd}/initializer-b.pid`; + const kernels: PythonKernel[] = []; + let initializationReached = false; + let executionA: Promise | undefined; + let executionB: Promise | undefined; + let ownerCleanup: Promise | undefined; + let reentrantCleanup: Promise | undefined; + let bodyFailed = false; + let bodyError: unknown; + let cleanupFailed = false; + let cleanupError: unknown; + const controller = new AbortController(); + const options = { + cwd, + sessionId: "shared-held-initializer", + kernelMode: "session" as const, + }; + controller.signal.addEventListener( + "abort", + () => { + reentrantCleanup = disposeKernelSessionsByOwner("owner-a"); + }, + { once: true }, + ); + vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockImplementation(async (...args) => { + if (initializationReached && args[0] === cwd) waiterAvailability.resolve(); + return await originalAvailability(...args); }); - - expect(startSpy).toHaveBeenCalledTimes(1); - - await disposeAllKernelSessions(); - expect(unconfirmedShutdown).toHaveBeenCalledTimes(1); - - // Re-executing the same session must reuse the retained kernel (no new start). - await executePython("2", { - cwd: "/tmp/unconfirmed-shutdown", - sessionId: "unconfirmed-shutdown-session", - kernelMode: "session", + vi.spyOn(PythonKernel, "start").mockImplementation(async startOptions => { + const kernel = await originalStart(startOptions); + kernels.push(kernel); + if (startOptions.cwd === cwd) { + const marked = await kernel.execute( + `from pathlib import Path\nimport os\nPath(${JSON.stringify(pidPath)}).write_text(str(os.getpid()))`, + ); + if (marked.status !== "ok") throw new Error("Could not mark the real shared initializer process"); + initializationReached = true; + initialized.resolve(kernel); + await release.promise; + } + return kernel; }); - expect(startSpy).toHaveBeenCalledTimes(1); - expect(kernel.execute).toHaveBeenCalledTimes(2); - - // Swap to a confirmed shutdown so afterEach can drain the retained session. - const confirmedShutdown = vi.fn(async (): Promise => ({ confirmed: true })); - kernel.shutdown = confirmedShutdown; - await disposeAllKernelSessions(); - expect(confirmedShutdown).toHaveBeenCalledTimes(1); - }); - - it("retains owner mapping when owner-scoped shutdown is not confirmed", async () => { - const kernel = new FakeKernel(); - const unconfirmedShutdown = vi.fn(async (): Promise => ({ confirmed: false })); - kernel.shutdown = unconfirmedShutdown; - vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); - vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); - - await executePython("1", { - cwd: "/tmp/unconfirmed-owner-shutdown", - sessionId: "unconfirmed-owner-shutdown-session", - kernelMode: "session", - kernelOwnerId: "owner-a", + try { + executionB = executePython(`from pathlib import Path\nPath(${JSON.stringify(readyPath)}).touch()`, { + ...options, + kernelOwnerId: "owner-b", + }); + const kernel = await initialized.promise; + const pid = await waitForProcessFile(pidPath); + expect(kernel.isAlive()).toBe(true); + + executionA = executePython("print('A waits on B initializer')", { + ...options, + kernelOwnerId: "owner-a", + signal: controller.signal, + }); + await waiterAvailability.promise; + await flushMicrotasks(12); + + ownerCleanup = disposeKernelSessionsByOwner("owner-a"); + controller.abort(new DOMException("reenter owner cleanup", "AbortError")); + let ownerCleanupSettled = false; + let reentrantCleanupSettled = false; + void ownerCleanup.then(() => { + ownerCleanupSettled = true; + }); + void reentrantCleanup?.then(() => { + reentrantCleanupSettled = true; + }); + expect(reentrantCleanup).toBeDefined(); + const resultA = await executionA; + expect(resultA.cancelled).toBe(true); + await flushMicrotasks(); + expect(ownerCleanupSettled).toBe(false); + expect(reentrantCleanupSettled).toBe(false); + expect(kernels).toHaveLength(1); + expect(kernel.isAlive()).toBe(true); + expect(isProcessAlive(pid)).toBe(true); + + release.resolve(); + const resultB = await executionB; + expect(resultB.exitCode).toBe(0); + await Promise.all([ownerCleanup, reentrantCleanup]); + expect(ownerCleanupSettled).toBe(true); + expect(reentrantCleanupSettled).toBe(true); + expect(kernels).toHaveLength(1); + expect(kernel.isAlive()).toBe(true); + expect(await Bun.file(readyPath).exists()).toBe(true); + + await disposeKernelSessionsByOwner("owner-b"); + await waitForProcessGone(pid); + expect(kernel.isAlive()).toBe(false); + } catch (error) { + bodyFailed = true; + bodyError = error; + } finally { + release.resolve(); + controller.abort(); + const settled = await Promise.allSettled([ + ...(ownerCleanup ? [ownerCleanup] : []), + ...(reentrantCleanup ? [reentrantCleanup] : []), + ...(executionA ? [executionA] : []), + ...(executionB ? [executionB] : []), + disposeKernelSessionsByOwner("owner-b"), + disposeAllKernelSessions(), + ]); + const failed = settled.find(result => result.status === "rejected"); + if (failed?.status === "rejected") { + cleanupFailed = true; + cleanupError = failed.reason; + } + } + if (bodyFailed) throw bodyError; + if (cleanupFailed) throw cleanupError; + }, 30_000); + + it("propagates and retries failed shutdown of a captured late initializer by its original label", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-late-shutdown-retry-"); + const cwd = tempDir.path(); + const initialized = Promise.withResolvers(); + const release = Promise.withResolvers(); + const pidFile = `${cwd}/late-kernel.pid`; + const kernels: PythonKernel[] = []; + const originalError = new Error("late kernel shutdown failed"); + let shutdownCalls = 0; + vi.spyOn(PythonKernel, "start").mockImplementation(async options => { + const kernel = await originalStart(options); + kernels.push(kernel); + if (options.cwd === cwd) { + const shutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async shutdownOptions => { + shutdownCalls += 1; + if (shutdownCalls === 1) throw originalError; + return await shutdown(shutdownOptions); + }; + const marked = await kernel.execute( + `from pathlib import Path\nimport os\nPath(${JSON.stringify(pidFile)}).write_text(str(os.getpid()))`, + ); + if (marked.status !== "ok") throw new Error("Could not mark the real late initializer process"); + initialized.resolve(kernel); + await release.promise; + } + return kernel; }); - - await disposeKernelSessionsByOwner("owner-a"); - expect(unconfirmedShutdown).toHaveBeenCalledTimes(1); - - const confirmedShutdown = vi.fn(async (): Promise => ({ confirmed: true })); - kernel.shutdown = confirmedShutdown; - await disposeKernelSessionsByOwner("owner-a"); - expect(confirmedShutdown).toHaveBeenCalledTimes(1); - }); + let execution: Promise | undefined; + let cleanup: Promise | undefined; + let pid: number | undefined; + try { + execution = executePython("print('late initializer must not publish')", { + cwd, + sessionId: "late-initializer-retry-session", + kernelMode: "session", + kernelOwnerId: "late-retry-owner", + }); + const kernel = await initialized.promise; + pid = await waitForProcessFile(pidFile); + cleanup = disposeKernelSessionsByOwner("late-retry-owner"); + release.resolve(); + await expect(cleanup).rejects.toBe(originalError); + const result = await execution; + expect(result.cancelled).toBe(true); + expect(isProcessAlive(pid)).toBe(true); + expect(kernel.isAlive()).toBe(true); + + await disposeKernelSessionsByOwner("late-retry-owner"); + await waitForProcessGone(pid); + expect(shutdownCalls).toBe(2); + expect(kernel.isAlive()).toBe(false); + } finally { + release.resolve(); + await execution?.catch(() => undefined); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("global cleanup joins preflight but leaves a same-name process created after its captured scope", async () => { + using tempDir = TempDir.createSync("@gjc-python-global-preflight-"); + const cwd = tempDir.path(); + const entered = Promise.withResolvers(); + const release = Promise.withResolvers(); + const kernels = trackStartedKernels(); + let executionA: Promise | undefined; + let globalCleanup: Promise | undefined; + holdAvailability(cwd, release.promise, entered.resolve); + try { + executionA = executePython("print('A must not launch')", { + cwd, + sessionId: "global-captured-session", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + await entered.promise; + globalCleanup = disposeAllKernelSessions(); + let cleanupSettled = false; + void globalCleanup.then(() => { + cleanupSettled = true; + }); + + const firstB = await executePython("print('B survives')", { + cwd, + sessionId: "global-captured-session", + kernelMode: "session", + kernelOwnerId: "owner-b", + }); + expect(firstB.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + expect(kernels[0].isAlive()).toBe(true); + expect(cleanupSettled).toBe(false); + + release.resolve(); + const resultA = await executionA; + await globalCleanup; + expect(resultA.cancelled).toBe(true); + + const secondB = await executePython("print('B still survives')", { + cwd, + sessionId: "global-captured-session", + kernelMode: "session", + kernelOwnerId: "owner-b", + }); + expect(secondB.exitCode).toBe(0); + expect(kernels).toHaveLength(1); + expect(kernels[0].isAlive()).toBe(true); + } finally { + release.resolve(); + await (globalCleanup ?? disposeAllKernelSessions()); + await executionA?.catch(() => undefined); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("settles a real in-flight kernel that ignores SIGINT within the finite outer timeout", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-sigint-"); + const kernels = trackStartedKernels(); + const controller = new AbortController(); + const readyFile = `${tempDir.path()}/ignoring-sigint-ready`; + let execution: Promise | undefined; + try { + execution = executePython( + `import signal, time\nfrom pathlib import Path\nsignal.signal(signal.SIGINT, signal.SIG_IGN)\nPath(${JSON.stringify(readyFile)}).touch()\ntime.sleep(60)`, + { + cwd: tempDir.path(), + sessionId: "ignore-sigint-session", + kernelMode: "session", + kernelOwnerId: "owner-a", + signal: controller.signal, + timeoutMs: 30_000, + }, + ); + await waitForFile(readyFile); + const abortedAt = Date.now(); + controller.abort(new DOMException("cancel ignored-SIGINT cell", "AbortError")); + const result = await execution; + expect(result.cancelled).toBe(true); + expect(Date.now() - abortedAt).toBeGreaterThanOrEqual(4_000); + expect(kernels).toHaveLength(1); + expect(kernels[0].isAlive()).toBe(false); + await disposeKernelSessionsByOwner("owner-a"); + } finally { + controller.abort(); + await execution?.catch(() => undefined); + await disposeAllKernelSessions(); + } + }, 30_000); + + it("global disposal cancels queued real executions before their Python code runs", async () => { + using tempDir = TempDir.createSync("@gjc-python-owner-queue-"); + const kernels = trackStartedKernels(); + const firstMarker = `${tempDir.path()}/first-running`; + const queuedMarker = `${tempDir.path()}/queued-must-not-run`; + let first: Promise | undefined; + let queued: Promise | undefined; + try { + first = executePython( + `from pathlib import Path\nimport time\nPath(${JSON.stringify(firstMarker)}).touch()\ntime.sleep(60)`, + { + cwd: tempDir.path(), + sessionId: "queued-disposal-session", + kernelMode: "session", + }, + ); + await waitForFile(firstMarker); + queued = executePython(`from pathlib import Path\nPath(${JSON.stringify(queuedMarker)}).touch()`, { + cwd: tempDir.path(), + sessionId: "queued-disposal-session", + kernelMode: "session", + }); + await flushMicrotasks(); + await disposeAllKernelSessions(); + const [firstResult, queuedResult] = await Promise.all([first, queued]); + expect(firstResult.cancelled).toBe(true); + expect(queuedResult.cancelled).toBe(true); + expect(await Bun.file(queuedMarker).exists()).toBe(false); + expect(kernels).toHaveLength(1); + expect(kernels[0].isAlive()).toBe(false); + } finally { + await disposeAllKernelSessions(); + await first?.catch(() => undefined); + await queued?.catch(() => undefined); + } + }, 30_000); }); diff --git a/packages/coding-agent/test/tools/python-tool-builtin.test.ts b/packages/coding-agent/test/tools/python-tool-builtin.test.ts index a60d3d126bd..c0720646f86 100644 --- a/packages/coding-agent/test/tools/python-tool-builtin.test.ts +++ b/packages/coding-agent/test/tools/python-tool-builtin.test.ts @@ -5,11 +5,12 @@ import { Agent, type AgentTool, type AgentToolResult } from "@gajae-code/agent-c import { getBundledModel } from "@gajae-code/ai/core"; import { ModelRegistry } from "@gajae-code/coding-agent/config/model-registry"; import { Settings } from "@gajae-code/coding-agent/config/settings"; -import * as pyExecutor from "@gajae-code/coding-agent/eval/py/executor"; -import { AgentSession } from "@gajae-code/coding-agent/session/agent-session"; +import { AgentSession, isSessionDisposalIncompleteError } from "@gajae-code/coding-agent/session/agent-session"; import { AuthStorage } from "@gajae-code/coding-agent/session/auth-storage"; import { SessionManager } from "@gajae-code/coding-agent/session/session-manager"; import { TempDir } from "@gajae-code/utils"; +import * as pyExecutor from "../../src/eval/py/executor"; +import * as pythonKernel from "../../src/eval/py/kernel"; import { sessionIpykernelsArtifactsDir, sessionIpykernelsDir } from "../../src/gjc-runtime/session-layout"; import { BUILTIN_TOOL_DESCRIPTORS, createTools, type ToolSession } from "../../src/tools"; import { PYTHON_TOOL_NAME, pythonKernelOwnerId } from "../../src/tools/python"; @@ -22,6 +23,24 @@ function textOf(result: AgentToolResult): string { return result.content.map(block => (block.type === "text" ? block.text : "")).join("\n"); } +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +async function waitForProcessGone(pid: number, timeoutMs = 10_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (!isProcessAlive(pid)) return true; + await Bun.sleep(25); + } + return !isProcessAlive(pid); +} + function pythonResult(overrides: Partial = {}): pyExecutor.PythonResult { const output = overrides.output ?? "ok"; return { @@ -41,24 +60,32 @@ function pythonResult(overrides: Partial = {}): pyExecu function makeToolSession(options: { cwd: string; + getCwd?: () => string; + getSessionFile?: () => string | null; getSessionId?: () => string | null; - registerSessionCleanup?: (cleanup: () => Promise | void) => void; + settings?: Settings; + registerSessionCleanup?: (cleanup: () => Promise | void) => (() => void) | void; + assertEvalExecutionAllowed?: () => void; + trackEvalExecution?: ToolSession["trackEvalExecution"]; }): ToolSession { const session: ToolSession = { - cwd: options.cwd, + get cwd() { + return options.getCwd?.() ?? options.cwd; + }, hasUI: false, - settings: Settings.isolated(), + settings: options.settings ?? Settings.isolated(), requireYieldTool: false, enableLsp: true, taskDepth: 0, - getSessionFile: () => null, + getSessionFile: options.getSessionFile ?? (() => null), getSessionSpawns: () => null, getSessionId: options.getSessionId ?? (() => TEST_SESSION_ID), }; + if (options.assertEvalExecutionAllowed) session.assertEvalExecutionAllowed = options.assertEvalExecutionAllowed; + if (options.trackEvalExecution) session.trackEvalExecution = options.trackEvalExecution; if (options.registerSessionCleanup) { session.registerSessionCleanup = cleanup => { - options.registerSessionCleanup?.(cleanup); - return () => {}; + return options.registerSessionCleanup?.(cleanup) ?? (() => {}); }; } return session; @@ -66,8 +93,13 @@ function makeToolSession(options: { async function loadPythonTool(options: { cwd: string; + getCwd?: () => string; + getSessionFile?: () => string | null; getSessionId?: () => string | null; - registerSessionCleanup?: (cleanup: () => Promise | void) => void; + settings?: Settings; + registerSessionCleanup?: (cleanup: () => Promise | void) => (() => void) | void; + assertEvalExecutionAllowed?: () => void; + trackEvalExecution?: ToolSession["trackEvalExecution"]; }): Promise { const tool = await BUILTIN_TOOL_DESCRIPTORS[PYTHON_TOOL_NAME].load(makeToolSession(options)); if (!tool) throw new Error("Expected the built-in Python tool to load"); @@ -109,20 +141,22 @@ async function transcriptRecords( async function createAgentSessionFixture(options: { cwd: string; toolRegistry: Map; + sessionManager?: SessionManager; + settings?: Settings; }): Promise<{ session: AgentSession; sessionManager: SessionManager; cleanup: () => Promise }> { const authStorage = await AuthStorage.create(path.join(options.cwd, "testauth.db")); authStorage.setRuntimeApiKey("anthropic", "test-key"); const modelRegistry = new ModelRegistry(authStorage); const model = getBundledModel("anthropic", "claude-sonnet-4-5"); if (!model) throw new Error("Expected bundled anthropic model to exist"); - const sessionManager = SessionManager.create(options.cwd, options.cwd); + const sessionManager = options.sessionManager ?? SessionManager.create(options.cwd, options.cwd); const agent = new Agent({ initialState: { model, systemPrompt: ["Test"], tools: [], messages: [] }, }); const session = new AgentSession({ agent, sessionManager, - settings: Settings.isolated(), + settings: options.settings ?? Settings.isolated(), modelRegistry, toolRegistry: options.toolRegistry, discoveryMode: "all", @@ -131,12 +165,62 @@ async function createAgentSessionFixture(options: { session, sessionManager, cleanup: async () => { - await session.dispose(); + try { + await session.dispose(); + } catch (error) { + if (!isSessionDisposalIncompleteError(error)) throw error; + } + await session.awaitDisposeCompletion(); authStorage.close(); }, }; } +async function createPythonToolSessionFixture(options: { + cwd: string; + getCwd: () => string; + getSessionFile: () => string | null; + getSessionId: () => string | null; + settings: Settings; + sessionManager?: SessionManager; +}): Promise<{ + session: AgentSession; + sessionManager: SessionManager; + pythonTool: AgentTool; + cleanup: () => Promise; +}> { + let fixture: Awaited> | undefined; + const toolSession = makeToolSession({ + cwd: options.cwd, + getCwd: options.getCwd, + getSessionFile: options.getSessionFile, + getSessionId: options.getSessionId, + settings: options.settings, + }); + toolSession.registerSessionCleanup = cleanup => { + if (!fixture) throw new Error("Python cleanup was registered before SDK session construction"); + return fixture.session.registerToolSessionTransitionCleanup(cleanup); + }; + toolSession.assertEvalExecutionAllowed = () => { + if (!fixture) throw new Error("Python execution was admitted before SDK session construction"); + fixture.session.assertEvalExecutionAllowed(); + }; + toolSession.trackEvalExecution = (execution, abortController) => { + if (!fixture) throw new Error("Python execution was tracked before SDK session construction"); + return fixture.session.trackEvalExecution(execution, abortController); + }; + const tools = await createTools(toolSession); + const pythonTool = tools.find(tool => tool.name === PYTHON_TOOL_NAME); + if (!pythonTool) throw new Error("Expected Python tool to be created"); + fixture = await createAgentSessionFixture({ + cwd: options.cwd, + toolRegistry: new Map([[PYTHON_TOOL_NAME, pythonTool]]), + sessionManager: options.sessionManager, + settings: options.settings, + }); + return { ...fixture, pythonTool }; +} + describe("builtin session Python tool", () => { const tempDirs: TempDir[] = []; const sessionCleanups: Array<() => Promise> = []; @@ -208,6 +292,252 @@ describe("builtin session Python tool", () => { expect(secondOptions.artifactsDir).toBe(sessionIpykernelsArtifactsDir(cwd, TEST_SESSION_ID)); }); + it("captures live cwd and session metadata before preflight and tracks through transcript append", async () => { + const cwdA = tempDir(); + const cwdB = tempDir(); + const settings = Settings.isolated(); + const availabilityStarted = Promise.withResolvers(); + const releaseAvailability = Promise.withResolvers(); + const kernelStarted = Promise.withResolvers(); + const shutdownStarted = Promise.withResolvers(); + const shutdownCompleted = Promise.withResolvers(); + const releaseShutdown = Promise.withResolvers(); + const appendStarted = Promise.withResolvers(); + const releaseAppend = Promise.withResolvers(); + const managerA = SessionManager.create(cwdA, cwdA); + const managerB = SessionManager.create(cwdB, cwdB); + let activeManager = managerA; + let clearSettled = false; + let observedKernelStartId: string | undefined; + const realAvailability = pythonKernel.checkPythonKernelAvailability.bind(pythonKernel); + let holdAvailability = true; + const availabilitySpy = vi + .spyOn(pythonKernel, "checkPythonKernelAvailability") + .mockImplementation(async (...args) => { + if (holdAvailability) { + holdAvailability = false; + availabilityStarted.resolve(); + await releaseAvailability.promise; + } + return await realAvailability(...args); + }); + const realStart = pythonKernel.PythonKernel.start.bind(pythonKernel.PythonKernel); + vi.spyOn(pythonKernel.PythonKernel, "start").mockImplementation(async options => { + const kernel = await realStart(options); + const originalShutdown = kernel.shutdown.bind(kernel); + kernel.shutdown = async shutdownOptions => { + shutdownStarted.resolve(); + await releaseShutdown.promise; + try { + return await originalShutdown(shutdownOptions); + } finally { + shutdownCompleted.resolve(); + } + }; + kernelStarted.resolve(); + return kernel; + }); + const realExecute = pyExecutor.executePython.bind(pyExecutor); + const executeSpy = vi.spyOn(pyExecutor, "executePython").mockImplementation((code, options) => { + const originalOnKernelStart = options?.onKernelStart; + return realExecute(code, { + ...options, + onKernelStart: kernelInstanceId => { + observedKernelStartId = kernelInstanceId; + originalOnKernelStart?.(kernelInstanceId); + }, + }); + }); + const realAppendFile = fs.appendFile.bind(fs); + vi.spyOn(fs, "appendFile").mockImplementation(async (filePath, data, options) => { + const result = await realAppendFile(filePath, data, options); + if ( + String(filePath).startsWith(sessionIpykernelsDir(cwdA, managerA.getSessionId())) && + String(filePath).endsWith("transcript.jsonl") + ) { + appendStarted.resolve(); + await releaseAppend.promise; + } + return result; + }); + let sessionForCleanup: Awaited> | undefined; + let executionForCleanup: Promise | undefined; + let clear: Promise | undefined; + let transition: Promise | undefined; + let transitionSettled = false; + let cleanupResults: PromiseSettledResult[] = []; + try { + const session = await createPythonToolSessionFixture({ + cwd: cwdA, + getCwd: () => activeManager.getCwd(), + getSessionFile: () => activeManager.getSessionFile() ?? null, + getSessionId: () => activeManager.getSessionId(), + settings, + sessionManager: managerA, + }); + sessionForCleanup = session; + const sessionFileA = managerA.getSessionFile(); + const sessionIdA = managerA.getSessionId(); + if (!sessionFileA) throw new Error("Expected an actual SDK session file"); + const pidFile = path.join(cwdA, "python-captured.pid"); + const executionCode = `import os\nwith open(${JSON.stringify(pidFile)}, "w") as pid_file:\n pid_file.write(str(os.getpid()))\nprint(os.getpid())`; + const execution = session.pythonTool.execute("python-captured-call", { code: executionCode }); + executionForCleanup = execution; + await availabilityStarted.promise; + expect(session.session.isEvalRunning).toBe(true); + activeManager = managerB; + releaseAvailability.resolve(); + await kernelStarted.promise; + await appendStarted.promise; + expect(session.session.isEvalRunning).toBe(true); + const observedPid = Number((await Bun.file(pidFile).text()).trim()); + expect(Number.isSafeInteger(observedPid) && observedPid > 0).toBe(true); + expect(isProcessAlive(observedPid)).toBe(true); + + activeManager = managerA; + clear = session.pythonTool.execute("python-captured-clear", { action: "clear" }).then(result => { + clearSettled = true; + return result; + }); + await shutdownStarted.promise; + await Bun.sleep(0); + expect(clearSettled).toBe(false); + + const options = executeSpy.mock.calls[0]?.[1]; + if (!options) throw new Error("Expected captured Python executor options"); + expect(availabilitySpy).toHaveBeenCalled(); + expect(options.cwd).toBe(cwdA); + expect(options.sessionFile).toBe(sessionFileA); + expect(options.sessionId).toBe(pythonKernelOwnerId(sessionIdA)); + expect(options.settings).toBe(settings); + expect(options.artifactsDir).toBe(sessionIpykernelsArtifactsDir(cwdA, sessionIdA)); + expect(observedKernelStartId).toBeString(); + const directories = await transcriptDirectories(cwdA, sessionIdA); + expect(directories).toHaveLength(1); + expect(directories[0]).toEndWith(`-${observedKernelStartId}`); + const transcriptPath = path.join(sessionIpykernelsDir(cwdA, sessionIdA), directories[0]!, "transcript.jsonl"); + const transcriptBytesWhileAppendHeld = new Uint8Array(await Bun.file(transcriptPath).arrayBuffer()); + expect(await transcriptDirectories(cwdB, sessionIdA)).toEqual([]); + expect(await transcriptDirectories(cwdB, managerB.getSessionId())).toEqual([]); + expect(await transcriptRecords(cwdA, sessionIdA, directories[0]!)).toEqual([ + expect.objectContaining({ code: executionCode, output: expect.stringContaining(String(observedPid)) }), + ]); + expect(clearSettled).toBe(false); + releaseShutdown.resolve(); + await shutdownCompleted.promise; + expect(await waitForProcessGone(observedPid)).toBe(true); + expect(clearSettled).toBe(false); + expect(session.session.isEvalRunning).toBe(true); + transition = session.session.newSession().then(result => { + transitionSettled = true; + return result; + }); + await Bun.sleep(0); + expect(transitionSettled).toBe(false); + expect(session.session.isEvalRunning).toBe(true); + releaseAppend.resolve(); + const result = await execution; + expect(result.isError).toBeUndefined(); + expect(textOf(result)).toContain(String(observedPid)); + expect(session.session.isEvalRunning).toBe(false); + const cleared = await clear; + expect(cleared.isError).toBeUndefined(); + await expect(transition).resolves.toBe(true); + await Bun.sleep(0); + expect(clearSettled).toBe(true); + expect(transitionSettled).toBe(true); + expect(session.session.isEvalRunning).toBe(false); + expect(new Uint8Array(await Bun.file(transcriptPath).arrayBuffer())).toEqual(transcriptBytesWhileAppendHeld); + } finally { + releaseAvailability.resolve(); + releaseShutdown.resolve(); + releaseAppend.resolve(); + const cleanupTasks: Promise[] = [sessionForCleanup ? sessionForCleanup.cleanup() : managerA.close()]; + if (executionForCleanup) cleanupTasks.push(executionForCleanup); + if (clear) cleanupTasks.push(clear); + if (transition) cleanupTasks.push(transition); + cleanupTasks.push(managerB.close()); + cleanupResults = await Promise.allSettled(cleanupTasks); + } + const cleanupFailures = cleanupResults.flatMap(result => (result.status === "rejected" ? [result.reason] : [])); + if (cleanupFailures.length === 1) throw cleanupFailures[0]; + if (cleanupFailures.length > 1) throw new AggregateError(cleanupFailures, "Python test cleanup failed."); + }, 30_000); + + it("clears and recreates a real Python generation under the actual session lifecycle", async () => { + const cwd = tempDir(); + const settings = Settings.isolated(); + const sessionManager = SessionManager.create(cwd, cwd); + const session = await createPythonToolSessionFixture({ + cwd, + getCwd: () => sessionManager.getCwd(), + getSessionFile: () => sessionManager.getSessionFile() ?? null, + getSessionId: () => sessionManager.getSessionId(), + settings, + sessionManager, + }); + const sessionId = sessionManager.getSessionId(); + let pidA: number | undefined; + let pidB: number | undefined; + let disposed = false; + let cleanupResults: PromiseSettledResult[] = []; + try { + const resultA = await session.pythonTool.execute("python-generation-a", { + code: "import os\ngeneration_a_marker = True\nprint(os.getpid())", + }); + expect(resultA.isError).toBeUndefined(); + pidA = Number(textOf(resultA).match(/\b\d+\b/)?.[0]); + expect(Number.isSafeInteger(pidA) && pidA > 0).toBe(true); + expect(isProcessAlive(pidA)).toBe(true); + + const directoriesA = await transcriptDirectories(cwd, sessionId); + expect(directoriesA).toHaveLength(1); + const bytesA = new Uint8Array( + await Bun.file( + path.join(sessionIpykernelsDir(cwd, sessionId), directoriesA[0]!, "transcript.jsonl"), + ).arrayBuffer(), + ); + + const clearResult = await session.pythonTool.execute("python-generation-clear", { action: "clear" }); + expect(clearResult.isError).toBeUndefined(); + expect(await waitForProcessGone(pidA)).toBe(true); + expect(session.session.isEvalRunning).toBe(false); + + const resultB = await session.pythonTool.execute("python-generation-b", { + code: "import os\nprint(os.getpid())\nprint('generation_a_marker' in globals())", + }); + expect(resultB.isError).toBeUndefined(); + pidB = Number(textOf(resultB).match(/\b\d+\b/)?.[0]); + expect(Number.isSafeInteger(pidB) && pidB > 0).toBe(true); + expect(isProcessAlive(pidB)).toBe(true); + expect(textOf(resultB)).toContain("False"); + + const directoriesB = await transcriptDirectories(cwd, sessionId); + expect(directoriesB).toHaveLength(2); + const transcriptBDirectory = directoriesB.find(directory => directory !== directoriesA[0]); + if (!transcriptBDirectory) throw new Error("Expected a fresh transcript directory after clear"); + const pathA = path.join(sessionIpykernelsDir(cwd, sessionId), directoriesA[0]!, "transcript.jsonl"); + const pathB = path.join(sessionIpykernelsDir(cwd, sessionId), transcriptBDirectory, "transcript.jsonl"); + const bytesB = new Uint8Array(await Bun.file(pathB).arrayBuffer()); + await session.cleanup(); + disposed = true; + expect(await waitForProcessGone(pidB)).toBe(true); + expect(new Uint8Array(await Bun.file(pathA).arrayBuffer())).toEqual(bytesA); + expect(new Uint8Array(await Bun.file(pathB).arrayBuffer())).toEqual(bytesB); + } finally { + cleanupResults = await Promise.allSettled([ + ...(!disposed ? [session.cleanup()] : []), + ...(pidA !== undefined ? [waitForProcessGone(pidA)] : []), + ...(pidB !== undefined ? [waitForProcessGone(pidB)] : []), + ]); + } + const cleanupFailures = cleanupResults.flatMap(result => (result.status === "rejected" ? [result.reason] : [])); + if (cleanupFailures.length === 1) throw cleanupFailures[0]; + if (cleanupFailures.length > 1) throw new AggregateError(cleanupFailures, "Python test cleanup failed."); + const processResults = cleanupResults.slice(disposed ? 0 : 1); + expect(processResults.every(result => result.status === "fulfilled" && result.value === true)).toBe(true); + }, 30_000); + it("records each executor callback lifetime in its own transcript directory", async () => { const cwd = tempDir(); let kernelId = "k1";