From cd999364e41564b5ac72c44bf994015c7f1ae3b4 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 00:30:47 +0530 Subject: [PATCH 1/7] refactor(claude): keep streaming sessions warm --- src/local-agent-adapters.ts | 85 ++--------- src/local-agent-claude/runtime.ts | 229 ++++++++++++++++++++++++++++++ src/local-agent-manager.ts | 1 + src/local-agent-runtime.ts | 1 + 4 files changed, 241 insertions(+), 75 deletions(-) create mode 100644 src/local-agent-claude/runtime.ts diff --git a/src/local-agent-adapters.ts b/src/local-agent-adapters.ts index b015f6a9a..dabd2ebd6 100644 --- a/src/local-agent-adapters.ts +++ b/src/local-agent-adapters.ts @@ -1,12 +1,14 @@ -import { spawnSync } from "node:child_process"; -import { resolve } from "node:path"; -import type { EffortLevel } from "@anthropic-ai/claude-agent-sdk"; import type { LocalAgentProvider } from "./local-agent-profiles.js"; import type { LocalAgentRunInput, LocalAgentRunResult } from "./local-agent-runtime.js"; import type { HarnessDriver, HarnessRuntime } from "./local-agent-runtime-pool.js"; import { createCodexHarnessDriver } from "./local-agent-codex/runtime.js"; import { createAcpHarnessDriver } from "./local-agent-acp/runtime.js"; import { createPiHarnessDriver } from "./local-agent-pi/runtime.js"; +import { + claudeCommandEnvironment, + createClaudeHarnessDriver, +} from "./local-agent-claude/runtime.js"; +export { claudeCommandEnvironment } from "./local-agent-claude/runtime.js"; export { resolveAcpModelConfigUpdate, resolveAcpThinkingConfigUpdate, @@ -62,86 +64,19 @@ class ClaudeLocalAgentAdapter implements LocalAgentAdapter { readonly provider = "claude" as const; async run(input: LocalAgentRunInput): Promise { - const { query } = await import("@anthropic-ai/claude-agent-sdk"); - const claudeExecutable = process.env.CLAUDE_COMMAND ?? resolveExecutable("claude"); - const messages = query({ - prompt: input.prompt, - options: { - cwd: input.workspace, - model: input.model, - ...(input.thinking ? { thinking: { type: "adaptive" } as const, effort: input.thinking as EffortLevel } : {}), - resume: input.providerSessionId, - permissionMode: "bypassPermissions", - allowDangerouslySkipPermissions: true, - env: claudeCommandEnvironment(process.env), - ...(claudeExecutable ? { pathToClaudeCodeExecutable: claudeExecutable } : {}), - }, - }); - - let providerSessionId = input.providerSessionId ?? null; - let finalResponse = ""; - const items: unknown[] = []; - for await (const message of messages) { - items.push(message); - const record = message as Record; - if (typeof record.session_id === "string") providerSessionId = record.session_id; - if (record.type === "result" && typeof record.result === "string") { - const resultError = claudeResultError(record); - if (resultError) throw new Error(resultError); - finalResponse = record.result; - } + const runtime = await createClaudeHarnessDriver().createRuntime(input); + try { + return await runtime.run({ ...input, agentId: input.agentId ?? "direct" }); + } finally { + await runtime.close(); } - - finalResponse = requireFinalResponse("Claude", finalResponse); - return { - provider: this.provider, - providerSessionId, - finalResponse, - items, - }; } } -function claudeResultError(record: Record): string | undefined { - const subtype = typeof record.subtype === "string" ? record.subtype : undefined; - const isError = record.is_error === true || subtype?.startsWith("error"); - if (!isError) return undefined; - const message = - directString(record.error) ?? - directString(record.message) ?? - directString(record.result) ?? - subtype ?? - "Claude returned an error result."; - return `Claude returned an error result: ${message}`; -} - function directString(value: unknown): string | undefined { return typeof value === "string" && value.trim() ? value.trim() : undefined; } -function resolveExecutable(command: string): string | undefined { - const result = spawnSync(process.platform === "win32" ? "where.exe" : "command", [ - ...(process.platform === "win32" ? [command] : ["-v", command]), - ], { - encoding: "utf8", - shell: process.platform !== "win32", - }); - const executable = result.stdout?.split(/\r?\n/).find((line) => line.trim()); - return executable?.trim() || undefined; -} - -export function claudeCommandEnvironment(env: NodeJS.ProcessEnv): NodeJS.ProcessEnv { - const next = { ...env }; - for (const key of [ - "CLAUDECODE", - "CLAUDE_CODE_ENTRYPOINT", - "CLAUDE_CODE_SSE_PORT", - "CLAUDE_AGENT_SDK_VERSION", - ]) { - delete next[key]; - } - return next; -} class OpencodeLocalAgentAdapter implements LocalAgentAdapter { readonly provider = "opencode" as const; diff --git a/src/local-agent-claude/runtime.ts b/src/local-agent-claude/runtime.ts new file mode 100644 index 000000000..9d6dd476b --- /dev/null +++ b/src/local-agent-claude/runtime.ts @@ -0,0 +1,229 @@ +import type { + EffortLevel, + Query, + SDKMessage, + SDKUserMessage, +} from "@anthropic-ai/claude-agent-sdk"; +import type { LocalAgentRunInput, LocalAgentRunResult } from "../local-agent-runtime.js"; +import type { HarnessDriver, HarnessRuntime } from "../local-agent-runtime-pool.js"; +import { resolveLocalAgentExecutable } from "../local-agent-path.js"; + +export interface ClaudeQueryLike { + next(): Promise>; + setModel(model?: string): Promise; + close(): void; +} + +export type ClaudeQueryFactory = ( + input: LocalAgentRunInput, + prompts: AsyncIterable, + resume: string | undefined, +) => ClaudeQueryLike; + +interface LiveClaudeQuery { + query: ClaudeQueryLike; + prompts: AsyncPromptQueue; + effort: string | undefined; + model: string | undefined; +} + +/** + * Keeps one streaming Claude SDK query warm for a logical DevSpace agent. + * If session-start-only configuration changes, only that query is replaced; + * the durable Claude session id is used to resume the conversation. + */ +export class ClaudeWarmRuntime implements HarnessRuntime { + private live: LiveClaudeQuery | undefined; + private providerSessionId: string | undefined; + private closed = false; + + constructor(private readonly createQuery: ClaudeQueryFactory) {} + + async run(input: LocalAgentRunInput): Promise { + if (this.closed) throw new Error("Claude runtime is closed."); + this.providerSessionId = input.providerSessionId ?? this.providerSessionId; + + const live = this.ensureQuery(input); + if (input.model !== undefined && input.model !== live.model) { + await live.query.setModel(input.model); + live.model = input.model; + } + + live.prompts.push(userMessage(input.prompt, this.providerSessionId)); + const items: unknown[] = []; + try { + for (;;) { + const next = await live.query.next(); + if (next.done) { + this.resetQuery(); + throw new Error("Claude session ended before returning a result."); + } + const message = next.value; + items.push(message); + if (typeof message.session_id === "string") { + this.providerSessionId = message.session_id; + } + if (message.type !== "result") continue; + const error = claudeResultError(message as unknown as Record); + if (error) throw new Error(error); + if (typeof message.result !== "string" || !message.result.trim()) { + throw new Error("Claude completed without a final response."); + } + return { + provider: "claude", + providerSessionId: this.providerSessionId ?? null, + finalResponse: message.result, + items, + }; + } + } catch (error) { + // A failed iterator is not safe to reuse. The next explicit turn will + // recreate the query and resume from the persisted provider session id. + this.resetQuery(); + throw error; + } + } + + isUsable(): boolean { + return !this.closed; + } + + async close(): Promise { + if (this.closed) return; + this.closed = true; + this.resetQuery(); + } + + private ensureQuery(input: LocalAgentRunInput): LiveClaudeQuery { + if (this.live && this.live.effort === input.thinking) return this.live; + this.resetQuery(); + const prompts = new AsyncPromptQueue(); + const query = this.createQuery(input, prompts, this.providerSessionId); + this.live = { + query, + prompts, + effort: input.thinking, + model: input.model, + }; + return this.live; + } + + private resetQuery(): void { + const live = this.live; + this.live = undefined; + if (!live) return; + live.prompts.close(); + live.query.close(); + } +} + +export function createClaudeHarnessDriver( + env: NodeJS.ProcessEnv = process.env, +): HarnessDriver { + return { + provider: "claude", + runtimeKey: (input) => { + if (!input.agentId) { + throw new Error("Claude pooled runtime requires a DevSpace agent id."); + } + return input.agentId; + }, + createRuntime: async () => { + const { query } = await import("@anthropic-ai/claude-agent-sdk"); + return new ClaudeWarmRuntime((input, prompts, resume) => { + const claudeExecutable = env.CLAUDE_COMMAND?.trim() + || resolveLocalAgentExecutable("claude", env); + return query({ + prompt: prompts, + options: { + cwd: input.workspace, + model: input.model, + ...(input.thinking + ? { + thinking: { type: "adaptive" } as const, + effort: input.thinking as EffortLevel, + } + : {}), + resume, + permissionMode: "bypassPermissions", + allowDangerouslySkipPermissions: true, + env: claudeCommandEnvironment(env), + ...(claudeExecutable ? { pathToClaudeCodeExecutable: claudeExecutable } : {}), + }, + }) as Query; + }); + }, + }; +} + +export function claudeCommandEnvironment(env: NodeJS.ProcessEnv): NodeJS.ProcessEnv { + const next = { ...env }; + for (const key of [ + "CLAUDECODE", + "CLAUDE_CODE_ENTRYPOINT", + "CLAUDE_CODE_SSE_PORT", + "CLAUDE_AGENT_SDK_VERSION", + ]) { + delete next[key]; + } + return next; +} + +function userMessage(prompt: string, sessionId: string | undefined): SDKUserMessage { + return { + type: "user", + message: { role: "user", content: prompt }, + parent_tool_use_id: null, + ...(sessionId ? { session_id: sessionId } : {}), + }; +} + +function claudeResultError(record: Record): string | undefined { + const subtype = typeof record.subtype === "string" ? record.subtype : undefined; + const isError = record.is_error === true || subtype?.startsWith("error"); + if (!isError) return undefined; + const message = + directString(record.error) ?? + directString(record.message) ?? + directString(record.result) ?? + subtype ?? + "Claude returned an error result."; + return `Claude returned an error result: ${message}`; +} + +function directString(value: unknown): string | undefined { + return typeof value === "string" && value.trim() ? value.trim() : undefined; +} + +class AsyncPromptQueue implements AsyncIterable { + private readonly values: T[] = []; + private readonly waiters: Array<(result: IteratorResult) => void> = []; + private closed = false; + + push(value: T): void { + if (this.closed) throw new Error("Claude prompt stream is closed."); + const waiter = this.waiters.shift(); + if (waiter) { + waiter({ done: false, value }); + return; + } + this.values.push(value); + } + + close(): void { + if (this.closed) return; + this.closed = true; + for (const waiter of this.waiters.splice(0)) waiter({ done: true, value: undefined }); + } + + [Symbol.asyncIterator](): AsyncIterator { + return { + next: async () => { + const value = this.values.shift(); + if (value !== undefined) return { done: false, value }; + if (this.closed) return { done: true, value: undefined }; + return new Promise>((resolve) => this.waiters.push(resolve)); + }, + }; + } +} diff --git a/src/local-agent-manager.ts b/src/local-agent-manager.ts index d0a5f5010..7ce5ff839 100644 --- a/src/local-agent-manager.ts +++ b/src/local-agent-manager.ts @@ -186,6 +186,7 @@ export class LocalAgentManager { const profile = profiles.find((candidate) => candidate.name === record.profileName); const prompt = profile ? profilePrompt(profile, turn.prompt) : rawProviderPrompt(record, turn.prompt); const result = await this.runProvider(record.provider, { + agentId: record.id, prompt, workspace: workspaceRoot, providerSessionId: record.providerSessionId, diff --git a/src/local-agent-runtime.ts b/src/local-agent-runtime.ts index 9548326d9..7a3670b12 100644 --- a/src/local-agent-runtime.ts +++ b/src/local-agent-runtime.ts @@ -1,6 +1,7 @@ export type LocalAgentWriteMode = "read_only" | "allowed" | "full_access"; export interface LocalAgentRunInput { + agentId?: string; prompt: string; workspace: string; providerSessionId?: string; From 06c68f59d61f3b4e57e5f7f414269bf24aa2f40f Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 00:30:47 +0530 Subject: [PATCH 2/7] perf(claude): pool per-agent runtimes --- src/local-agent-runtime-registry.ts | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/src/local-agent-runtime-registry.ts b/src/local-agent-runtime-registry.ts index e2b7d50d4..e472c1dd3 100644 --- a/src/local-agent-runtime-registry.ts +++ b/src/local-agent-runtime-registry.ts @@ -4,6 +4,7 @@ import { createLocalAgentAdapter, createOpencodeHarnessDriver } from "./local-ag import { createCodexHarnessDriver } from "./local-agent-codex/runtime.js"; import { createAcpHarnessDriver } from "./local-agent-acp/runtime.js"; import { createPiHarnessDriver } from "./local-agent-pi/runtime.js"; +import { createClaudeHarnessDriver } from "./local-agent-claude/runtime.js"; import { HarnessRuntimePool, type HarnessDriver, @@ -16,6 +17,7 @@ interface LocalAgentRuntimeRegistryOptions { cursorDriver?: HarnessDriver; copilotDriver?: HarnessDriver; piDriver?: HarnessDriver; + claudeDriver?: HarnessDriver; } /** @@ -29,6 +31,7 @@ export class LocalAgentRuntimeRegistry { private readonly cursorDriver: HarnessDriver; private readonly copilotDriver: HarnessDriver; private readonly piDriver: HarnessDriver; + private readonly claudeDriver: HarnessDriver; constructor(options: LocalAgentRuntimeRegistryOptions = {}) { this.pool = options.pool ?? new HarnessRuntimePool(); @@ -37,6 +40,7 @@ export class LocalAgentRuntimeRegistry { this.cursorDriver = options.cursorDriver ?? createAcpHarnessDriver("cursor", ["cursor-agent", "acp"]); this.copilotDriver = options.copilotDriver ?? createAcpHarnessDriver("copilot", ["copilot", "--acp"]); this.piDriver = options.piDriver ?? createPiHarnessDriver(); + this.claudeDriver = options.claudeDriver ?? createClaudeHarnessDriver(); } async run( @@ -58,6 +62,9 @@ export class LocalAgentRuntimeRegistry { if (provider === "pi") { return this.pool.run(this.piDriver, input); } + if (provider === "claude") { + return this.pool.run(this.claudeDriver, input); + } return createLocalAgentAdapter(provider).run(input); } From 46121718d7f90f7a6d4137fed2571052b7fe01ae Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 00:30:47 +0530 Subject: [PATCH 3/7] test(claude): cover warm session reuse --- docs/agent-profile-schema.md | 2 +- package.json | 2 +- src/local-agent-claude/runtime.test.ts | 102 +++++++++++++++++++++++++ src/local-agent-manager.test.ts | 1 + 4 files changed, 105 insertions(+), 2 deletions(-) create mode 100644 src/local-agent-claude/runtime.test.ts diff --git a/docs/agent-profile-schema.md b/docs/agent-profile-schema.md index c956f6270..0fac62840 100644 --- a/docs/agent-profile-schema.md +++ b/docs/agent-profile-schema.md @@ -72,7 +72,7 @@ Unsupported or custom providers are rejected. DevSpace maps providers to their native integration: - `codex`: the user's Codex CLI through `codex app-server` -- `claude`: Claude Code SDK +- `claude`: warm Claude Code SDK streaming session - `opencode`: OpenCode SDK - `pi`: embedded Pi `AgentSession` runtime - `cursor`: ACP diff --git a/package.json b/package.json index 444a31950..429fe06da 100644 --- a/package.json +++ b/package.json @@ -28,7 +28,7 @@ "dev": "node scripts/dev-server.mjs", "postinstall": "node scripts/fix-node-pty-permissions.mjs", "start": "node dist/cli.js serve", - "test": "tsx src/config.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-codex/runtime.test.ts && tsx src/local-agent-pi/runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/local-agent-manager.test.ts && tsx src/local-agent-control.test.ts && tsx src/local-agent-runtime-pool.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts", + "test": "tsx src/config.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-codex/runtime.test.ts && tsx src/local-agent-pi/runtime.test.ts && tsx src/local-agent-claude/runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/local-agent-manager.test.ts && tsx src/local-agent-control.test.ts && tsx src/local-agent-runtime-pool.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit" }, "keywords": [], diff --git a/src/local-agent-claude/runtime.test.ts b/src/local-agent-claude/runtime.test.ts new file mode 100644 index 000000000..7084fb280 --- /dev/null +++ b/src/local-agent-claude/runtime.test.ts @@ -0,0 +1,102 @@ +import assert from "node:assert/strict"; +import type { SDKMessage, SDKUserMessage } from "@anthropic-ai/claude-agent-sdk"; +import { + ClaudeWarmRuntime, + type ClaudeQueryLike, +} from "./runtime.js"; + +class FakeClaudeQuery implements ClaudeQueryLike { + readonly models: Array = []; + closed = false; + private readonly prompts: AsyncIterator; + private turn = 0; + + constructor( + prompts: AsyncIterable, + readonly sessionId: string, + ) { + this.prompts = prompts[Symbol.asyncIterator](); + } + + async next(): Promise> { + const prompt = await this.prompts.next(); + if (prompt.done) return { done: true, value: undefined }; + this.turn += 1; + const content = prompt.value.message.content; + const text = typeof content === "string" ? content : JSON.stringify(content); + return { + done: false, + value: { + type: "result", + subtype: "success", + is_error: false, + result: `response:${text}`, + session_id: this.sessionId, + } as SDKMessage, + }; + } + + async setModel(model?: string): Promise { + this.models.push(model); + } + + close(): void { + this.closed = true; + } +} + +const queries: FakeClaudeQuery[] = []; +const resumes: Array = []; +const efforts: Array = []; +const runtime = new ClaudeWarmRuntime((input, prompts, resume) => { + resumes.push(resume); + efforts.push(input.thinking); + const query = new FakeClaudeQuery(prompts, resume ?? `claude_${queries.length + 1}`); + queries.push(query); + return query; +}); + +try { + const first = await runtime.run({ + agentId: "agt_claude", + workspace: "/tmp/project", + prompt: "first", + model: "sonnet", + thinking: "high", + }); + const second = await runtime.run({ + agentId: "agt_claude", + workspace: "/tmp/project", + prompt: "second", + providerSessionId: first.providerSessionId ?? undefined, + model: "opus", + thinking: "high", + }); + + assert.equal(queries.length, 1, "same agent and effort should reuse one live Claude query"); + assert.equal(first.providerSessionId, "claude_1"); + assert.equal(second.providerSessionId, "claude_1"); + assert.equal(first.finalResponse, "response:first"); + assert.equal(second.finalResponse, "response:second"); + assert.deepEqual(queries[0]?.models, ["opus"], "model changes should use the live query control channel"); + + const changedEffort = await runtime.run({ + agentId: "agt_claude", + workspace: "/tmp/project", + prompt: "third", + providerSessionId: second.providerSessionId ?? undefined, + model: "opus", + thinking: "xhigh", + }); + + assert.equal(queries.length, 2, "session-start-only effort changes should replace only the live query"); + assert.equal(queries[0]?.closed, true); + assert.deepEqual(resumes, [undefined, "claude_1"]); + assert.deepEqual(efforts, ["high", "xhigh"]); + assert.equal(changedEffort.providerSessionId, "claude_1"); + assert.equal(changedEffort.finalResponse, "response:third"); +} finally { + await runtime.close(); +} + +assert.equal(queries.at(-1)?.closed, true); diff --git a/src/local-agent-manager.test.ts b/src/local-agent-manager.test.ts index 16cd972cf..a3ab5d1e0 100644 --- a/src/local-agent-manager.test.ts +++ b/src/local-agent-manager.test.ts @@ -80,6 +80,7 @@ try { await waitFor(() => calls.length === 2); assert.equal(calls[0]?.provider, "codex"); + assert.equal(calls[0]?.input.agentId, first.id); assert.equal(calls[0]?.input.prompt, "first"); assert.equal(calls[0]?.input.model, "gpt-test"); assert.equal(calls[1]?.input.prompt, "second"); From 7491476a06b17cbdaab14108192e53db65d532a4 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 00:35:03 +0530 Subject: [PATCH 4/7] perf(agents): reclaim idle provider sessions --- src/local-agent-acp/runtime.ts | 91 ++++++++++++++++++--------- src/local-agent-codex/runtime.test.ts | 9 +++ src/local-agent-codex/runtime.ts | 21 +++++++ src/local-agent-pi/runtime.test.ts | 12 ++++ src/local-agent-pi/runtime.ts | 36 ++++++++++- src/local-agent-runtime-pool.test.ts | 7 +++ src/local-agent-runtime-pool.ts | 2 + 7 files changed, 146 insertions(+), 32 deletions(-) diff --git a/src/local-agent-acp/runtime.ts b/src/local-agent-acp/runtime.ts index 9157d8fd9..2cf7a3213 100644 --- a/src/local-agent-acp/runtime.ts +++ b/src/local-agent-acp/runtime.ts @@ -22,6 +22,7 @@ import { const STDERR_LIMIT = 32_000; const SHUTDOWN_GRACE_MS = 2_000; +const ACP_SESSION_IDLE_MS = 5 * 60 * 1_000; type AcpProvider = "cursor" | "copilot"; @@ -29,8 +30,13 @@ interface PendingAcpTurn { textParts: string[]; } +interface AcpSessionEntry extends AcpSessionConfigState { + active: boolean; + lastUsedAt: number; +} + export class AcpHarnessRuntime implements HarnessRuntime { - private readonly sessions = new Map(); + private readonly sessions = new Map(); private readonly pendingTurns = new Map(); private closed = false; @@ -49,33 +55,38 @@ export class AcpHarnessRuntime implements HarnessRuntime { if (this.pendingTurns.has(session.sessionId)) { throw new Error(`${this.provider} ACP session ${session.sessionId} already has a turn in progress.`); } - - if (input.model) { - await this.connection.agent.request("session/set_config_option", { - ...resolveAcpModelConfigUpdate(session, input.model, this.provider), - }); - } - if (input.thinking) { - await this.connection.agent.request("session/set_config_option", { - ...resolveAcpThinkingConfigUpdate(session, input.thinking, this.provider), - }); - } - - const pending: PendingAcpTurn = { textParts: [] }; - this.pendingTurns.set(session.sessionId, pending); + session.active = true; try { - await this.connection.agent.request("session/prompt", { - sessionId: session.sessionId, - prompt: [{ type: "text", text: input.prompt }], - }); - return { - provider: this.provider, - providerSessionId: session.sessionId, - finalResponse: pending.textParts.join("").trim(), - items: [], - }; + if (input.model) { + await this.connection.agent.request("session/set_config_option", { + ...resolveAcpModelConfigUpdate(session, input.model, this.provider), + }); + } + if (input.thinking) { + await this.connection.agent.request("session/set_config_option", { + ...resolveAcpThinkingConfigUpdate(session, input.thinking, this.provider), + }); + } + + const pending: PendingAcpTurn = { textParts: [] }; + this.pendingTurns.set(session.sessionId, pending); + try { + await this.connection.agent.request("session/prompt", { + sessionId: session.sessionId, + prompt: [{ type: "text", text: input.prompt }], + }); + return { + provider: this.provider, + providerSessionId: session.sessionId, + finalResponse: pending.textParts.join("").trim(), + items: [], + }; + } finally { + this.pendingTurns.delete(session.sessionId); + } } finally { - this.pendingTurns.delete(session.sessionId); + session.active = false; + session.lastUsedAt = Date.now(); } } catch (error) { const detail = this.stderr().trim(); @@ -100,6 +111,19 @@ export class AcpHarnessRuntime implements HarnessRuntime { && this.child.signalCode === null; } + async reapIdleSessions(now: number): Promise { + if (!this.initializeResponse.agentCapabilities?.sessionCapabilities?.close) return; + for (const [sessionId, session] of this.sessions) { + if (session.active || now - session.lastUsedAt < ACP_SESSION_IDLE_MS) continue; + try { + await this.connection.agent.request("session/close", { sessionId }); + this.sessions.delete(sessionId); + } catch { + // The connection remains useful; retry this session on a later reap. + } + } + } + async close(): Promise { if (this.closed) return; this.closed = true; @@ -118,22 +142,25 @@ export class AcpHarnessRuntime implements HarnessRuntime { } } - private async ensureSession(input: LocalAgentRunInput): Promise { + private async ensureSession(input: LocalAgentRunInput): Promise { if (input.providerSessionId) { const existing = this.sessions.get(input.providerSessionId); if (existing) return existing; const resumed = await this.resumeSession(input.providerSessionId, input.workspace); - this.sessions.set(input.providerSessionId, resumed); - return resumed; + const entry = withActivity(resumed); + this.sessions.set(input.providerSessionId, entry); + return entry; } const response = await this.connection.agent.request("session/new", { cwd: input.workspace, mcpServers: [], }) as NewSessionResponse; - const session: AcpSessionConfigState = { + const session: AcpSessionEntry = { sessionId: response.sessionId, configOptions: response.configOptions, + active: false, + lastUsedAt: Date.now(), }; this.sessions.set(session.sessionId, session); return session; @@ -163,6 +190,10 @@ export class AcpHarnessRuntime implements HarnessRuntime { } } +function withActivity(session: AcpSessionConfigState): AcpSessionEntry { + return { ...session, active: false, lastUsedAt: Date.now() }; +} + export function createAcpHarnessDriver( provider: AcpProvider, command: readonly [string, ...string[]], diff --git a/src/local-agent-codex/runtime.test.ts b/src/local-agent-codex/runtime.test.ts index d58c83c6d..d891b7265 100644 --- a/src/local-agent-codex/runtime.test.ts +++ b/src/local-agent-codex/runtime.test.ts @@ -47,6 +47,7 @@ class FakeCodexConnection implements CodexAppServerConnection { }); return { turn: { id: turnId } }; } + if (method === "thread/unsubscribe") return {}; throw new Error(`Unexpected request: ${method}`); } @@ -136,6 +137,14 @@ try { const finalTurn = connection.requests.filter((request) => request.method === "turn/start").at(-1); assert.equal(asRecord(finalTurn?.params)?.effort, "xhigh"); + + await runtime.reapIdleSessions(Date.now() + 5 * 60 * 1_000 + 1); + const unsubscribes = connection.requests.filter((request) => request.method === "thread/unsubscribe"); + assert.deepEqual( + unsubscribes.map((request) => asRecord(request.params)?.threadId).sort(), + ["thread_1", "thread_2"], + "idle Codex threads should be released while the shared App Server stays alive", + ); } finally { await runtime.close(); } diff --git a/src/local-agent-codex/runtime.ts b/src/local-agent-codex/runtime.ts index 6f4beb687..64dd1501f 100644 --- a/src/local-agent-codex/runtime.ts +++ b/src/local-agent-codex/runtime.ts @@ -7,6 +7,7 @@ import { import { resolveCodexCommand } from "./command.js"; type CodexSandboxMode = "read-only" | "workspace-write" | "danger-full-access"; +const CODEX_THREAD_IDLE_MS = 5 * 60 * 1_000; interface PendingTurn { turnId?: string; @@ -29,6 +30,7 @@ interface TurnCompletion { */ export class CodexAppServerRuntime implements HarnessRuntime { private readonly pendingTurns = new Map(); + private readonly threadActivity = new Map(); private readonly unsubscribe: () => void; private readonly unsubscribeClose: () => void; private closed = false; @@ -52,6 +54,9 @@ export class CodexAppServerRuntime implements HarnessRuntime { throw new Error(`Codex thread ${threadId} already has a turn in progress.`); } + const activity = this.threadActivity.get(threadId) ?? { active: false, lastUsedAt: Date.now() }; + activity.active = true; + this.threadActivity.set(threadId, activity); const pending = createPendingTurn(); this.pendingTurns.set(threadId, pending); try { @@ -84,6 +89,9 @@ export class CodexAppServerRuntime implements HarnessRuntime { } catch (error) { if (this.pendingTurns.get(threadId) === pending) this.pendingTurns.delete(threadId); throw error; + } finally { + activity.active = false; + activity.lastUsedAt = Date.now(); } } @@ -91,6 +99,18 @@ export class CodexAppServerRuntime implements HarnessRuntime { return !this.closed && this.connection.isUsable(); } + async reapIdleSessions(now: number): Promise { + for (const [threadId, activity] of this.threadActivity) { + if (activity.active || now - activity.lastUsedAt < CODEX_THREAD_IDLE_MS) continue; + try { + await this.connection.request("thread/unsubscribe", { threadId }); + this.threadActivity.delete(threadId); + } catch { + // Keep it tracked so a later reap can retry if the connection recovers. + } + } + } + async close(): Promise { if (this.closed) return; this.closed = true; @@ -98,6 +118,7 @@ export class CodexAppServerRuntime implements HarnessRuntime { this.unsubscribeClose(); const error = new Error("Codex app-server runtime closed."); this.rejectPendingTurns(error); + this.threadActivity.clear(); await this.connection.close(); } diff --git a/src/local-agent-pi/runtime.test.ts b/src/local-agent-pi/runtime.test.ts index 7eea2ac47..59f959684 100644 --- a/src/local-agent-pi/runtime.test.ts +++ b/src/local-agent-pi/runtime.test.ts @@ -106,6 +106,18 @@ try { }), /completed without a final response/, ); + + await runtime.reapIdleSessions(Date.now() + 5 * 60 * 1_000 + 1); + assert.equal(sessions.get("pi_1")?.disposed, true); + assert.equal(sessions.get("pi_2")?.disposed, true); + + const resumedAfterReap = await runtime.run({ + workspace: "/tmp/a", + prompt: "after reap", + providerSessionId: "pi_1", + }); + assert.equal(created, 3, "an evicted Pi session should be recreated from its durable session id"); + assert.equal(resumedAfterReap.providerSessionId, "pi_1"); } finally { await runtime.close(); } diff --git a/src/local-agent-pi/runtime.ts b/src/local-agent-pi/runtime.ts index 1c71de7b2..d9b869c58 100644 --- a/src/local-agent-pi/runtime.ts +++ b/src/local-agent-pi/runtime.ts @@ -1,6 +1,8 @@ import type { LocalAgentRunInput, LocalAgentRunResult } from "../local-agent-runtime.js"; import type { HarnessDriver, HarnessRuntime } from "../local-agent-runtime-pool.js"; +const PI_SESSION_IDLE_MS = 5 * 60 * 1_000; + interface PiSessionLike { readonly sessionId: string; readonly messages: unknown[]; @@ -20,6 +22,8 @@ interface PiSessionEntry { workspace: string; session: PiSessionLike; tail: Promise; + active: boolean; + lastUsedAt: number; } type PiSessionSlot = @@ -37,7 +41,15 @@ export class PiHarnessRuntime implements HarnessRuntime { const entry = await this.ensureSession(input); const turn = entry.tail .catch(() => undefined) - .then(() => this.runTurn(entry.session, input)); + .then(async () => { + entry.active = true; + try { + return await this.runTurn(entry.session, input); + } finally { + entry.active = false; + entry.lastUsedAt = Date.now(); + } + }); entry.tail = turn.then(() => undefined, () => undefined); return turn; } @@ -68,6 +80,20 @@ export class PiHarnessRuntime implements HarnessRuntime { return !this.closed; } + async reapIdleSessions(now: number): Promise { + for (const [sessionId, slot] of this.sessions) { + if (slot.status !== "ready") continue; + const entry = slot.entry; + if (entry.active || now - entry.lastUsedAt < PI_SESSION_IDLE_MS) continue; + try { + entry.session.dispose(); + this.sessions.delete(sessionId); + } catch { + // Keep the entry so a later reap can retry disposal. + } + } + } + async close(): Promise { if (this.closed) return; this.closed = true; @@ -141,7 +167,13 @@ export class PiHarnessRuntime implements HarnessRuntime { } function createPiSessionEntry(workspace: string, session: PiSessionLike): PiSessionEntry { - return { workspace, session, tail: Promise.resolve() }; + return { + workspace, + session, + tail: Promise.resolve(), + active: false, + lastUsedAt: Date.now(), + }; } export function createPiHarnessDriver(): HarnessDriver { diff --git a/src/local-agent-runtime-pool.test.ts b/src/local-agent-runtime-pool.test.ts index b43652029..b31ff663d 100644 --- a/src/local-agent-runtime-pool.test.ts +++ b/src/local-agent-runtime-pool.test.ts @@ -5,6 +5,7 @@ import type { LocalAgentRunInput, LocalAgentRunResult } from "./local-agent-runt class FakeRuntime implements HarnessRuntime { readonly prompts: string[] = []; + readonly reaps: number[] = []; closed = false; private sessionCount = 0; @@ -26,6 +27,10 @@ class FakeRuntime implements HarnessRuntime { return !this.closed; } + async reapIdleSessions(now: number): Promise { + this.reaps.push(now); + } + async close(): Promise { this.closed = true; } @@ -75,10 +80,12 @@ try { now = 9; await pool.reapIdle(); assert.equal(runtimes[0]?.closed, false); + assert.deepEqual(runtimes[0]?.reaps, [9]); now = 10; await pool.reapIdle(); assert.equal(runtimes[0]?.closed, true, "idle runtimes should be reclaimable without touching durable session ids"); + assert.deepEqual(runtimes[0]?.reaps, [9, 10]); await registry.run("opencode", input("/tmp/a", "third", second.providerSessionId ?? undefined)); assert.equal(created, 2, "a later run should recreate an evicted runtime"); diff --git a/src/local-agent-runtime-pool.ts b/src/local-agent-runtime-pool.ts index 3e9b5a54e..4fecfec8f 100644 --- a/src/local-agent-runtime-pool.ts +++ b/src/local-agent-runtime-pool.ts @@ -6,6 +6,7 @@ const DEFAULT_REAP_INTERVAL_MS = 30 * 1_000; export interface HarnessRuntime { run(input: LocalAgentRunInput): Promise; isUsable(): boolean; + reapIdleSessions?(now: number): Promise; close(): Promise; } @@ -81,6 +82,7 @@ export class HarnessRuntimePool { const closing: HarnessRuntime[] = []; for (const [key, slot] of this.slots) { if (slot.status !== "ready" || slot.activeRuns > 0) continue; + await slot.runtime.reapIdleSessions?.(now).catch(() => undefined); if (now - slot.lastUsedAt < this.idleMs) continue; this.slots.delete(key); closing.push(slot.runtime); From 06045726998293e515ac966af6f95e1d13a02fa5 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 07:44:20 +0530 Subject: [PATCH 5/7] fix(agents): serialize runtime maintenance with turns --- src/local-agent-runtime-pool.test.ts | 71 ++++++++++++++++++++++++++++ src/local-agent-runtime-pool.ts | 49 ++++++++++++++----- 2 files changed, 108 insertions(+), 12 deletions(-) diff --git a/src/local-agent-runtime-pool.test.ts b/src/local-agent-runtime-pool.test.ts index b31ff663d..842470119 100644 --- a/src/local-agent-runtime-pool.test.ts +++ b/src/local-agent-runtime-pool.test.ts @@ -52,6 +52,17 @@ class FailingRuntime extends FakeRuntime { } } +class BlockingReapRuntime extends FakeRuntime { + readonly reapStarted = deferred(); + readonly finishReap = deferred(); + + override async reapIdleSessions(now: number): Promise { + this.reaps.push(now); + this.reapStarted.resolve(); + await this.finishReap.promise; + } +} + let now = 0; let created = 0; const runtimes: FakeRuntime[] = []; @@ -144,6 +155,51 @@ try { } } +{ + let raceNow = 0; + let raceCreates = 0; + const raceRuntimes: FakeRuntime[] = []; + const raceDriver: HarnessDriver = { + provider: "opencode", + runtimeKey: () => "reap-race", + async createRuntime() { + raceCreates += 1; + const runtime = raceCreates === 1 + ? new BlockingReapRuntime("reap-runtime-1") + : new FakeRuntime(`reap-runtime-${raceCreates}`); + raceRuntimes.push(runtime); + return runtime; + }, + }; + const racePool = new HarnessRuntimePool({ idleMs: 10, reapIntervalMs: 0, now: () => raceNow }); + + try { + await racePool.run(raceDriver, input("/tmp/a", "seed")); + const firstRuntime = raceRuntimes[0] as BlockingReapRuntime; + raceNow = 10; + const reaping = racePool.reapIdle(); + await firstRuntime.reapStarted.promise; + + const concurrentRun = racePool.run(raceDriver, input("/tmp/a", "during reap")); + await immediate(); + assert.deepEqual( + firstRuntime.prompts, + ["seed"], + "a provider turn must not overlap runtime/session maintenance", + ); + + firstRuntime.finishReap.resolve(); + await reaping; + const result = await concurrentRun; + + assert.equal(firstRuntime.closed, true, "the previously idle runtime may be reclaimed after maintenance"); + assert.equal(raceCreates, 2, "a waiting turn should reacquire a fresh runtime after idle reclamation"); + assert.equal(result.finalResponse, "reap-runtime-2:during reap"); + } finally { + await racePool.shutdown(); + } +} + function input(workspace: string, prompt: string, providerSessionId?: string): LocalAgentRunInput { return { workspace, @@ -162,3 +218,18 @@ function countingDriver(provider: string, next: () => number): HarnessDriver { }, }; } + +function deferred(): { + promise: Promise; + resolve(value: T): void; +} { + let resolve!: (value: T) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + +function immediate(): Promise { + return new Promise((resolve) => setImmediate(resolve)); +} diff --git a/src/local-agent-runtime-pool.ts b/src/local-agent-runtime-pool.ts index 4fecfec8f..c5a896518 100644 --- a/src/local-agent-runtime-pool.ts +++ b/src/local-agent-runtime-pool.ts @@ -26,6 +26,7 @@ type RuntimeSlot = runtime: HarnessRuntime; activeRuns: number; lastUsedAt: number; + maintenance?: Promise; }; interface HarnessRuntimePoolOptions { @@ -64,8 +65,7 @@ export class HarnessRuntimePool { ): Promise { if (this.closing) throw new Error("Harness runtime pool is shutting down."); const key = `${driver.provider}\0${driver.runtimeKey(input)}`; - const slot = await this.acquire(key, driver, input); - slot.activeRuns += 1; + const slot = await this.acquireForRun(key, driver, input); try { return await slot.runtime.run(input); } finally { @@ -79,15 +79,20 @@ export class HarnessRuntimePool { } async reapIdle(now = this.now()): Promise { - const closing: HarnessRuntime[] = []; + const maintenance: Promise[] = []; for (const [key, slot] of this.slots) { - if (slot.status !== "ready" || slot.activeRuns > 0) continue; - await slot.runtime.reapIdleSessions?.(now).catch(() => undefined); - if (now - slot.lastUsedAt < this.idleMs) continue; - this.slots.delete(key); - closing.push(slot.runtime); + if (slot.status !== "ready" || slot.activeRuns > 0 || slot.maintenance) continue; + const task = this.maintainSlot(key, slot, now); + slot.maintenance = task; + maintenance.push((async () => { + try { + await task; + } finally { + if (slot.maintenance === task) slot.maintenance = undefined; + } + })()); } - await Promise.allSettled(closing.map((runtime) => runtime.close())); + await Promise.allSettled(maintenance); } async shutdown(): Promise { @@ -99,11 +104,12 @@ export class HarnessRuntimePool { await Promise.allSettled(slots.map(async (slot) => { const runtime = slot.status === "ready" ? slot.runtime : await slot.promise; + if (slot.status === "ready") await slot.maintenance?.catch(() => undefined); await runtime.close(); })); } - private async acquire( + private async acquireForRun( key: string, driver: HarnessDriver, input: LocalAgentRunInput, @@ -112,7 +118,14 @@ export class HarnessRuntimePool { if (this.closing) throw new Error("Harness runtime pool is shutting down."); const existing = this.slots.get(key); if (existing?.status === "ready") { - if (existing.runtime.isUsable()) return existing; + if (existing.maintenance) { + await existing.maintenance.catch(() => undefined); + continue; + } + if (existing.runtime.isUsable()) { + existing.activeRuns += 1; + return existing; + } this.slots.delete(key); await existing.runtime.close().catch(() => undefined); continue; @@ -137,7 +150,7 @@ export class HarnessRuntimePool { const ready: Extract = { status: "ready", runtime, - activeRuns: 0, + activeRuns: 1, lastUsedAt: this.now(), }; this.slots.set(key, ready); @@ -148,4 +161,16 @@ export class HarnessRuntimePool { } } } + + private async maintainSlot( + key: string, + slot: Extract, + now: number, + ): Promise { + await slot.runtime.reapIdleSessions?.(now).catch(() => undefined); + if (this.slots.get(key) !== slot || slot.activeRuns > 0) return; + if (now - slot.lastUsedAt < this.idleMs) return; + this.slots.delete(key); + await slot.runtime.close(); + } } From 22601285d7eef669a922e0ada786252d821ece25 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 07:57:45 +0530 Subject: [PATCH 6/7] fix(claude): bind warm runtimes to one workspace --- src/local-agent-claude/runtime.test.ts | 12 ++++++++++++ src/local-agent-claude/runtime.ts | 5 +++++ 2 files changed, 17 insertions(+) diff --git a/src/local-agent-claude/runtime.test.ts b/src/local-agent-claude/runtime.test.ts index 7084fb280..a08e99939 100644 --- a/src/local-agent-claude/runtime.test.ts +++ b/src/local-agent-claude/runtime.test.ts @@ -95,6 +95,18 @@ try { assert.deepEqual(efforts, ["high", "xhigh"]); assert.equal(changedEffort.providerSessionId, "claude_1"); assert.equal(changedEffort.finalResponse, "response:third"); + + await assert.rejects( + runtime.run({ + agentId: "agt_claude", + workspace: "/tmp/other-project", + prompt: "wrong workspace", + providerSessionId: changedEffort.providerSessionId ?? undefined, + thinking: "xhigh", + }), + /belongs to workspace \/tmp\/project, not \/tmp\/other-project/, + ); + assert.equal(queries.length, 2, "workspace mismatch must not create or reuse a query in another cwd"); } finally { await runtime.close(); } diff --git a/src/local-agent-claude/runtime.ts b/src/local-agent-claude/runtime.ts index 9d6dd476b..f6b843990 100644 --- a/src/local-agent-claude/runtime.ts +++ b/src/local-agent-claude/runtime.ts @@ -35,12 +35,17 @@ interface LiveClaudeQuery { export class ClaudeWarmRuntime implements HarnessRuntime { private live: LiveClaudeQuery | undefined; private providerSessionId: string | undefined; + private workspace: string | undefined; private closed = false; constructor(private readonly createQuery: ClaudeQueryFactory) {} async run(input: LocalAgentRunInput): Promise { if (this.closed) throw new Error("Claude runtime is closed."); + if (this.workspace && this.workspace !== input.workspace) { + throw new Error(`Claude agent runtime belongs to workspace ${this.workspace}, not ${input.workspace}.`); + } + this.workspace ??= input.workspace; this.providerSessionId = input.providerSessionId ?? this.providerSessionId; const live = this.ensureQuery(input); From fbc19e6887889ae6ead4f5e795ae0473313a18ee Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 12 Aug 2026 07:57:45 +0530 Subject: [PATCH 7/7] fix(agents): defer maintenance until runtimes are idle --- src/local-agent-runtime-pool.test.ts | 43 ++++++++++++++++++++++++++++ src/local-agent-runtime-pool.ts | 40 +++++++++++++++++++------- 2 files changed, 73 insertions(+), 10 deletions(-) diff --git a/src/local-agent-runtime-pool.test.ts b/src/local-agent-runtime-pool.test.ts index 842470119..26a3f2a3a 100644 --- a/src/local-agent-runtime-pool.test.ts +++ b/src/local-agent-runtime-pool.test.ts @@ -63,6 +63,19 @@ class BlockingReapRuntime extends FakeRuntime { } } +class BlockingRunRuntime extends FakeRuntime { + readonly runStarted = deferred(); + readonly finishRun = deferred(); + + override async run(input: LocalAgentRunInput): Promise { + if (input.prompt === "hold") { + this.runStarted.resolve(); + await this.finishRun.promise; + } + return super.run(input); + } +} + let now = 0; let created = 0; const runtimes: FakeRuntime[] = []; @@ -200,6 +213,36 @@ try { } } +{ + let deferredNow = 0; + const runtime = new BlockingRunRuntime("deferred-reap"); + const deferredDriver: HarnessDriver = { + provider: "opencode", + runtimeKey: () => "deferred-reap", + createRuntime: async () => runtime, + }; + const deferredPool = new HarnessRuntimePool({ idleMs: 100, reapIntervalMs: 0, now: () => deferredNow }); + + try { + const activeRun = deferredPool.run(deferredDriver, input("/tmp/a", "hold")); + await runtime.runStarted.promise; + deferredNow = 10; + await deferredPool.reapIdle(); + assert.deepEqual(runtime.reaps, [], "maintenance must not overlap an active provider turn"); + + runtime.finishRun.resolve(); + await activeRun; + await immediate(); + assert.deepEqual( + runtime.reaps, + [10], + "a reap requested during activity should run when the runtime next becomes idle", + ); + } finally { + await deferredPool.shutdown(); + } +} + function input(workspace: string, prompt: string, providerSessionId?: string): LocalAgentRunInput { return { workspace, diff --git a/src/local-agent-runtime-pool.ts b/src/local-agent-runtime-pool.ts index c5a896518..4ad6184e6 100644 --- a/src/local-agent-runtime-pool.ts +++ b/src/local-agent-runtime-pool.ts @@ -27,6 +27,7 @@ type RuntimeSlot = activeRuns: number; lastUsedAt: number; maintenance?: Promise; + reapRequested: boolean; }; interface HarnessRuntimePoolOptions { @@ -74,6 +75,14 @@ export class HarnessRuntimePool { if (slot.activeRuns === 0 && !slot.runtime.isUsable()) { if (this.slots.get(key) === slot) this.slots.delete(key); await slot.runtime.close().catch(() => undefined); + } else if ( + slot.activeRuns === 0 + && slot.reapRequested + && !slot.maintenance + && this.slots.get(key) === slot + ) { + slot.reapRequested = false; + this.beginMaintenance(key, slot, this.now()); } } } @@ -81,16 +90,13 @@ export class HarnessRuntimePool { async reapIdle(now = this.now()): Promise { const maintenance: Promise[] = []; for (const [key, slot] of this.slots) { - if (slot.status !== "ready" || slot.activeRuns > 0 || slot.maintenance) continue; - const task = this.maintainSlot(key, slot, now); - slot.maintenance = task; - maintenance.push((async () => { - try { - await task; - } finally { - if (slot.maintenance === task) slot.maintenance = undefined; - } - })()); + if (slot.status !== "ready" || slot.maintenance) continue; + if (slot.activeRuns > 0) { + slot.reapRequested = true; + continue; + } + slot.reapRequested = false; + maintenance.push(this.beginMaintenance(key, slot, now)); } await Promise.allSettled(maintenance); } @@ -152,6 +158,7 @@ export class HarnessRuntimePool { runtime, activeRuns: 1, lastUsedAt: this.now(), + reapRequested: false, }; this.slots.set(key, ready); return ready; @@ -173,4 +180,17 @@ export class HarnessRuntimePool { this.slots.delete(key); await slot.runtime.close(); } + + private beginMaintenance( + key: string, + slot: Extract, + now: number, + ): Promise { + const task = this.maintainSlot(key, slot, now); + slot.maintenance = task; + void task.finally(() => { + if (slot.maintenance === task) slot.maintenance = undefined; + }).catch(() => undefined); + return task; + } }