From 2b480f7da1e5f18c9bd05d10287c559bb7412c0a Mon Sep 17 00:00:00 2001 From: Carl Date: Sun, 12 Apr 2026 14:45:19 +0200 Subject: [PATCH] fix(coding-agent): addressed python cleanup review findings fixed retained-kernel restart and owner cleanup edge cases during recovery and disposal tracked async user_python hooks during disposal-sensitive execution paths and hardened startup warmup tracking strengthened cleanup and kernel lifecycle regressions to remove deadlocks, false positives, and timing flakes --- packages/coding-agent/src/ipy/executor.ts | 22 ++-- packages/coding-agent/src/sdk.ts | 3 +- .../coding-agent/src/session/agent-session.ts | 28 ++--- packages/coding-agent/src/tools/index.ts | 42 +++++-- .../test/agent-session-python-cleanup.test.ts | 75 ++++++++++--- .../core/python-executor-lifecycle.test.ts | 12 +- .../python-executor-owner-cleanup.test.ts | 96 +++++++++++++--- .../test/core/python-executor-session.test.ts | 2 +- .../core/python-executor.lifecycle.test.ts | 7 +- .../test/core/python-kernel-session.test.ts | 8 +- .../test/core/python-kernel.lifecycle.test.ts | 105 ++++++++++-------- 11 files changed, 278 insertions(+), 122 deletions(-) diff --git a/packages/coding-agent/src/ipy/executor.ts b/packages/coding-agent/src/ipy/executor.ts index ebc0b3964..cd019eca3 100644 --- a/packages/coding-agent/src/ipy/executor.ts +++ b/packages/coding-agent/src/ipy/executor.ts @@ -540,8 +540,8 @@ export async function disposeAllKernelSessions(): Promise { export async function disposeKernelSessionsByOwner(ownerId: string): Promise { const sessionsToDispose: KernelSession[] = []; - for (const session of Array.from(kernelSessions.values())) { - if (session.disposing || !session.ownerIds.delete(ownerId)) continue; + for (const session of new Set([...kernelSessions.values(), ...disposingKernelSessions.values()])) { + if (!session.ownerIds.delete(ownerId)) continue; if (session.ownerIds.size === 0) { sessionsToDispose.push(session); } @@ -641,12 +641,14 @@ function isResourceExhaustionError(error: unknown): boolean { ); } -function clearDisposingKernelSessionTracking(): void { - for (const session of disposingKernelSessions.values()) { +function clearSharedGatewayDisposingKernelSessionTracking(): void { + for (const session of Array.from(disposingKernelSessions.values())) { + if (!session.kernel.isSharedGateway) continue; if (session.heartbeatTimer) { clearInterval(session.heartbeatTimer); session.heartbeatTimer = undefined; } + disposingKernelSessions.delete(session); session.resolveDisposeCapacity?.(); session.resolveDisposeCapacity = undefined; session.disposeCapacityPromise = undefined; @@ -656,8 +658,8 @@ function clearDisposingKernelSessionTracking(): void { session.disposeResultPromise = undefined; session.disposeResultTimeoutMs = undefined; session.nextDisposalRetryAt = undefined; + session.kernelInvalidatedByRecovery = false; } - disposingKernelSessions.clear(); } function markLiveKernelSessionsForRecovery(): void { @@ -676,7 +678,7 @@ async function recoverFromResourceExhaustion(): Promise { logger.warn("Resource exhaustion detected, recovering by restarting shared gateway"); stopCleanupTimer(); markLiveKernelSessionsForRecovery(); - clearDisposingKernelSessionTracking(); + clearSharedGatewayDisposingKernelSessionTracking(); await shutdownSharedGateway(); syncCleanupTimer(); } @@ -750,11 +752,17 @@ async function restartKernelSession( requireRemainingTimeoutMs(options.deadlineMs); try { if (!session.kernelInvalidatedByRecovery) { + const deadKernel = session.dead || !session.kernel.isAlive(); const shutdownTimeoutMs = requireRemainingTimeoutMs(options.deadlineMs); const shutdownResult = await session.kernel.shutdown({ signal: options.signal, timeoutMs: shutdownTimeoutMs }); - if (!shutdownResult.confirmed) { + if (!shutdownResult.confirmed && !deadKernel) { throw new Error("Failed to confirm crashed kernel shutdown before restart"); } + if (!shutdownResult.confirmed) { + logger.warn("Proceeding with retained kernel restart after unconfirmed dead-kernel shutdown", { + sessionId: session.id, + }); + } } const env: Record | undefined = options.sessionFile ? { PI_SESSION_FILE: options.sessionFile } diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 1b4c838c1..075a911a7 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -914,7 +914,8 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} getSessionFile: () => sessionManager.getSessionFile() ?? null, getPythonKernelOwnerId: () => pythonKernelOwnerId, assertPythonExecutionAllowed: () => session?.assertPythonExecutionAllowed(), - trackPythonExecution: (execution, abortController) => session.trackPythonExecution(execution, abortController), + trackPythonExecution: (execution, abortController) => + session ? session.trackPythonExecution(execution, abortController) : execution, getSessionId: () => sessionManager.getSessionId?.() ?? null, getSessionSpawns: () => options.spawns ?? "*", getModelString: () => (hasExplicitModel && model ? formatModelString(model) : undefined), diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index e4abfdcff..802225566 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -5688,22 +5688,22 @@ export class AgentSession { const cwd = this.sessionManager.getCwd(); this.assertPythonExecutionAllowed(); - if (this.#extensionRunner?.hasHandlers("user_python")) { - const hookResult = await this.#extensionRunner.emitUserPython({ - type: "user_python", - code, - excludeFromContext, - cwd, - }); - this.assertPythonExecutionAllowed(); - if (hookResult?.result) { - this.recordPythonResult(code, hookResult.result, options); - return hookResult.result; - } - } - const abortController = new AbortController(); const execution = (async (): Promise => { + if (this.#extensionRunner?.hasHandlers("user_python")) { + const hookResult = await this.#extensionRunner.emitUserPython({ + type: "user_python", + code, + excludeFromContext, + cwd, + }); + this.assertPythonExecutionAllowed(); + if (hookResult?.result) { + this.recordPythonResult(code, hookResult.result, options); + return hookResult.result; + } + } + // Use the same session ID as the Python tool for kernel sharing const sessionFile = this.sessionManager.getSessionFile(); const sessionId = sessionFile ? `session:${sessionFile}:cwd:${cwd}` : `cwd:${cwd}`; diff --git a/packages/coding-agent/src/tools/index.ts b/packages/coding-agent/src/tools/index.ts index b98ba7a17..34b9c82bb 100644 --- a/packages/coding-agent/src/tools/index.ts +++ b/packages/coding-agent/src/tools/index.ts @@ -8,7 +8,7 @@ import type { Settings } from "../config/settings"; import { EditTool } from "../edit"; import type { Skill } from "../extensibility/skills"; import type { InternalUrlRouter } from "../internal-urls"; -import { getPreludeDocs, warmPythonEnvironment } from "../ipy/executor"; +import { getPreludeDocs, resetPreludeDocsCache, warmPythonEnvironment } from "../ipy/executor"; import { checkPythonKernelAvailability } from "../ipy/kernel"; import { LspTool } from "../lsp"; import type { DiscoverableMCPSearchIndex, DiscoverableMCPTool } from "../mcp/discoverable-tool-metadata"; @@ -308,6 +308,8 @@ export async function createTools(session: ToolSession, toolNames?: string[]): P const isTestEnv = isBunTestRuntime(); const forcePythonWarmup = session.forcePythonWarmup === true; const skipPythonWarm = (isTestEnv && !forcePythonWarmup) || $flag("PI_PYTHON_SKIP_CHECK"); + const cachedPreludeDocs = getPreludeDocs(); + const shouldWarmPython = !skipPythonWarm && (forcePythonWarmup || cachedPreludeDocs.length === 0); if (shouldCheckPython) { const availability = await logger.time("createTools:pythonCheck", checkPythonKernelAvailability, session.cwd); pythonAvailable = availability.ok; @@ -315,20 +317,38 @@ export async function createTools(session: ToolSession, toolNames?: string[]): P logger.warn("Python kernel unavailable, falling back to bash", { reason: availability.reason, }); - } else if (!skipPythonWarm && getPreludeDocs().length === 0) { + } else if (shouldWarmPython) { const sessionFile = session.getSessionFile?.() ?? undefined; const kernelOwnerId = session.getPythonKernelOwnerId?.() ?? undefined; const warmSessionId = sessionFile ? `session:${sessionFile}:cwd:${session.cwd}` : `cwd:${session.cwd}`; + const warmupAbortController = new AbortController(); try { - await logger.time( - "createTools:warmPython", - warmPythonEnvironment, - session.cwd, - warmSessionId, - session.settings.get("python.sharedGateway"), - sessionFile, - kernelOwnerId, - ); + session.assertPythonExecutionAllowed?.(); + if (forcePythonWarmup && cachedPreludeDocs.length > 0) { + resetPreludeDocsCache(); + } + const warmupExecution = session.trackPythonExecution + ? logger.time( + "createTools:warmPython", + warmPythonEnvironment, + session.cwd, + warmSessionId, + session.settings.get("python.sharedGateway"), + sessionFile, + kernelOwnerId, + warmupAbortController.signal, + ) + : logger.time( + "createTools:warmPython", + warmPythonEnvironment, + session.cwd, + warmSessionId, + session.settings.get("python.sharedGateway"), + sessionFile, + kernelOwnerId, + ); + await (session.trackPythonExecution?.(warmupExecution, warmupAbortController) ?? warmupExecution); + session.assertPythonExecutionAllowed?.(); } catch (err) { logger.warn("Failed to warm Python environment", { error: err instanceof Error ? err.message : String(err), 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 c137d0d49..dd358b38f 100644 --- a/packages/coding-agent/test/agent-session-python-cleanup.test.ts +++ b/packages/coding-agent/test/agent-session-python-cleanup.test.ts @@ -103,13 +103,22 @@ const createSession = async ( const stubPythonWarmup = () => vi.spyOn(pythonExecutor, "warmPythonEnvironment").mockResolvedValue({ ok: true, docs: [] }); -const createWarmupKernel = (docs: PreludeHelper[] = []) => ({ - introspectPrelude: vi.fn().mockResolvedValue(docs), - execute: vi.fn(async () => OK_EXECUTION), - ping: vi.fn(async () => true), - isAlive: () => true, - shutdown: vi.fn(async () => ({ confirmed: true })), -}); +const createWarmupKernel = (docs: PreludeHelper[] = []) => { + let alive = true; + return { + introspectPrelude: vi.fn().mockResolvedValue(docs), + execute: vi.fn(async () => { + if (!alive) throw new Error("Expected warmup kernel to be restarted after shutdown"); + return OK_EXECUTION; + }), + ping: vi.fn(async () => alive), + isAlive: () => alive, + shutdown: vi.fn(async () => { + alive = false; + return { confirmed: true }; + }), + }; +}; describe("AgentSession python cleanup", () => { const tempDirs: string[] = []; @@ -174,6 +183,19 @@ describe("AgentSession python cleanup", () => { expect(warmedKernel.shutdown).toHaveBeenCalledTimes(1); expect(unrelatedKernel.shutdown).not.toHaveBeenCalled(); + const replacementKernel = createWarmupKernel(); + startSpy.mockResolvedValueOnce(replacementKernel as unknown as PythonKernelInstance); + await pythonExecutor.executePython("print('fresh warmup before')", { + cwd, + sessionId: `cwd:${cwd}`, + kernelMode: "session", + kernelOwnerId: "fresh-owner-before", + }); + expect(startSpy).toHaveBeenCalledTimes(3); + expect(replacementKernel.execute).toHaveBeenCalledTimes(1); + expect(warmedKernel.shutdown).toHaveBeenCalledTimes(1); + expect(warmedKernel.execute).not.toHaveBeenCalled(); + await pythonExecutor.executePython("print('still alive before')", { cwd: unrelatedCwd, sessionId: "unrelated-before-session", @@ -181,7 +203,7 @@ describe("AgentSession python cleanup", () => { kernelOwnerId: "other-owner", }); - expect(startSpy).toHaveBeenCalledTimes(2); + expect(startSpy).toHaveBeenCalledTimes(3); expect(unrelatedKernel.execute).toHaveBeenCalledTimes(2); }); @@ -235,6 +257,19 @@ describe("AgentSession python cleanup", () => { expect(warmedKernel.shutdown).toHaveBeenCalledTimes(1); expect(unrelatedKernel.shutdown).not.toHaveBeenCalled(); + const replacementKernel = createWarmupKernel(); + startSpy.mockResolvedValueOnce(replacementKernel as unknown as PythonKernelInstance); + await pythonExecutor.executePython("print('fresh warmup after')", { + cwd, + sessionId: `cwd:${cwd}`, + kernelMode: "session", + kernelOwnerId: "fresh-owner-after", + }); + expect(startSpy).toHaveBeenCalledTimes(3); + expect(replacementKernel.execute).toHaveBeenCalledTimes(1); + expect(warmedKernel.shutdown).toHaveBeenCalledTimes(1); + expect(warmedKernel.execute).not.toHaveBeenCalled(); + await pythonExecutor.executePython("print('still alive after')", { cwd: unrelatedCwd, sessionId: "unrelated-after-session", @@ -242,7 +277,7 @@ describe("AgentSession python cleanup", () => { kernelOwnerId: "other-owner", }); - expect(startSpy).toHaveBeenCalledTimes(2); + expect(startSpy).toHaveBeenCalledTimes(3); expect(unrelatedKernel.execute).toHaveBeenCalledTimes(2); }); @@ -639,11 +674,18 @@ describe("AgentSession python cleanup", () => { const session = await createSession(tempDir, cwd, { extensions: [hookExtension] }); const execution = session.executePython("print('late after hook')"); await hookStarted.promise; - await session.dispose(); + let disposed = false; + const disposeSession = session.dispose().then(() => { + disposed = true; + }); + await Bun.sleep(0); + expect(disposed).toBe(false); releaseHook.resolve(); await expect(execution).rejects.toThrow("Python execution is unavailable while session disposal is in progress"); + await disposeSession; + expect(disposed).toBe(true); expect(executeSpy).not.toHaveBeenCalled(); - }); + }, 10000); it("rejects async user_python hook results after dispose begins", async () => { const { tempDir, cwd } = createTempProject(); @@ -687,12 +729,19 @@ describe("AgentSession python cleanup", () => { const session = await createSession(tempDir, cwd, { extensions: [hookExtension] }); const execution = session.executePython("print('late hook result')"); await hookStarted.promise; - await session.dispose(); + let disposed = false; + const disposeSession = session.dispose().then(() => { + disposed = true; + }); + await Bun.sleep(0); + expect(disposed).toBe(false); releaseHook.resolve(); await expect(execution).rejects.toThrow("Python execution is unavailable while session disposal is in progress"); + await disposeSession; + expect(disposed).toBe(true); expect(executeSpy).not.toHaveBeenCalled(); expect(session.messages.some(message => message.role === "pythonExecution")).toBe(false); - }); + }, 10000); it("rejects Python tool starts once dispose begins", async () => { const { tempDir, cwd } = createTempProject(); diff --git a/packages/coding-agent/test/core/python-executor-lifecycle.test.ts b/packages/coding-agent/test/core/python-executor-lifecycle.test.ts index 1897a7513..146880d78 100644 --- a/packages/coding-agent/test/core/python-executor-lifecycle.test.ts +++ b/packages/coding-agent/test/core/python-executor-lifecycle.test.ts @@ -100,25 +100,23 @@ describe("executePython lifecycle", () => { expect(kernelNext.execute).toHaveBeenCalledTimes(1); }); - it("retries a dead session restart after an unconfirmed shutdown", async () => { + it("restarts dead retained sessions even when shutdown confirmation is missing", async () => { const kernel = new FakeKernel(OK_RESULT); const kernelNext = new FakeKernel(OK_RESULT); kernel.alive = false; - kernel.shutdown.mockResolvedValueOnce({ confirmed: false }).mockResolvedValueOnce({ confirmed: true }); + kernel.shutdown.mockResolvedValueOnce({ confirmed: false }); vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); const startSpy = vi .spyOn(pythonKernel.PythonKernel, "start") .mockResolvedValueOnce(kernel as unknown as pythonKernel.PythonKernel) .mockResolvedValueOnce(kernelNext as unknown as pythonKernel.PythonKernel); - await expect( - executePython("1 + 1", { kernelMode: "session", sessionId: "retry-dead-session", cwd: getProjectDir() }), - ).rejects.toThrow("Failed to confirm crashed kernel shutdown before restart"); + await executePython("1 + 1", { kernelMode: "session", sessionId: "retry-dead-session", cwd: getProjectDir() }); await executePython("2 + 2", { kernelMode: "session", sessionId: "retry-dead-session", cwd: getProjectDir() }); expect(startSpy).toHaveBeenCalledTimes(2); - expect(kernel.shutdown).toHaveBeenCalledTimes(2); + expect(kernel.shutdown).toHaveBeenCalledTimes(1); expect(kernel.execute).toHaveBeenCalledTimes(0); - expect(kernelNext.execute).toHaveBeenCalledTimes(1); + expect(kernelNext.execute).toHaveBeenCalledTimes(2); }); }); 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 b844b81e4..20994d7f7 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 @@ -39,6 +39,12 @@ class FakeKernel { } } +async function flushMicrotasks(turns = 6): Promise { + for (let turn = 0; turn < turns; turn += 1) { + await Promise.resolve(); + } +} + afterEach(async () => { await disposeAllKernelSessions(); resetPreludeDocsCache(); @@ -240,6 +246,53 @@ describe("python executor owner cleanup", () => { expect(replacementKernel.shutdown).toHaveBeenCalledTimes(1); }); + it("waits with the owner-cleanup timeout when the last owner is removed from an already-disposing session", async () => { + vi.useFakeTimers(); + try { + const kernel = new FakeKernel(); + const shutdownConfirmation = Promise.withResolvers(); + kernel.shutdown = vi.fn(() => shutdownConfirmation.promise); + vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true }); + const startSpy = vi.spyOn(PythonKernel, "start").mockResolvedValue(kernel as unknown as PythonKernelInstance); + + await executePython("print('owner-a')", { + cwd: "/tmp/disposing-owner-cleanup-session", + sessionId: "disposing-owner-cleanup-session", + kernelMode: "session", + kernelOwnerId: "owner-a", + }); + + let globalCleanupResolved = false; + const globalCleanup = disposeAllKernelSessions().finally(() => { + globalCleanupResolved = true; + }); + await flushMicrotasks(); + expect(kernel.shutdown).toHaveBeenCalledTimes(1); + expect(globalCleanupResolved).toBe(false); + + let ownerCleanupResolved = false; + const ownerCleanup = disposeKernelSessionsByOwner("owner-a").finally(() => { + ownerCleanupResolved = true; + }); + await flushMicrotasks(); + expect(ownerCleanupResolved).toBe(false); + expect(kernel.shutdown).toHaveBeenCalledTimes(1); + + vi.advanceTimersByTime(2_000); + await ownerCleanup; + expect(ownerCleanupResolved).toBe(true); + expect(globalCleanupResolved).toBe(false); + expect(kernel.shutdown).toHaveBeenCalledTimes(1); + + shutdownConfirmation.resolve({ confirmed: true }); + await globalCleanup; + expect(globalCleanupResolved).toBe(true); + expect(startSpy).toHaveBeenCalledTimes(1); + } finally { + vi.useRealTimers(); + } + }); + it("returns a cancelled result when a dead session restart shutdown times out", async () => { const kernel = new FakeKernel(); kernel.alive = false; @@ -264,7 +317,7 @@ describe("python executor owner cleanup", () => { cwd: "/tmp/restart-timeout-session", sessionId: "restart-timeout-session", kernelMode: "session", - timeoutMs: 25, + timeoutMs: 100, }); expect(result.cancelled).toBe(true); @@ -272,7 +325,7 @@ describe("python executor owner cleanup", () => { expect(kernel.shutdown).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: expect.any(Number) })); expect(startSpy).toHaveBeenCalledTimes(1); }); - it("clears stuck tracked disposals during resource-exhaustion recovery", async () => { + it("keeps local owner-cleanup disposals counted during resource-exhaustion recovery", async () => { vi.useFakeTimers(); try { const staleKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel()]; @@ -319,15 +372,28 @@ describe("python executor owner cleanup", () => { expect(startSpy).toHaveBeenCalledTimes(5); expect(recoveredKernel.execute).toHaveBeenCalledTimes(1); - await executePython("print('later')", { + const blockedExecution = executePython("print('later')", { cwd: "/tmp/recovery-after-emfile-later", sessionId: "recovery-session-later", kernelMode: "session", - deadlineMs: Date.now() + 50, }); - expect(startSpy).toHaveBeenCalledTimes(6); + await flushMicrotasks(); + expect(startSpy).toHaveBeenCalledTimes(5); expect(recoveredKernel.shutdown).not.toHaveBeenCalled(); + expect(laterKernel.execute).not.toHaveBeenCalled(); + + staleShutdownDeferreds[0]!.resolve({ confirmed: true }); + await blockedExecution; + expect(startSpy).toHaveBeenCalledTimes(6); expect(laterKernel.execute).toHaveBeenCalledTimes(1); + + for (const deferred of staleShutdownDeferreds.slice(1)) { + deferred.resolve({ confirmed: true }); + } + await flushMicrotasks(); + await disposeAllKernelSessions(); + expect(recoveredKernel.shutdown).toHaveBeenCalledTimes(1); + expect(laterKernel.shutdown).toHaveBeenCalledTimes(1); } finally { vi.useRealTimers(); } @@ -358,10 +424,10 @@ describe("python executor owner cleanup", () => { } let ownerCleanupResolved = false; - const ownerCleanup = disposeKernelSessionsByOwner("owner-a").then(() => { + const ownerCleanup = disposeKernelSessionsByOwner("owner-a").finally(() => { ownerCleanupResolved = true; }); - await Promise.resolve(); + await flushMicrotasks(); for (const kernel of retainedKernels) { expect(kernel.shutdown).toHaveBeenCalledWith({ timeoutMs: 2_000 }); @@ -701,20 +767,16 @@ describe("python executor owner cleanup", () => { }); await globalExecutionStarted.promise; - const ownerCleanup = Promise.race([ - disposeKernelSessionsByOwner("owner-a").then(() => "disposed-owner" as const), - new Promise<"timeout">(resolve => setTimeout(() => resolve("timeout"), 50)), - ]); - await expect(ownerCleanup).resolves.toBe("disposed-owner"); + const ownerCleanup = disposeKernelSessionsByOwner("owner-a"); + await flushMicrotasks(); expect(ownerKernel.shutdown).toHaveBeenCalledTimes(1); expect(globalKernel.shutdown).not.toHaveBeenCalled(); + await ownerCleanup; - const globalCleanup = Promise.race([ - disposeAllKernelSessions().then(() => "disposed-all" as const), - new Promise<"timeout">(resolve => setTimeout(() => resolve("timeout"), 50)), - ]); - await expect(globalCleanup).resolves.toBe("disposed-all"); + const globalCleanup = disposeAllKernelSessions(); + await flushMicrotasks(); expect(globalKernel.shutdown).toHaveBeenCalledTimes(1); + await globalCleanup; }); it("attaches cached warmup sessions to newly provided owners", async () => { diff --git a/packages/coding-agent/test/core/python-executor-session.test.ts b/packages/coding-agent/test/core/python-executor-session.test.ts index df34a0571..d8a45eb4b 100644 --- a/packages/coding-agent/test/core/python-executor-session.test.ts +++ b/packages/coding-agent/test/core/python-executor-session.test.ts @@ -25,7 +25,7 @@ class FakeKernel { return this.alive; } - async shutdown(): Promise<{ confirmed: boolean }> { + async shutdown(): Promise { this.shutdownCalls += 1; this.alive = false; return { confirmed: true }; diff --git a/packages/coding-agent/test/core/python-executor.lifecycle.test.ts b/packages/coding-agent/test/core/python-executor.lifecycle.test.ts index e472e0591..c4db1d55e 100644 --- a/packages/coding-agent/test/core/python-executor.lifecycle.test.ts +++ b/packages/coding-agent/test/core/python-executor.lifecycle.test.ts @@ -3,6 +3,7 @@ import { disposeAllKernelSessions, executePython } from "@oh-my-pi/pi-coding-age import { type KernelExecuteOptions, type KernelExecuteResult, + type KernelShutdownResult, PythonKernel, } from "@oh-my-pi/pi-coding-agent/ipy/kernel"; @@ -34,7 +35,7 @@ class FakeKernel { return this.#result; } - async shutdown(): Promise<{ confirmed: boolean }> { + async shutdown(): Promise { this.shutdownCalls += 1; this.#alive = false; return { confirmed: true }; @@ -173,11 +174,11 @@ describe("executePython session lifecycle", () => { return kernels.shift() as unknown as PythonKernel; }; - kernelA.shutdown = async () => { + kernelA.shutdown = async (): Promise => { shutdownCount += 1; return { confirmed: true }; }; - kernelB.shutdown = async () => { + kernelB.shutdown = async (): Promise => { shutdownCount += 1; return { confirmed: true }; }; diff --git a/packages/coding-agent/test/core/python-kernel-session.test.ts b/packages/coding-agent/test/core/python-kernel-session.test.ts index 4f17d12b3..5d80e7007 100644 --- a/packages/coding-agent/test/core/python-kernel-session.test.ts +++ b/packages/coding-agent/test/core/python-kernel-session.test.ts @@ -1,6 +1,10 @@ import { afterEach, beforeEach, describe, expect, it } from "bun:test"; import { disposeAllKernelSessions, executePython } from "@oh-my-pi/pi-coding-agent/ipy/executor"; -import type { KernelExecuteOptions, KernelExecuteResult } from "@oh-my-pi/pi-coding-agent/ipy/kernel"; +import type { + KernelExecuteOptions, + KernelExecuteResult, + KernelShutdownResult, +} from "@oh-my-pi/pi-coding-agent/ipy/kernel"; import { PythonKernel } from "@oh-my-pi/pi-coding-agent/ipy/kernel"; import { TempDir } from "@oh-my-pi/pi-utils"; @@ -20,7 +24,7 @@ class FakeKernel { return { status: "ok", cancelled: false, timedOut: false, stdinRequested: false }; } - async shutdown(): Promise<{ confirmed: boolean }> { + async shutdown(): Promise { this.shutdownCalls += 1; this.alive = false; return { confirmed: true }; diff --git a/packages/coding-agent/test/core/python-kernel.lifecycle.test.ts b/packages/coding-agent/test/core/python-kernel.lifecycle.test.ts index 45cf6493b..59cdb0054 100644 --- a/packages/coding-agent/test/core/python-kernel.lifecycle.test.ts +++ b/packages/coding-agent/test/core/python-kernel.lifecycle.test.ts @@ -77,6 +77,21 @@ const createFakeProcess = (): Subprocess => { return { pid: 999999, exited } as Subprocess; }; +const expectResolvesWithin = async (promise: Promise, timeoutMs: number, message: string): Promise => { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(message)), timeoutMs); + timer.unref?.(); + }), + ]); + } finally { + if (timer !== undefined) clearTimeout(timer); + } +}; + describe("PythonKernel gateway lifecycle", () => { const originalWebSocket = globalThis.WebSocket; const originalGatewayUrl = Bun.env.PI_PYTHON_GATEWAY_URL; @@ -339,39 +354,37 @@ describe("PythonKernel gateway lifecycle", () => { ); }); - it("treats a retry against an already-missing kernel as confirmed shutdown", async () => { + it("treats initial 404 and 410 shutdown responses as confirmed", async () => { using _runtime = stubKernelRuntime(); vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({ url: "http://127.0.0.1:9999", isShared: true, }); - let deleteCalls = 0; - using _hook = hookFetch((input, init) => { - const url = String(input); - env.fetchCalls.push({ url, init }); - if (url.endsWith("/api/kernels") && init?.method === "POST") { - return createResponse({ ok: true, json: { id: "kernel-retry-delete" } }) as unknown as Response; - } - if (url.endsWith("/api/kernels/kernel-retry-delete") && init?.method === "DELETE") { - deleteCalls += 1; - if (deleteCalls === 1) { - return createResponse({ ok: false, status: 503, text: "not yet" }) as unknown as Response; + for (const status of [404, 410]) { + let deleteCalls = 0; + using _hook = hookFetch((input, init) => { + const url = String(input); + env.fetchCalls.push({ url, init }); + if (url.endsWith("/api/kernels") && init?.method === "POST") { + return createResponse({ ok: true, json: { id: `kernel-missing-${status}` } }) as unknown as Response; } - return createResponse({ ok: false, status: 404, text: "gone" }) as unknown as Response; - } - return createResponse({ ok: true }) as unknown as Response; - }); + if (url.endsWith(`/api/kernels/kernel-missing-${status}`) && init?.method === "DELETE") { + deleteCalls += 1; + return createResponse({ ok: false, status, text: "gone" }) as unknown as Response; + } + return createResponse({ ok: true }) as unknown as Response; + }); - const kernel = await PythonKernel.start({ cwd: tempDir.path() }); + const kernel = await PythonKernel.start({ cwd: tempDir.path() }); - await expect(kernel.shutdown()).resolves.toEqual({ confirmed: false }); - expect(kernel.isAlive()).toBe(false); - expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED); - - await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true }); - await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true }); - expect(deleteCalls).toBe(2); + await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true }); + expect(deleteCalls).toBe(1); + expect(kernel.isAlive()).toBe(false); + expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED); + await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true }); + expect(deleteCalls).toBe(1); + } }); it("returns unconfirmed when shutdown times out and can confirm on retry", async () => { @@ -394,17 +407,18 @@ describe("PythonKernel gateway lifecycle", () => { if (deleteCalls === 1) { firstDeleteStarted.resolve(); return new Promise((_, reject) => { - const waitForAbort = () => { - if (init.signal?.aborted) { - firstDeleteAborted.resolve(); - const reason = init.signal.reason; - reject(reason instanceof Error ? reason : new Error("Python kernel shutdown timed out")); - return; - } - const poll = setTimeout(waitForAbort, 5); - poll.unref?.(); + const abortSignal = init.signal; + if (!abortSignal) return; + const rejectOnAbort = () => { + firstDeleteAborted.resolve(); + const reason = abortSignal.reason; + reject(reason instanceof Error ? reason : new Error("Python kernel shutdown timed out")); }; - waitForAbort(); + if (abortSignal.aborted) { + rejectOnAbort(); + return; + } + abortSignal.addEventListener("abort", rejectOnAbort, { once: true }); }); } return createResponse({ ok: false, status: 404, text: "gone" }) as unknown as Response; @@ -413,18 +427,17 @@ describe("PythonKernel gateway lifecycle", () => { }); const kernel = await PythonKernel.start({ cwd: tempDir.path() }); const shutdownPromise = kernel.shutdown({ timeoutMs: 25 }); - await firstDeleteStarted.promise; - const pending = Symbol("pending"); - const settled = await Promise.race([ - shutdownPromise, - new Promise(resolve => { - const timer = setTimeout(() => resolve(pending), 250); - timer.unref?.(); - }), - ]); - await firstDeleteAborted.promise; - expect(settled).not.toBe(pending); - expect(settled).toEqual({ confirmed: false }); + await expectResolvesWithin(firstDeleteStarted.promise, 250, "kernel shutdown never issued a delete request"); + await expectResolvesWithin( + firstDeleteAborted.promise, + 500, + "timed out waiting for the first delete request to abort", + ); + await expect( + expectResolvesWithin(shutdownPromise, 500, "kernel shutdown did not settle after timing out"), + ).resolves.toEqual({ + confirmed: false, + }); expect(kernel.isAlive()).toBe(false); expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED); await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });