Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 13 additions & 16 deletions apps/server/src/application/bootstrap/server-bootstrap.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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),
});
Expand Down Expand Up @@ -940,9 +937,18 @@ async function shutdown(): Promise<void> {
shutdownCoordinator.setPhase("stop agent sessions");
await agentService.stopAll();

let shutdownFailure: unknown = null;
const captureCleanupFailure = async (cleanup: () => Promise<void> | void): Promise<void> => {
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();
Expand All @@ -951,15 +957,6 @@ async function shutdown(): Promise<void> {
// 3. Dispose settings file watcher
settingsService.dispose();

let shutdownFailure: unknown = null;
const captureCleanupFailure = async (cleanup: () => Promise<void> | void): Promise<void> => {
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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -123,15 +123,15 @@ 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 () => {
const http = fakeHttp({ getSessionStatus: vi.fn(async () => ({})) });
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 () => {
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -251,14 +251,14 @@ 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 () => {
const http = fakeHttp({ subscribeEvents: openSubscribe([IDLE, IDLE, IDLE]) });
const { provider, submitted } = testProvider(http as never);
await provider.sendTurn(turnRequest());
expect(endedOutcomes(submitted)).toEqual(["completed"]);
provider.shutdown();
await provider.shutdown();
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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" },
Expand All @@ -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();
});
});

Expand All @@ -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 () => {
Expand All @@ -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();
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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();
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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[] = [];
Expand All @@ -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 },
});
Expand Down Expand Up @@ -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<void>((resolve) => { signalTerminationStarted = resolve; });
const terminationAllowed = new Promise<void>((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,
Expand Down Expand Up @@ -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 () => {
Expand Down Expand Up @@ -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 () => {
Expand Down Expand Up @@ -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 () => {
Expand Down Expand Up @@ -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();
});
});
Loading
Loading