diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 8da6ffe4357d..e1acd044eb2f 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -1,4 +1,9 @@ import { revertCodexThread } from "../../provider/CodexThreadRevert.ts"; +import { + noteCodexMcpStartup, + waitForCodexMcpCatalogBeforeTurn, + type CodexMcpStartupCatalog, +} from "../../provider/CodexMcpCatalog.ts"; import { historyResponseItems } from "../ContextHandoffBudget.ts"; import { makeProviderTextDeltaCoalescer } from "./ProviderTextDeltaCoalescer.ts"; import { @@ -1664,6 +1669,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ).hasSubagents = true; }; const pendingRootTurns = yield* Ref.make(new Map()); + const mcpStartupByThread = yield* Ref.make(new Map()); const turnWaiters = yield* Ref.make(new Map>()); const subagentThreads = yield* Ref.make(new Map()); const subagentModels = new Map(); @@ -3708,6 +3714,16 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi return { node, request, turnItem }; }); + yield* client.handleServerNotification("mcpServer/startupStatus/updated", (payload) => + Ref.update(mcpStartupByThread, (catalog) => + noteCodexMcpStartup(catalog, { + threadId: payload.threadId ?? null, + name: payload.name, + status: payload.status, + }), + ).pipe(Effect.asVoid), + ); + yield* client.handleServerNotification("item/agentMessage/delta", (payload) => Effect.gen(function* () { const context = (yield* Ref.get(activeTurns)).get(payload.turnId); @@ -5547,6 +5563,13 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi startTurn: (turnInput) => Effect.gen(function* () { const threadId = yield* getNativeThreadId(turnInput.providerThread); + // turn/start snapshots the tool catalog. Wait out servers that + // have already reported `starting` so the snapshot is not limited + // to whichever stdio server finished first. + yield* waitForCodexMcpCatalogBeforeTurn({ + catalog: mcpStartupByThread, + threadId, + }); const codexInput = turnInput.restartContinuationOfRunId === undefined diff --git a/apps/server/src/provider/CodexMcpCatalog.test.ts b/apps/server/src/provider/CodexMcpCatalog.test.ts new file mode 100644 index 000000000000..f6e7e1197253 --- /dev/null +++ b/apps/server/src/provider/CodexMcpCatalog.test.ts @@ -0,0 +1,181 @@ +import { assert, describe, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Ref from "effect/Ref"; +import { TestClock } from "effect/testing"; + +import { + codexMcpCatalogSnapshot, + codexMcpStartupStatusesForThread, + noteCodexMcpStartup, + reduceCodexMcpStartupStatus, + waitForCodexMcpCatalogBeforeTurn, + type CodexMcpStartupCatalog, +} from "./CodexMcpCatalog.ts"; + +const emptyCatalog: CodexMcpStartupCatalog = new Map(); + +/** Record one startup phase in a catalog under test. */ +function note( + catalog: CodexMcpStartupCatalog, + threadId: string | null, + name: string, + status: "starting" | "ready" | "failed" | "cancelled", +) { + return noteCodexMcpStartup(catalog, { threadId, name, status }); +} + +describe("Codex MCP catalog startup", () => { + it("lets a later ready replace a spurious cancelled status", () => { + const cancelled = reduceCodexMcpStartupStatus(new Map(), { + name: "slow", + status: "cancelled", + }); + const ready = reduceCodexMcpStartupStatus(cancelled, { name: "slow", status: "ready" }); + assert.equal(ready.get("slow"), "ready"); + + const afterReady = reduceCodexMcpStartupStatus(ready, { name: "slow", status: "cancelled" }); + assert.equal(afterReady.get("slow"), "ready"); + assert.equal(codexMcpCatalogSnapshot(afterReady), "settled"); + }); + + it("keeps a partial catalog pending while any server is still starting", () => { + const catalog = note( + note(emptyCatalog, "thread-1", "alpha", "ready"), + "thread-1", + "beta", + "starting", + ); + assert.equal( + codexMcpCatalogSnapshot(codexMcpStartupStatusesForThread(catalog, "thread-1")), + "pending", + ); + assert.equal( + codexMcpCatalogSnapshot(codexMcpStartupStatusesForThread(catalog, "thread-2")), + "idle", + ); + }); + + it.effect("does not delay a turn when no MCP server has reported startup", () => + Effect.gen(function* () { + const catalog = yield* Ref.make(emptyCatalog); + const outcome = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + }); + assert.equal(outcome, "idle"); + }), + ); + + it.effect("holds the turn until a slower MCP server leaves starting", () => + Effect.gen(function* () { + const catalog = yield* Ref.make( + note(note(emptyCatalog, "thread-1", "alpha", "ready"), "thread-1", "beta", "starting"), + ); + let turnStarted = false; + const fiber = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + timeout: "30 seconds", + pollInterval: "50 millis", + }).pipe( + Effect.tap(() => + Effect.sync(() => { + turnStarted = true; + }), + ), + Effect.forkChild, + ); + + yield* TestClock.adjust("2 seconds"); + assert.equal(turnStarted, false); + assert.equal(fiber.pollUnsafe(), undefined); + + yield* Ref.update(catalog, (current) => note(current, "thread-1", "beta", "ready")); + yield* TestClock.adjust("50 millis"); + assert.equal(yield* Fiber.join(fiber), "settled"); + assert.equal(turnStarted, true); + }), + ); + + it.effect("keeps waiting through a spurious cancelled status until ready", () => + Effect.gen(function* () { + const catalog = yield* Ref.make( + note(note(emptyCatalog, "thread-1", "alpha", "ready"), "thread-1", "beta", "cancelled"), + ); + const fiber = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + timeout: "30 seconds", + pollInterval: "50 millis", + }).pipe(Effect.forkChild); + + yield* TestClock.adjust("2 seconds"); + assert.equal(fiber.pollUnsafe(), undefined); + + yield* Ref.update(catalog, (current) => note(current, "thread-1", "beta", "ready")); + yield* TestClock.adjust("50 millis"); + assert.equal(yield* Fiber.join(fiber), "settled"); + }), + ); + + it.effect("starts the turn when a server stays starting through the timeout", () => + Effect.gen(function* () { + const catalog = yield* Ref.make(note(emptyCatalog, "thread-1", "hung", "starting")); + const fiber = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + timeout: "200 millis", + pollInterval: "50 millis", + }).pipe(Effect.forkChild); + yield* TestClock.adjust("100 millis"); + assert.equal(fiber.pollUnsafe(), undefined); + yield* TestClock.adjust("100 millis"); + assert.equal(yield* Fiber.join(fiber), "timedOut"); + + const again = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + timeout: "200 millis", + pollInterval: "50 millis", + }); + assert.equal(again, "settled"); + assert.equal( + codexMcpStartupStatusesForThread(yield* Ref.get(catalog), "thread-1").get("hung"), + "unavailable", + ); + }), + ); + + it.effect("applies a threadless startup update to the turn's thread", () => + Effect.gen(function* () { + const catalog = yield* Ref.make(note(emptyCatalog, null, "shared", "starting")); + const fiber = yield* waitForCodexMcpCatalogBeforeTurn({ + catalog, + threadId: "thread-1", + timeout: "5 seconds", + pollInterval: "50 millis", + }).pipe(Effect.forkChild); + + yield* TestClock.adjust("100 millis"); + assert.equal(fiber.pollUnsafe(), undefined); + + yield* Ref.update(catalog, (current) => note(current, null, "shared", "ready")); + yield* TestClock.adjust("50 millis"); + assert.equal(yield* Fiber.join(fiber), "settled"); + }), + ); + + it("treats a failed server as settled alongside servers that are ready", () => { + const catalog = note( + note(emptyCatalog, "thread-1", "alpha", "ready"), + "thread-1", + "broken", + "failed", + ); + assert.equal( + codexMcpCatalogSnapshot(codexMcpStartupStatusesForThread(catalog, "thread-1")), + "settled", + ); + }); +}); diff --git a/apps/server/src/provider/CodexMcpCatalog.ts b/apps/server/src/provider/CodexMcpCatalog.ts new file mode 100644 index 000000000000..d55dc2e649bb --- /dev/null +++ b/apps/server/src/provider/CodexMcpCatalog.ts @@ -0,0 +1,202 @@ +import * as Clock from "effect/Clock"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Ref from "effect/Ref"; + +/** + * Per-server MCP startup phase observed from `mcpServer/startupStatus/updated`. + * + * `unavailable` is local: a bounded wait gave up on `starting` or `cancelled` + * so a later turn is not held for a server that never finished. + */ +export type CodexMcpStartupPhase = "starting" | "ready" | "failed" | "cancelled" | "unavailable"; + +/** Startup phases for one Codex thread, keyed by MCP server name. */ +export type CodexMcpStartupStatuses = ReadonlyMap; + +/** + * Startup phases for every thread in one app-server session. + * Notifications that omit `threadId` are stored under {@link CODEX_MCP_GLOBAL_THREAD_KEY}. + */ +export type CodexMcpStartupCatalog = ReadonlyMap; + +/** Bucket for startup notifications that are not scoped to a thread. */ +export const CODEX_MCP_GLOBAL_THREAD_KEY = "*"; + +/** + * How long a turn waits for MCP servers to leave `starting` before Codex + * snapshots the tool catalog. Long enough for the slow stdio servers in + * the partial-catalog reports, short enough that a hung server still yields. + */ +export const CODEX_MCP_CATALOG_STARTUP_TIMEOUT: Duration.Input = "30 seconds"; + +/** How often a waiting turn re-reads startup phases. */ +export const CODEX_MCP_CATALOG_POLL_INTERVAL: Duration.Input = "50 millis"; + +const EMPTY_MCP_STARTUP_STATUSES: CodexMcpStartupStatuses = new Map(); + +export type CodexMcpCatalogWait = "idle" | "settled" | "timedOut"; + +/** + * Fold one startup status into a server map. + * + * A later `ready` replaces `cancelled`. Codex can emit a spurious `cancelled` + * before `ready` for a server that was started once. `cancelled` does not + * replace `ready`. A fresh `starting` opens another round. + */ +export function reduceCodexMcpStartupStatus( + statuses: CodexMcpStartupStatuses, + update: { readonly name: string; readonly status: CodexMcpStartupPhase }, +): CodexMcpStartupStatuses { + const name = update.name.trim(); + if (name.length === 0) return statuses; + const current = statuses.get(name); + if (update.status === "cancelled" && current === "ready") return statuses; + if (current === update.status) return statuses; + const next = new Map(statuses); + next.set(name, update.status); + return next; +} + +/** + * Record one `mcpServer/startupStatus/updated` notification on its thread. + * Returns the same catalog when the phase does not change. + */ +export function noteCodexMcpStartup( + catalog: CodexMcpStartupCatalog, + update: { + readonly threadId: string | null; + readonly name: string; + readonly status: CodexMcpStartupPhase; + }, +): CodexMcpStartupCatalog { + const name = update.name.trim(); + if (name.length === 0) return catalog; + const threadKey = + update.threadId === null || update.threadId.length === 0 + ? CODEX_MCP_GLOBAL_THREAD_KEY + : update.threadId; + const current = catalog.get(threadKey) ?? EMPTY_MCP_STARTUP_STATUSES; + const reduced = reduceCodexMcpStartupStatus(current, { name, status: update.status }); + if (reduced === current) return catalog; + const next = new Map(catalog); + next.set(threadKey, reduced); + return next; +} + +/** + * Phases that apply to `threadId`, with that thread's own updates winning + * over notifications that carried no thread id. + */ +export function codexMcpStartupStatusesForThread( + catalog: CodexMcpStartupCatalog, + threadId: string, +): CodexMcpStartupStatuses { + const globalStatuses = catalog.get(CODEX_MCP_GLOBAL_THREAD_KEY); + const threadStatuses = catalog.get(threadId); + if (globalStatuses === undefined) return threadStatuses ?? EMPTY_MCP_STARTUP_STATUSES; + if (threadStatuses === undefined) return globalStatuses; + const merged = new Map(globalStatuses); + for (const [name, phase] of threadStatuses) merged.set(name, phase); + return merged; +} + +/** + * `pending` while any server is `starting` or `cancelled`. + * `cancelled` stays pending so a later `ready` can still win. + * `idle` means this thread has not reported any MCP server yet. + */ +export function codexMcpCatalogSnapshot( + statuses: CodexMcpStartupStatuses, +): "idle" | "pending" | "settled" { + if (statuses.size === 0) return "idle"; + for (const phase of statuses.values()) { + if (phase === "starting" || phase === "cancelled") return "pending"; + } + return "settled"; +} + +/** + * Stop blocking on servers that were still `starting` or `cancelled` when + * the wait budget ran out. A later `starting` or `ready` replaces this. + */ +export function releaseTimedOutCodexMcpStartup( + catalog: CodexMcpStartupCatalog, + threadId: string, +): CodexMcpStartupCatalog { + const statuses = codexMcpStartupStatusesForThread(catalog, threadId); + let next = catalog; + for (const [name, phase] of statuses) { + if (phase === "starting" || phase === "cancelled") { + next = noteCodexMcpStartup(next, { threadId, name, status: "unavailable" }); + } + } + return next; +} + +/** + * Resolve when the thread's MCP startup phases are safe to snapshot. + * + * Returns immediately when no server is `starting` or `cancelled`, so a turn + * with a finished catalog does not sleep. On timeout the caller still starts + * the turn; the servers that were pending are marked `unavailable` for this + * thread so the same hang does not consume the budget again. + */ +export const awaitCodexMcpCatalog = Effect.fn("awaitCodexMcpCatalog")(function* (input: { + readonly readStatuses: Effect.Effect; + readonly timeout?: Duration.Input; + readonly pollInterval?: Duration.Input; +}) { + const timeoutMillis = Duration.toMillis(input.timeout ?? CODEX_MCP_CATALOG_STARTUP_TIMEOUT); + const pollInterval = input.pollInterval ?? CODEX_MCP_CATALOG_POLL_INTERVAL; + const startedAt = yield* Clock.currentTimeMillis; + while (true) { + const statuses = yield* input.readStatuses; + const snapshot = codexMcpCatalogSnapshot(statuses); + if (snapshot !== "pending") return snapshot; + const elapsed = (yield* Clock.currentTimeMillis) - startedAt; + if (elapsed >= timeoutMillis) return "timedOut" as const; + yield* Effect.sleep(pollInterval); + } +}); + +/** + * Hold a Codex turn until every MCP server that has reported in for this + * thread has left `starting`. `turn/start` snapshots the model-facing tool + * catalog, and a server still starting is omitted from that snapshot. + */ +export const waitForCodexMcpCatalogBeforeTurn = Effect.fn("waitForCodexMcpCatalogBeforeTurn")( + function* (input: { + readonly catalog: Ref.Ref; + readonly threadId: string; + readonly timeout?: Duration.Input; + readonly pollInterval?: Duration.Input; + }) { + const readStatuses = Ref.get(input.catalog).pipe( + Effect.map((catalog) => codexMcpStartupStatusesForThread(catalog, input.threadId)), + ); + const initialSnapshot = codexMcpCatalogSnapshot(yield* readStatuses); + if (initialSnapshot !== "pending") return initialSnapshot; + + const outcome = yield* awaitCodexMcpCatalog({ + readStatuses, + ...(input.timeout === undefined ? {} : { timeout: input.timeout }), + ...(input.pollInterval === undefined ? {} : { pollInterval: input.pollInterval }), + }); + if (outcome !== "timedOut") return outcome; + + const pending = [ + ...codexMcpStartupStatusesForThread(yield* Ref.get(input.catalog), input.threadId), + ] + .filter(([, phase]) => phase === "starting" || phase === "cancelled") + .map(([name]) => name); + yield* Effect.logWarning( + "Codex MCP servers were still starting when the catalog wait ended; the turn will start with whatever is ready.", + { threadId: input.threadId, pending }, + ); + yield* Ref.update(input.catalog, (catalog) => + releaseTimedOutCodexMcpStartup(catalog, input.threadId), + ); + return outcome; + }, +);