diff --git a/apps/server/src/application/bootstrap/server-bootstrap.ts b/apps/server/src/application/bootstrap/server-bootstrap.ts index f87480d20..5975addb5 100644 --- a/apps/server/src/application/bootstrap/server-bootstrap.ts +++ b/apps/server/src/application/bootstrap/server-bootstrap.ts @@ -109,7 +109,7 @@ import { ProviderUsageWarmupService } from "../../features/providers/availabilit import { ProviderRegistry } from "../../features/providers/composition/provider-registry.js"; import { ProviderEventIngress } from "../../features/providers/composition/provider-event-ingress.js"; import type { CursorProviderBoundary } from "@mcode/providers"; -import { ModelCacheService } from "../../features/providers/models/model-cache-service.js"; +import { ModelCacheService, startupModelProviderIds } from "../../features/providers/models/model-cache-service.js"; import { DiffSummaryService } from "../../features/projects/diffs/summaries/diff-summary-service.js"; import { RecapService } from "../../features/agents/recap/recap-service.js"; import { seedAgentRuntimeWorkspace } from "../../runtime/startup/dev-agent-seed.js"; @@ -571,11 +571,8 @@ providerAvailability // blocking `codex --version` spawnSync. warmCodexVersionGate(); providerUsageWarmup.warmEnabledProviders(true); - // Warm the model cache once after CLI verification has gated which providers - // are usable. Triggering this per WS connect would spam refreshes; running - // it once at startup is sufficient because ModelCacheService also refreshes - // lazily on stale reads (stale-while-revalidate). - void modelCacheService.refreshAll().catch((err: unknown) => { + const modelProviders = startupModelProviderIds(providerAvailability.listAvailability()); + void modelCacheService.refreshProviders(modelProviders).catch((err: unknown) => { logger.warn("Model cache startup refresh failed", { error: err instanceof Error ? err.message : String(err), }); @@ -940,9 +937,18 @@ async function shutdown(): Promise { shutdownCoordinator.setPhase("stop agent sessions"); await agentService.stopAll(); + let shutdownFailure: unknown = null; + const captureCleanupFailure = async (cleanup: () => Promise | void): Promise => { + try { + await cleanup(); + } catch (error) { + shutdownFailure ??= error; + } + }; + // 2. Shutdown provider registry shutdownCoordinator.setPhase("shutdown providers"); - await providerRegistry.shutdown(); + await captureCleanupFailure(() => providerRegistry.shutdown()); shutdownCoordinator.setPhase("shutdown provider event workers"); providerEventIngress.shutdown(); browserAutomationBroker.shutdown(); @@ -951,15 +957,6 @@ async function shutdown(): Promise { // 3. Dispose settings file watcher settingsService.dispose(); - let shutdownFailure: unknown = null; - const captureCleanupFailure = async (cleanup: () => Promise | void): Promise => { - try { - await cleanup(); - } catch (error) { - shutdownFailure ??= error; - } - }; - // 6. Contain Project command sessions before their Terminal dependency shuts down. shutdownCoordinator.setPhase("shutdown Project commands"); await captureCleanupFailure(() => projectActionService.dispose()); diff --git a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-idle.test.ts b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-idle.test.ts index 67a24b7d6..fbb1a23ee 100644 --- a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-idle.test.ts +++ b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-idle.test.ts @@ -123,7 +123,7 @@ describe("OpenCodeProvider idle confirmation", () => { await provider.sendTurn(turnRequest()); expect(http.getSessionStatus).toHaveBeenCalledTimes(2); expect(endedOutcomes(submitted)).toEqual(["completed"]); - provider.shutdown(); + await provider.shutdown(); }); it("treats a drained session as quiet, not as a poll failure", async () => { @@ -131,7 +131,7 @@ describe("OpenCodeProvider idle confirmation", () => { const { provider, submitted } = testProvider(http as never); await provider.sendTurn(turnRequest()); expect(endedOutcomes(submitted)).toEqual(["completed"]); - provider.shutdown(); + await provider.shutdown(); }); it("restarts confirmation when mapped activity arrives between idles", async () => { @@ -166,7 +166,7 @@ describe("OpenCodeProvider idle confirmation", () => { await sending; expect(endedOutcomes(submitted)).toEqual(["completed"]); expect(polls).toBeGreaterThanOrEqual(3); - provider.shutdown(); + await provider.shutdown(); }); it("abandons confirmation while the session reports busy", async () => { @@ -189,7 +189,7 @@ describe("OpenCodeProvider idle confirmation", () => { await sending; expect(endedOutcomes(submitted)).toEqual(["completed"]); expect(http.getSessionStatus.mock.calls.length).toBeGreaterThanOrEqual(3); - provider.shutdown(); + await provider.shutdown(); }); it("never settles while a permission card is pending", async () => { @@ -216,7 +216,7 @@ describe("OpenCodeProvider idle confirmation", () => { expect(provider.resolvePermission("per_1", "allow")).toBe(true); await sending; expect(endedOutcomes(submitted)).toEqual(["completed"]); - provider.shutdown(); + await provider.shutdown(); }); it("settles errored after repeated status poll failures", async () => { @@ -231,7 +231,7 @@ describe("OpenCodeProvider idle confirmation", () => { const events = submittedEvents(submitted); expect(events.filter((event) => event.type === "error")).toHaveLength(1); expect(endedOutcomes(submitted)).toEqual(["errored"]); - provider.shutdown(); + await provider.shutdown(); }); it("aborts a hung status poll when Stop cancels the turn", async () => { @@ -251,7 +251,7 @@ describe("OpenCodeProvider idle confirmation", () => { expect(statusSignal?.aborted).toBe(true); expect(endedOutcomes(submitted)).toEqual(["cancelled"]); - provider.shutdown(); + await provider.shutdown(); }); it("emits one terminal outcome for duplicate idles", async () => { @@ -259,6 +259,6 @@ describe("OpenCodeProvider idle confirmation", () => { const { provider, submitted } = testProvider(http as never); await provider.sendTurn(turnRequest()); expect(endedOutcomes(submitted)).toEqual(["completed"]); - provider.shutdown(); + await provider.shutdown(); }); }); diff --git a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-permissions.test.ts b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-permissions.test.ts index 2fd0daec1..b8e4fe887 100644 --- a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-permissions.test.ts +++ b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-permissions.test.ts @@ -173,7 +173,7 @@ describe("OpenCodeProvider permission flow", () => { .filter((event) => event.type === "ended") .map((event) => event.outcome); expect(outcomes).toEqual(["completed"]); - provider.shutdown(); + await provider.shutdown(); }); it("relays reject and resolves unknown ids as false", async () => { @@ -191,7 +191,7 @@ describe("OpenCodeProvider permission flow", () => { ); expect(resolved).toEqual([{ requestId: "per_1", decision: "deny" }]); expect(provider.resolvePermission("nope", "allow")).toBe(false); - provider.shutdown(); + await provider.shutdown(); }); it("relays exact question selections, rejects invalid answers locally, and rejects on deny", async () => { @@ -245,7 +245,7 @@ describe("OpenCodeProvider permission flow", () => { "http://127.0.0.1:4096", "ses_1", "que_2", "v2", expect.objectContaining({ signal: expect.any(AbortSignal) }), ); expect(http.rejectQuestion).toHaveBeenCalledTimes(1); - provider.shutdown(); + await provider.shutdown(); }); it("keeps a failed reply answerable instead of stalling the turn", async () => { @@ -263,7 +263,7 @@ describe("OpenCodeProvider permission flow", () => { await sending; expect(http.replyPermission).toHaveBeenCalledTimes(2); expect(provider.listPendingPermissions("thread-1")).toHaveLength(0); - provider.shutdown(); + await provider.shutdown(); }); it("drains pending cards as cancelled on stop", async () => { @@ -287,7 +287,7 @@ describe("OpenCodeProvider permission flow", () => { await sending; expect(resolved).toEqual([{ requestId: "per_1", decision: "cancelled" }]); expect(provider.listPendingPermissions("thread-1")).toHaveLength(0); - provider.shutdown(); + await provider.shutdown(); }); it("cancels an in-flight reply without a second local settlement", async () => { @@ -320,7 +320,7 @@ describe("OpenCodeProvider permission flow", () => { expect(resolved).toEqual([{ requestId: "per_1", decision: "cancelled" }]); expect(provider.resolvePermission("per_1", "allow")).toBe(false); expect(http.replyPermission).toHaveBeenCalledTimes(1); - provider.shutdown(); + await provider.shutdown(); }); it("drains a replayed ask after stream termination so it can card again", async () => { @@ -346,7 +346,7 @@ describe("OpenCodeProvider permission flow", () => { { requestId: "per_1", decision: "cancelled" }, { requestId: "per_1", decision: "cancelled" }, ]); - provider.shutdown(); + await provider.shutdown(); }); it("drains a pending ask on provider failure and shutdown", async () => { @@ -361,7 +361,7 @@ describe("OpenCodeProvider permission flow", () => { const retry = provider.sendTurn({ ...turnRequest(), turnId: "turn-2", turnExecutionId: "66666666-6666-4666-8666-666666666666" }); await vi.waitFor(() => expect(provider.listPendingPermissions("thread-1")).toHaveLength(1)); - provider.shutdown(); + await provider.shutdown(); await retry; expect(resolved).toEqual([ { requestId: "per_1", decision: "cancelled" }, @@ -385,7 +385,7 @@ describe("OpenCodeProvider permission flow", () => { const events = submittedEvents(submitted); expect(events.filter((event) => event.type === "system" && event.subtype === "sdk_session_invalidated")).toHaveLength(1); expect(events.filter((event) => event.type === "ended").map((event) => event.outcome)).toEqual(["cancelled"]); - provider.shutdown(); + await provider.shutdown(); }); }); @@ -410,7 +410,7 @@ describe("OpenCodeProvider notice dedup", () => { && event.requestedModel === "anthropic/claude-sonnet-4-6" && event.actualModel === "anthropic/claude-sonnet-4-6" ))).toHaveLength(1); - provider.shutdown(); + await provider.shutdown(); }); it("surfaces one diagnostic row for a malformed ask without carding", async () => { @@ -429,6 +429,6 @@ describe("OpenCodeProvider notice dedup", () => { .filter((event) => event.type === "system") .map((event) => event.subtype); expect(subtypes.filter((subtype) => subtype === "provider.notice.malformed-request")).toHaveLength(1); - provider.shutdown(); + await provider.shutdown(); }); }); diff --git a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-resume.test.ts b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-resume.test.ts index b4e4ef8a8..ffe33f46f 100644 --- a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-resume.test.ts +++ b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider-resume.test.ts @@ -106,7 +106,7 @@ describe("OpenCodeProvider resume cursor", () => { expect(http.promptAsync).toHaveBeenCalledWith("http://127.0.0.1:4096", "ses_kept", expect.anything()); const events = submittedEvents(submitted); expect(events).toContainEqual(expect.objectContaining({ subtype: "sdk_session_id:ses_kept" })); - provider.shutdown(); + await provider.shutdown(); }); it("ignores an unknown cursor version and starts fresh", async () => { @@ -117,7 +117,7 @@ describe("OpenCodeProvider resume cursor", () => { })); expect(http.createSession).toHaveBeenCalledTimes(1); expect(http.promptAsync).toHaveBeenCalledWith("http://127.0.0.1:4096", "ses_brand_new", expect.anything()); - provider.shutdown(); + await provider.shutdown(); }); it("starts fresh with a visible notice when the adopted session is gone (deleted upstream)", async () => { @@ -137,7 +137,7 @@ describe("OpenCodeProvider resume cursor", () => { type: "system", subtype: "sdk_session_invalidated", })); - provider.shutdown(); + await provider.shutdown(); }); it("starts fresh with a visible notice on a prompt-time 404 race", async () => { @@ -155,7 +155,7 @@ describe("OpenCodeProvider resume cursor", () => { type: "system", subtype: "sdk_session_invalidated", })); - provider.shutdown(); + await provider.shutdown(); }); it("leaves other threads alone when one upstream session is recreated", async () => { @@ -182,7 +182,7 @@ describe("OpenCodeProvider resume cursor", () => { expect(http.createSession).toHaveBeenCalledTimes(1); const prompted = http.promptAsync.mock.calls.map(([, sessionId]) => sessionId); expect(prompted).toEqual(["ses_a", "ses_fresh_other", "ses_a"]); - provider.shutdown(); + await provider.shutdown(); }); it("only a confirmed 404 starts fresh; other verify failures propagate without reset", async () => { @@ -195,6 +195,6 @@ describe("OpenCodeProvider resume cursor", () => { await provider.sendTurn(turnRequest({ resumeFrom: "ses_live" })); expect(http.createSession).not.toHaveBeenCalled(); expect(http.promptAsync).not.toHaveBeenCalled(); - provider.shutdown(); + await provider.shutdown(); }); }); diff --git a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider.test.ts b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider.test.ts index 5032ffe75..350fdf51b 100644 --- a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider.test.ts +++ b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-provider.test.ts @@ -2,9 +2,10 @@ import "reflect-metadata"; import { describe, expect, it, vi } from "vitest"; import { OpenCodeProvider, toOpenCodeModelRef } from "../opencode-provider.js"; import { OpenCodeServerPool } from "../opencode-server-pool.js"; +import type { OpenCodeHttpClient } from "../opencode-http-client.js"; import type { TurnRequest } from "@mcode/contracts"; -function testProvider(http: never, pool: OpenCodeServerPool) { +function testProvider(http: OpenCodeHttpClient | undefined, pool: OpenCodeServerPool) { const settingsService = { get: () => ({ provider: { cli: { opencode: "opencode" } } }) }; const envService = { getEnv: () => ({}) }; const submitted: unknown[] = []; @@ -20,7 +21,7 @@ function testProvider(http: never, pool: OpenCodeServerPool) { const provider = new OpenCodeProvider(settingsService as never, envService as never, host as never); provider.configureTestSeams({ pool, - http: http as never, + ...(http ? { http } : {}), probeCli: async () => ({ binaryPath: "opencode", version: "test" }), idleConfirm: { intervalMs: 5, requiredPolls: 2, timeoutMs: 500, maxPollErrors: 2 }, }); @@ -59,6 +60,35 @@ describe("toOpenCodeModelRef", () => { }); describe("OpenCodeProvider minimal turn", () => { + it("does not finish shutdown until its server has terminated", async () => { + let releaseTermination = () => {}; + let signalTerminationStarted = () => {}; + const terminationStarted = new Promise((resolve) => { signalTerminationStarted = resolve; }); + const terminationAllowed = new Promise((resolve) => { releaseTermination = resolve; }); + const pool = new OpenCodeServerPool({ + spawn: () => ({ pid: 4242, on: () => {}, off: () => {}, kill: () => true }), + waitForHealth: async () => {}, + terminateTree: async () => { + signalTerminationStarted(); + await terminationAllowed; + }, + findFreePort: async () => 4096, + now: () => Date.now(), + env: () => ({}), + }); + await pool.acquire({ binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }); + const { provider } = testProvider(undefined, pool); + + let completed = false; + const shutdown = provider.shutdown().then(() => { completed = true; }); + await terminationStarted; + expect(completed).toBe(false); + + releaseTermination(); + await shutdown; + expect(completed).toBe(true); + }); + it("streams a reply to completion and returns the pool server warm", async () => { const pool = new OpenCodeServerPool({ spawn: () => ({ pid: 1, on: () => {}, off: () => {}, kill: () => true }) as never, @@ -92,7 +122,7 @@ describe("OpenCodeProvider minimal turn", () => { }); expect(pool.size).toBe(1); expect(submitted.length).toBeGreaterThan(0); - provider.shutdown(); + await provider.shutdown(); }); it("stop aborts the upstream session and settles as cancelled with no further output", async () => { @@ -129,7 +159,7 @@ describe("OpenCodeProvider minimal turn", () => { await sending; expect(http.abortSession).toHaveBeenCalledTimes(1); expect(pool.size).toBe(1); - provider.shutdown(); + await provider.shutdown(); }); it("starts fresh once when the adopted upstream session is gone (404)", async () => { @@ -166,7 +196,7 @@ describe("OpenCodeProvider minimal turn", () => { .turns.get("mcode-thread-1")!.upstreamSessionId = "ses_stale"; await provider.sendTurn(turnRequest()); expect(http.createSession).toHaveBeenCalledTimes(2); - provider.shutdown(); + await provider.shutdown(); }); it("discardSession drops the adopted upstream session", async () => { @@ -194,6 +224,6 @@ describe("OpenCodeProvider minimal turn", () => { await provider.discardSession("mcode-thread-1"); await provider.sendTurn(turnRequest()); expect(http.createSession).toHaveBeenCalledTimes(2); - provider.shutdown(); + await provider.shutdown(); }); }); diff --git a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-server-pool.test.ts b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-server-pool.test.ts index c6b5b2596..67d4fc508 100644 --- a/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-server-pool.test.ts +++ b/apps/server/src/features/providers/adapters/opencode/__tests__/opencode-server-pool.test.ts @@ -18,6 +18,7 @@ function fakeDeps(overrides: Record = {}) { spawn: vi.fn(() => child), waitForHealth: vi.fn(async () => {}), terminateTree: vi.fn(async () => {}), + waitForPortClosed: vi.fn(async () => {}), findFreePort: vi.fn(async () => 4096), now: vi.fn(() => 1_000), env: vi.fn(() => ({})), @@ -63,16 +64,19 @@ describe("OpenCodeServerPool", () => { await pool.shutdown(); }); - it("falls back to a direct kill when tree termination fails", async () => { - const { deps, child } = fakeDeps({ terminateTree: vi.fn(async () => { throw new Error("taskkill failed"); }) }); + it("keeps ownership and reports failure when tree termination fails", async () => { + const terminateTree = vi.fn(async () => { throw new Error("taskkill failed"); }); + const { deps, child } = fakeDeps({ terminateTree }); const pool = new OpenCodeServerPool(deps); const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; await pool.acquire(key); pool.release(key); - await pool.closeIdle(1_000 + OPENCODE_POOL_IDLE_TTL_MS + 1); - expect(child.kill).toHaveBeenCalledTimes(1); - expect(pool.size).toBe(0); + await expect(pool.closeIdle(1_000 + OPENCODE_POOL_IDLE_TTL_MS + 1)).rejects.toThrow("taskkill failed"); + expect(child.kill).not.toHaveBeenCalled(); + expect(pool.size).toBe(1); + terminateTree.mockImplementation(async () => {}); await pool.shutdown(); + expect(pool.size).toBe(0); }); it("cleans up without orphans when the child exits unexpectedly", async () => { @@ -80,12 +84,151 @@ describe("OpenCodeServerPool", () => { const pool = new OpenCodeServerPool(deps); const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; await pool.acquire(key); + pool.release(key); for (const listener of listeners.get("exit") ?? []) listener(null, null); + await vi.waitFor(() => expect(pool.size).toBe(0)); expect(pool.size).toBe(0); expect(child.kill).not.toHaveBeenCalled(); await pool.shutdown(); }); + it("waits for health before sharing a server with a concurrent acquire", async () => { + let finishHealth: (() => void) | undefined; + const health = new Promise((resolve) => { finishHealth = resolve; }); + const { deps } = fakeDeps({ waitForHealth: vi.fn(() => health) }); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + const first = pool.acquire(key); + await vi.waitFor(() => expect(pool.entryFor(key)?.ready).toBe(false)); + let secondResolved = false; + const second = pool.acquire(key).then((entry) => { secondResolved = true; return entry; }); + await Promise.resolve(); + expect(secondResolved).toBe(false); + finishHealth?.(); + const [firstEntry, secondEntry] = await Promise.all([first, second]); + expect(firstEntry).toBe(secondEntry); + expect(firstEntry.refs).toBe(2); + expect(firstEntry.ready).toBe(true); + expect(deps.spawn).toHaveBeenCalledTimes(1); + await pool.shutdown(); + }); + + it("waits for closing to finish before acquiring a replacement server", async () => { + let finishTermination: (() => void) | undefined; + const termination = new Promise((resolve) => { finishTermination = resolve; }); + const terminateTree = vi.fn().mockImplementationOnce(() => termination).mockResolvedValue(undefined); + const { deps } = fakeDeps({ terminateTree }); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + await pool.acquire(key); + const closing = pool.close(key); + await vi.waitFor(() => expect(terminateTree).toHaveBeenCalledTimes(1)); + let acquired = false; + const replacement = pool.acquire(key).then((entry) => { acquired = true; return entry; }); + await Promise.resolve(); + expect(acquired).toBe(false); + expect(deps.spawn).toHaveBeenCalledTimes(1); + finishTermination?.(); + await closing; + const entry = await replacement; + expect(entry.ready).toBe(true); + expect(entry.refs).toBe(1); + expect(deps.spawn).toHaveBeenCalledTimes(2); + await pool.shutdown(); + }); + + it("keeps an active server owned when its shell wrapper exits", async () => { + const { deps, listeners } = fakeDeps(); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + await pool.acquire(key); + for (const listener of listeners.get("exit") ?? []) listener(null, null); + expect(pool.size).toBe(1); + expect(pool.entryFor(key)?.refs).toBe(1); + await pool.shutdown(); + expect(pool.size).toBe(0); + }); + + it("waits for a pending start and refuses new work during shutdown", async () => { + let finishHealth: (() => void) | undefined; + const health = new Promise((resolve) => { finishHealth = resolve; }); + const { deps } = fakeDeps({ waitForHealth: vi.fn(() => health) }); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + const acquire = pool.acquire(key); + await vi.waitFor(() => expect(deps.spawn).toHaveBeenCalledTimes(1)); + const shutdown = pool.shutdown(); + await expect(pool.acquire({ ...key, cwd: "/w/b" })).rejects.toThrow("shutting down"); + expect(pool.size).toBe(1); + finishHealth?.(); + await expect(acquire).rejects.toThrow("shutting down"); + await shutdown; + expect(pool.size).toBe(0); + expect(deps.terminateTree).toHaveBeenCalledWith(4242); + }); + + it("does not spawn after shutdown overtakes free-port allocation", async () => { + let finishPort: ((port: number) => void) | undefined; + const port = new Promise((resolve) => { finishPort = resolve; }); + const { deps } = fakeDeps({ findFreePort: vi.fn(() => port) }); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + const acquire = pool.acquire(key); + const shutdown = pool.shutdown(); + finishPort?.(4096); + await expect(acquire).rejects.toThrow("shutting down"); + await shutdown; + expect(deps.spawn).not.toHaveBeenCalled(); + expect(pool.size).toBe(0); + }); + + it("owns the shell and descendants in a scope before declaring ready", async () => { + const events: string[] = []; + const scope = { + ready: true, + assign: vi.fn((_pid: number) => { events.push("assign"); return { ok: true }; }), + reconcile: vi.fn(async (_pid: number) => { events.push("reconcile"); return { ok: true }; }), + terminate: vi.fn(() => { events.push("terminate"); return { ok: true }; }), + waitForEmpty: vi.fn(async () => { events.push("empty"); return { ok: true }; }), + close: vi.fn(() => { events.push("scope-close"); }), + }; + const { deps } = fakeDeps({ + createScope: () => scope, + waitForHealth: vi.fn(async () => { events.push("health"); }), + waitForPortClosed: vi.fn(async () => { events.push("port-closed"); }), + }); + const pool = new OpenCodeServerPool(deps); + await pool.acquire({ binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }); + await pool.shutdown(); + expect(events).toEqual(["assign", "reconcile", "health", "terminate", "empty", "port-closed", "scope-close"]); + expect(deps.terminateTree).not.toHaveBeenCalled(); + }); + + it("does not spawn when its Windows process scope is unavailable", async () => { + const close = vi.fn(); + const { deps } = fakeDeps({ createScope: () => ({ ready: false, close }) }); + const pool = new OpenCodeServerPool(deps); + await expect(pool.acquire({ binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" })) + .rejects.toThrow("process scope unavailable"); + expect(deps.spawn).not.toHaveBeenCalled(); + expect(close).toHaveBeenCalledTimes(1); + await pool.shutdown(); + }); + + it("retains ownership when the port remains open after termination", async () => { + const waitForPortClosed = vi.fn(async () => { throw new Error("port 4096 remained open"); }); + const { deps } = fakeDeps({ waitForPortClosed }); + const pool = new OpenCodeServerPool(deps); + const key = { binaryPath: "opencode", cwd: "/w/a", hostname: "127.0.0.1" }; + await pool.acquire(key); + await expect(pool.close(key)).rejects.toThrow("port 4096 remained open"); + expect(pool.size).toBe(1); + waitForPortClosed.mockImplementation(async () => {}); + await pool.close(key); + expect(pool.size).toBe(0); + await pool.shutdown(); + }); + it("rejects the acquire when the child fails to spawn instead of crashing", async () => { const { deps, listeners } = fakeDeps({ waitForHealth: vi.fn(async () => { diff --git a/apps/server/src/features/providers/adapters/opencode/opencode-provider.ts b/apps/server/src/features/providers/adapters/opencode/opencode-provider.ts index d353571e0..623e5935f 100644 --- a/apps/server/src/features/providers/adapters/opencode/opencode-provider.ts +++ b/apps/server/src/features/providers/adapters/opencode/opencode-provider.ts @@ -356,7 +356,7 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP this.turns.delete(sessionId); } - shutdown(): void { + async shutdown(): Promise { for (const [sessionId, state] of this.turns) { state.aborted = true; this.drainPendingForSession(sessionId); @@ -364,9 +364,7 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP } this.drainAllPending(); this.turns.clear(); - void this.pool.shutdown().catch((err: unknown) => { - logger.warn("OpenCode pool shutdown failed", { error: String(err) }); - }); + await this.pool.shutdown(); } /** diff --git a/apps/server/src/features/providers/adapters/opencode/opencode-server-pool.ts b/apps/server/src/features/providers/adapters/opencode/opencode-server-pool.ts index 890b34500..e75993669 100644 --- a/apps/server/src/features/providers/adapters/opencode/opencode-server-pool.ts +++ b/apps/server/src/features/providers/adapters/opencode/opencode-server-pool.ts @@ -1,6 +1,8 @@ import * as NodeChildProcess from "node:child_process"; import * as NodeNet from "node:net"; import { logger } from "@mcode/shared"; +import { hostRuntime } from "@mcode/shared/node/host-runtime"; +import { WindowsProcessScopeFactory, type WindowsProcessScopeResult } from "../../../../runtime/process/containment/windows-process-scope.js"; /** Isolation boundary: never share a server across working directories. */ export type OpenCodePoolKey = Readonly<{ binaryPath: string; cwd: string; hostname: string }>; @@ -32,6 +34,15 @@ export interface OpenCodePoolProcess { } /** Injectable pool seams: spawning, health, termination, ports, and time. */ +export interface OpenCodeProcessScope { + readonly ready: boolean; + assign(pid: number): WindowsProcessScopeResult; + reconcile(pid: number): Promise; + terminate(exitCode?: number): WindowsProcessScopeResult; + waitForEmpty(timeoutMs: number): Promise; + close(): void; +} + export interface OpenCodePoolDeps { spawn(binaryPath: string, args: string[], cwd: string, env: Record): OpenCodePoolProcess; waitForHealth(baseUrl: string, timeoutMs: number, signal: AbortSignal): Promise; @@ -39,12 +50,45 @@ export interface OpenCodePoolDeps { findFreePort(hostname: string): Promise; now(): number; env(): Record; + createScope?(): OpenCodeProcessScope | null; + waitForPortClosed?(hostname: string, port: number): Promise; +} + +interface OwnedServer { + readonly entry: OpenCodePoolEntry; + readonly child: OpenCodePoolProcess; + readonly scope: OpenCodeProcessScope | null; } /** Idle time after the last release before an unreferenced server is closed. */ export const OPENCODE_POOL_IDLE_TTL_MS = 5 * 60 * 1_000; /** Longest wait for a fresh serve to answer health before startup fails. */ export const OPENCODE_POOL_READY_TIMEOUT_MS = 20_000; +const OPENCODE_POOL_CLOSE_TIMEOUT_MS = 5_000; + +async function defaultWaitForPortClosed(hostname: string, port: number): Promise { + const host = hostname === "0.0.0.0" ? "127.0.0.1" : hostname; + const deadline = Date.now() + OPENCODE_POOL_CLOSE_TIMEOUT_MS; + do { + const closed = await new Promise((resolve, reject) => { + const socket = NodeNet.createConnection({ host, port }); + socket.setTimeout(250); + socket.once("connect", () => { socket.destroy(); resolve(false); }); + socket.once("error", (error: NodeJS.ErrnoException) => { + socket.destroy(); + if (error.code === "ECONNREFUSED") resolve(true); + else reject(error); + }); + socket.once("timeout", () => { + socket.destroy(); + reject(new Error(`Timed out checking OpenCode serve port ${port}`)); + }); + }); + if (closed) return; + await new Promise((resolve) => setTimeout(resolve, 50)); + } while (Date.now() < deadline); + throw new Error(`OpenCode serve port ${port} remained open after termination`); +} async function defaultWaitForHealth(baseUrl: string, timeoutMs: number, signal: AbortSignal): Promise { const deadline = Date.now() + timeoutMs; @@ -52,9 +96,10 @@ async function defaultWaitForHealth(baseUrl: string, timeoutMs: number, signal: for (;;) { if (signal.aborted) throw new Error("OpenCode serve startup aborted"); try { - const res = await fetch(`${baseUrl}/global/health`); + const res = await fetch(`${baseUrl}/global/health`, { signal }); if (res.ok) return; } catch { + if (signal.aborted) throw new Error("OpenCode serve startup aborted"); // Not ready yet; keep polling below the deadline. } if (Date.now() >= deadline) throw new Error(`OpenCode serve did not become ready at ${baseUrl}`); @@ -88,10 +133,13 @@ async function defaultFindFreePort(hostname: string): Promise { * watches for unexpected exits, and proves process-tree termination on close. */ export class OpenCodeServerPool { - private readonly entries = new Map(); - private readonly children = new Map(); + private readonly servers = new Map(); private readonly pending = new Map>(); + private readonly starting = new Map(); + private readonly closing = new Map>(); private evictionTimer: ReturnType | null = null; + private stopping = false; + private shutdownPromise: Promise | null = null; constructor(private readonly deps: OpenCodePoolDeps) {} @@ -108,6 +156,10 @@ export class OpenCodeServerPool { }), waitForHealth: defaultWaitForHealth, findFreePort: defaultFindFreePort, + createScope: platform === "win32" + ? () => new WindowsProcessScopeFactory({ platform: "win32", architecture: hostRuntime.architecture }).create() + : () => null, + waitForPortClosed: defaultWaitForPortClosed, now: () => Date.now(), env: () => ({ ...process.env }) as Record, ...rest, @@ -115,46 +167,54 @@ export class OpenCodeServerPool { } get size(): number { - return this.entries.size; + return this.servers.size; } entryFor(key: OpenCodePoolKey): OpenCodePoolEntry | undefined { - return this.entries.get(openCodePoolKeyText(key)); + return this.servers.get(openCodePoolKeyText(key))?.entry; } /** Acquire (or spawn) the server for one working directory. Shares across threads. */ async acquire(key: OpenCodePoolKey): Promise { + if (this.stopping) throw new Error("OpenCode serve pool is shutting down"); this.ensureEvictionTimer(); const text = openCodePoolKeyText(key); - const existing = this.entries.get(text); - if (existing) { - existing.refs += 1; - existing.lastUsedAt = this.deps.now(); - return existing; - } - const inFlight = this.pending.get(text); - if (inFlight) { - const entry = await inFlight; - entry.refs += 1; - entry.lastUsedAt = this.deps.now(); - return entry; + for (;;) { + if (this.stopping) throw new Error("OpenCode serve pool is shutting down"); + const closing = this.closing.get(text); + if (closing) { + await closing; + continue; + } + const inFlight = this.pending.get(text); + if (inFlight) { + await this.waitForStart(text, inFlight); + continue; + } + const existing = this.servers.get(text)?.entry; + if (existing) { + existing.refs += 1; + existing.lastUsedAt = this.deps.now(); + return existing; + } + const started = this.start(key); + this.pending.set(text, started); + await this.waitForStart(text, started); } - const started = this.start(key); - this.pending.set(text, started); + } + + private async waitForStart(text: string, started: Promise): Promise { try { - const entry = await started; - entry.refs += 1; - entry.lastUsedAt = this.deps.now(); - return entry; + await started; } finally { - this.pending.delete(text); + if (this.pending.get(text) === started) this.pending.delete(text); } } /** Release one reference; the entry stays warm until the TTL or shutdown. */ release(key: OpenCodePoolKey): void { const text = openCodePoolKeyText(key); - const entry = this.entries.get(text); + const entry = this.servers.get(text)?.entry; if (!entry) return; entry.refs = Math.max(0, entry.refs - 1); entry.lastUsedAt = this.deps.now(); @@ -163,79 +223,148 @@ export class OpenCodeServerPool { /** Close one entry with proven tree termination, regardless of ref count. */ async close(key: OpenCodePoolKey): Promise { const text = openCodePoolKeyText(key); - const entry = this.entries.get(text); - if (!entry) return; - this.entries.delete(text); - const child = this.children.get(text); - this.children.delete(text); - await this.terminateEntry(entry, child); + const existing = this.closing.get(text); + if (existing) return existing; + const closing = this.closeServer(text); + this.closing.set(text, closing); + try { + await closing; + } finally { + this.closing.delete(text); + } + } + + private async closeServer(text: string): Promise { + const pending = this.pending.get(text); + if (pending) await Promise.allSettled([pending]); + const server = this.servers.get(text); + if (!server) return; + await this.terminateEntry(server); + this.servers.delete(text); } /** Close idle (refs at zero past TTL) entries with proven termination. */ async closeIdle(now = this.deps.now(), ttlMs = OPENCODE_POOL_IDLE_TTL_MS): Promise { const closed: string[] = []; - for (const [text, entry] of this.entries) { + for (const [text, { entry }] of this.servers) { if (entry.refs > 0) continue; if (now - entry.lastUsedAt <= ttlMs) continue; - this.entries.delete(text); - const child = this.children.get(text); - this.children.delete(text); - await this.terminateEntry(entry, child); + await this.close(entry.key); closed.push(text); } return closed; } async shutdown(): Promise { + if (this.shutdownPromise) return this.shutdownPromise; + this.stopping = true; if (this.evictionTimer) { clearInterval(this.evictionTimer); this.evictionTimer = null; } - await Promise.all([...this.entries.keys()].map(async (text) => { - const entry = this.entries.get(text); - if (!entry) return; - this.entries.delete(text); - const child = this.children.get(text); - this.children.delete(text); - await this.terminateEntry(entry, child); - })); + for (const controller of this.starting.values()) controller.abort(); + this.shutdownPromise = this.finishShutdown().catch((error: unknown) => { + this.shutdownPromise = null; + throw error; + }); + return this.shutdownPromise; + } + + private async finishShutdown(): Promise { + await Promise.allSettled(this.pending.values()); + const closures = await Promise.allSettled([...this.servers.values()].map(({ entry }) => this.close(entry.key))); + const failures = closures.filter((result): result is PromiseRejectedResult => result.status === "rejected") + .map((result) => result.reason); + if (failures.length > 0) throw new AggregateError(failures, "OpenCode serve pool shutdown failed"); } private async start(key: OpenCodePoolKey): Promise { const port = await this.deps.findFreePort(key.hostname); + if (this.stopping) throw new Error("OpenCode serve pool is shutting down"); const baseUrl = `http://${key.hostname}:${port}`; - const child = this.deps.spawn(key.binaryPath, ["serve", "--port", String(port), "--hostname", key.hostname], key.cwd, this.deps.env()); + const scope = this.createScope(); + const child = this.spawnChild(key, port, scope); logger.info("OpenCode serve spawn", { binaryPath: key.binaryPath, cwd: key.cwd, port }); const text = openCodePoolKeyText(key); const entry: OpenCodePoolEntry = { key, port, baseUrl, pid: child.pid ?? null, refs: 0, lastUsedAt: this.deps.now(), ready: false, }; - this.entries.set(text, entry); - this.children.set(text, child); + const server: OwnedServer = { entry, child, scope }; + this.servers.set(text, server); + const controller = new AbortController(); + this.starting.set(text, controller); + const spawnFailure: { error: Error | null } = { error: null }; + const onSpawnError = (error: Error) => { spawnFailure.error = error; }; + child.on("error", onSpawnError); const onExit = () => { - // Unexpected exits clean up without orphan processes; late waiters fail fast. - if (this.entries.get(text) === entry) this.entries.delete(text); - this.children.delete(text); + // The shell can exit before its server. Keep ownership until cleanup is proven. + if (entry.ready && entry.refs === 0) { + void this.close(key).catch((error: unknown) => logger.error("OpenCode serve exit cleanup failed", { + port, error: error instanceof Error ? error.message : String(error), + })); + } }; child.on("exit", onExit); - const controller = new AbortController(); const guard = setTimeout(() => controller.abort(), OPENCODE_POOL_READY_TIMEOUT_MS + 5_000); try { + await this.establishScope(server); + if (spawnFailure.error) throw new Error(`OpenCode serve failed to start: ${spawnFailure.error.message}`); await this.awaitServeReady(child, baseUrl, controller.signal); } catch (error) { child.off("exit", onExit); - this.entries.delete(text); - this.children.delete(text); - await this.terminateEntry(entry, child); + await this.cleanupFailedStart(text, server, error); throw error; } finally { clearTimeout(guard); + this.starting.delete(text); + child.off("error", onSpawnError); + } + if (this.stopping) { + await this.terminateEntry(server); + this.servers.delete(text); + throw new Error("OpenCode serve pool is shutting down"); } entry.ready = true; entry.lastUsedAt = this.deps.now(); return entry; } + private createScope(): OpenCodeProcessScope | null { + const scope = this.deps.createScope?.() ?? null; + if (scope && !scope.ready) { + scope.close(); + throw new Error("OpenCode serve process scope unavailable"); + } + return scope; + } + + private spawnChild(key: OpenCodePoolKey, port: number, scope: OpenCodeProcessScope | null): OpenCodePoolProcess { + try { + return this.deps.spawn(key.binaryPath, ["serve", "--port", String(port), "--hostname", key.hostname], key.cwd, this.deps.env()); + } catch (error) { + scope?.close(); + throw error; + } + } + + private async establishScope({ entry, scope }: OwnedServer): Promise { + if (!scope) return; + if (entry.pid == null) throw new Error("OpenCode serve process PID unavailable"); + const assigned = scope.assign(entry.pid); + if (!assigned.ok) throw new Error(`OpenCode serve process scope assignment failed: ${assigned.error}`); + const reconciled = await scope.reconcile(entry.pid); + if (!reconciled.ok) throw new Error(`OpenCode serve process scope reconciliation failed: ${reconciled.error}`); + } + + private async cleanupFailedStart(text: string, server: OwnedServer, startupError: unknown): Promise { + try { + await this.terminateEntry(server); + this.servers.delete(text); + } catch (cleanupError) { + throw new AggregateError([startupError, cleanupError], "OpenCode serve startup and cleanup failed"); + } + } + /** * Wait for the serve health endpoint while also watching the child for an * early `error` (e.g. missing binary). Without the `error` listener a spawn @@ -264,32 +393,36 @@ export class OpenCodeServerPool { }); } - private async terminateEntry(entry: OpenCodePoolEntry, child: OpenCodePoolProcess | undefined): Promise { - // Proven tree termination goes first. A best-effort kill of the direct - // child beforehand would orphan grandchildren on Windows (the tracked pid - // is the shell wrapper), making the later snapshot-based verification - // vacuously succeed while the real server keeps running. - if (entry.pid != null) { - try { - await this.deps.terminateTree(entry.pid); - return; - } catch (error) { - logger.warn("OpenCode serve tree termination failed; trying direct kill", { - pid: entry.pid, - error: error instanceof Error ? error.message : String(error), - }); - } + private async terminateEntry({ entry, child, scope }: OwnedServer): Promise { + if (scope) { + await this.terminateScope(scope, entry, child); + } else if (entry.pid != null) { + await this.deps.terminateTree(entry.pid); + } else if (!child.kill("SIGTERM")) { + throw new Error("OpenCode serve child could not be terminated"); } - try { - child?.kill("SIGTERM"); - } catch { - // Best effort; an exited process is not an orphan. + await this.deps.waitForPortClosed?.(entry.key.hostname, entry.port); + scope?.close(); + } + + private async terminateScope(scope: OpenCodeProcessScope, entry: OpenCodePoolEntry, child: OpenCodePoolProcess): Promise { + const terminated = scope.terminate(0); + if (!terminated.ok) { + if (entry.pid != null) await this.deps.terminateTree(entry.pid); + else if (!child.kill("SIGTERM")) throw new Error(terminated.error ?? "OpenCode serve scope termination failed"); + return; } + const emptied = await scope.waitForEmpty(OPENCODE_POOL_CLOSE_TIMEOUT_MS); + if (!emptied.ok) throw new Error(emptied.error ?? "OpenCode serve process scope remained non-empty"); } private ensureEvictionTimer(): void { if (this.evictionTimer) return; - this.evictionTimer = setInterval(() => void this.closeIdle(), 60_000); + this.evictionTimer = setInterval(() => { + void this.closeIdle().catch((error: unknown) => logger.error("OpenCode serve idle cleanup failed", { + error: error instanceof Error ? error.message : String(error), + })); + }, 60_000); this.evictionTimer.unref?.(); } } diff --git a/apps/server/src/features/providers/models/__tests__/model-cache-service.test.ts b/apps/server/src/features/providers/models/__tests__/model-cache-service.test.ts index 6cfb7a2de..91b2b0fbc 100644 --- a/apps/server/src/features/providers/models/__tests__/model-cache-service.test.ts +++ b/apps/server/src/features/providers/models/__tests__/model-cache-service.test.ts @@ -14,12 +14,19 @@ vi.mock("../../../../application/transport/push.js", () => ({ import type { Database } from "bun:sqlite"; import { openMemoryDatabase } from "../../../../runtime/persistence/sqlite/database.js"; import { ModelCacheRepo } from "../persistence/model-cache-repo.js"; -import { ModelCacheService } from "../model-cache-service.js"; -import type { ProviderModelInfo, IProviderRegistry } from "@mcode/contracts"; +import { ModelCacheService, startupModelProviderIds } from "../model-cache-service.js"; +import type { ProviderAvailability, ProviderId, ProviderModelInfo, IProviderRegistry } from "@mcode/contracts"; -function makeProvider(models: ProviderModelInfo[]) { +function availableProvider(id: ProviderId, enabled: boolean): ProviderAvailability { return { - id: "test-provider", + id, enabled, hasAdapter: true, beta: false, comingSoon: false, capabilities: [], + cli: { status: "found", resolvedPath: id, configuredPath: "" }, + }; +} + +function makeProvider(models: ProviderModelInfo[], id = "test-provider") { + return { + id, listModels: vi.fn().mockResolvedValue(models), sendTurn: vi.fn(), cancelSession: vi.fn(), @@ -55,6 +62,27 @@ describe("ModelCacheService", () => { db.close(); }); + it("warms enabled providers without starting idle OpenCode", async () => { + const claude = makeProvider([{ id: "claude-model", name: "Claude Model" }], "claude"); + const opencode = makeProvider([{ id: "open-model", name: "Open Model" }], "opencode"); + const codex = makeProvider([{ id: "codex-model", name: "Codex Model" }], "codex"); + const registry = makeRegistry(new Map([["claude", claude], ["opencode", opencode], ["codex", codex]])); + const service = new ModelCacheService(repo, registry); + + const startupProviders = startupModelProviderIds([ + availableProvider("claude", true), + availableProvider("opencode", true), + availableProvider("codex", false), + ]); + expect(startupProviders).toEqual(["claude"]); + await service.refreshProviders(startupProviders); + + expect(service.getCached("claude")).toEqual([{ id: "claude-model", name: "Claude Model" }]); + expect(service.getCached("opencode")).toBeUndefined(); + expect(opencode.listModels).not.toHaveBeenCalled(); + expect(codex.listModels).not.toHaveBeenCalled(); + }); + it("returns cached models without calling provider when cache is fresh", async () => { const models: ProviderModelInfo[] = [{ id: "m1", name: "Model 1" }]; repo.upsert("test-provider", models); diff --git a/apps/server/src/features/providers/models/model-cache-service.ts b/apps/server/src/features/providers/models/model-cache-service.ts index 7e2c38ac0..c5a5495c5 100644 --- a/apps/server/src/features/providers/models/model-cache-service.ts +++ b/apps/server/src/features/providers/models/model-cache-service.ts @@ -10,13 +10,20 @@ import { inject, injectable } from "tsyringe"; import { logger } from "@mcode/shared"; -import type { ProviderModelInfo, IProviderRegistry } from "@mcode/contracts"; +import type { ProviderAvailability, ProviderId, ProviderModelInfo, IProviderRegistry } from "@mcode/contracts"; import { broadcast } from "../../../application/transport/push.js"; import { ModelCacheRepo } from "./persistence/model-cache-repo.js"; /** How long a cached entry is considered "fresh" (no background refresh). */ const CACHE_FRESH_MS = 60 * 60 * 1000; // 1 hour +/** Select providers whose model lists can be warmed without starting OpenCode. */ +export function startupModelProviderIds(providers: readonly ProviderAvailability[]): ProviderId[] { + return providers + .filter((provider) => provider.id !== "opencode" && provider.enabled && provider.hasAdapter && !provider.comingSoon && provider.cli.status !== "not_found") + .map((provider) => provider.id); +} + /** * Parse a SQLite `datetime('now')` string as UTC. SQLite returns * `YYYY-MM-DD HH:MM:SS` in UTC with no timezone marker; the JS `Date` @@ -147,15 +154,11 @@ export class ModelCacheService { } } - /** - * Refreshes all providers that support model listing. - * Called on WS connect to ensure the cache stays warm. - */ - async refreshAll(): Promise { - const providers = this.registry.resolveAll(); - const promises = providers.map((p) => - this.refreshProvider(p.id).catch((err) => { - logger.warn("Model refresh failed", { providerId: p.id, err: String(err) }); + /** Refresh model lists for the providers selected by the caller. */ + async refreshProviders(providerIds: readonly ProviderId[]): Promise { + const promises = providerIds.map((providerId) => + this.refreshProvider(providerId).catch((err) => { + logger.warn("Model refresh failed", { providerId, err: String(err) }); }), ); await Promise.allSettled(promises);