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
This commit is contained in:
@@ -540,8 +540,8 @@ export async function disposeAllKernelSessions(): Promise<void> {
|
||||
|
||||
export async function disposeKernelSessionsByOwner(ownerId: string): Promise<void> {
|
||||
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<void> {
|
||||
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<string, string> | undefined = options.sessionFile
|
||||
? { PI_SESSION_FILE: options.sessionFile }
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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<PythonResult> => {
|
||||
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}`;
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -39,6 +39,12 @@ class FakeKernel {
|
||||
}
|
||||
}
|
||||
|
||||
async function flushMicrotasks(turns = 6): Promise<void> {
|
||||
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<KernelShutdownResult>();
|
||||
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 () => {
|
||||
|
||||
@@ -25,7 +25,7 @@ class FakeKernel {
|
||||
return this.alive;
|
||||
}
|
||||
|
||||
async shutdown(): Promise<{ confirmed: boolean }> {
|
||||
async shutdown(): Promise<pythonKernel.KernelShutdownResult> {
|
||||
this.shutdownCalls += 1;
|
||||
this.alive = false;
|
||||
return { confirmed: true };
|
||||
|
||||
@@ -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<KernelShutdownResult> {
|
||||
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<KernelShutdownResult> => {
|
||||
shutdownCount += 1;
|
||||
return { confirmed: true };
|
||||
};
|
||||
kernelB.shutdown = async () => {
|
||||
kernelB.shutdown = async (): Promise<KernelShutdownResult> => {
|
||||
shutdownCount += 1;
|
||||
return { confirmed: true };
|
||||
};
|
||||
|
||||
@@ -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<KernelShutdownResult> {
|
||||
this.shutdownCalls += 1;
|
||||
this.alive = false;
|
||||
return { confirmed: true };
|
||||
|
||||
@@ -77,6 +77,21 @@ const createFakeProcess = (): Subprocess => {
|
||||
return { pid: 999999, exited } as Subprocess;
|
||||
};
|
||||
|
||||
const expectResolvesWithin = async <T>(promise: Promise<T>, timeoutMs: number, message: string): Promise<T> => {
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
return await Promise.race([
|
||||
promise,
|
||||
new Promise<T>((_, 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<Response>((_, 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<typeof pending>(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 });
|
||||
|
||||
Reference in New Issue
Block a user