From 45a5748641dadbc551cdc744b12f9633a58b8960 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Sat, 27 Jun 2026 11:54:11 -0700 Subject: [PATCH 001/132] test(ai): add failing umans usage provider suite --- packages/ai/test/umans-usage.test.ts | 165 +++++++++++++++++++++++++++ 1 file changed, 165 insertions(+) create mode 100644 packages/ai/test/umans-usage.test.ts diff --git a/packages/ai/test/umans-usage.test.ts b/packages/ai/test/umans-usage.test.ts new file mode 100644 index 000000000..e4010e206 --- /dev/null +++ b/packages/ai/test/umans-usage.test.ts @@ -0,0 +1,165 @@ +import { describe, expect, it } from "bun:test"; +import type { FetchImpl } from "@oh-my-pi/pi-ai/types"; +import { umansUsageProvider } from "../src/usage/umans"; + +const DEFAULT_BASE_URL = "https://api.code.umans.ai"; + +function umansPayload(overrides: Record = {}): Record { + return { + plan: { display_name: "Code Max" }, + limits: { + requests: { limit: 200, hard_cap: 400, burst_pct: 1.0, window_seconds: 18000 }, + concurrency: { limit: 4, hard_cap: 8, burst_pct: 1.0 }, + }, + usage: { + requests_in_window: 48, + remaining_requests: 152, + concurrent_sessions: 1, + tokens_in: 1_200_000, + tokens_out: 340_000, + priority: { low: false, boxed_until: null, reason: null }, + }, + ...overrides, + }; +} + +function fakeFetch(payload: unknown, status = 200): FetchImpl { + const fn = async () => + new Response(JSON.stringify(payload), { + status, + headers: { "content-type": "application/json" }, + }); + return fn as unknown as typeof fetch; +} + +function fetchRecorder(calls: Array<{ url: string; headers: Record }>, payload: unknown, status = 200): FetchImpl { + const fn = async (input: string | URL | Request, init?: RequestInit) => { + calls.push({ + url: String(input), + headers: (init?.headers as Record) ?? {}, + }); + return new Response(JSON.stringify(payload), { + status, + headers: { "content-type": "application/json" }, + }); + }; + return fn as unknown as typeof fetch; +} + +describe("umans usage provider", () => { + it("parses the rolling 5h request window into a UsageLimit with used/remaining/fraction", async () => { + const report = await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test", accountId: "acct-1", email: "u@example.com" }, + }, + { fetch: fakeFetch(umansPayload()) }, + ); + expect(report).not.toBeNull(); + const requests = report?.limits.find(l => l.id === "umans:requests"); + expect(requests).toBeDefined(); + expect(requests?.amount.used).toBe(48); + expect(requests?.amount.limit).toBe(200); + expect(requests?.amount.remaining).toBe(152); + expect(requests?.amount.usedFraction).toBeCloseTo(0.24, 5); + expect(requests?.amount.remainingFraction).toBeCloseTo(0.76, 5); + expect(requests?.amount.unit).toBe("requests"); + // Rolling window: no fabricated reset timestamp. + expect(requests?.window?.resetsAt).toBeUndefined(); + expect(requests?.window?.durationMs).toBe(18000_000); + expect(requests?.window?.label).toBe("rolling 5h"); + }); + + it("emits a concurrency limit from limits.concurrency", async () => { + const report = await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fakeFetch(umansPayload()) }, + ); + const concurrency = report?.limits.find(l => l.id === "umans:concurrency"); + expect(concurrency).toBeDefined(); + expect(concurrency?.amount.used).toBe(1); + expect(concurrency?.amount.limit).toBe(4); + expect(concurrency?.amount.unit).toBe("requests"); + }); + + it("sends Authorization: Bearer to the default base URL", async () => { + const calls: Array<{ url: string; headers: Record }> = []; + await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fetchRecorder(calls, umansPayload()) }, + ); + expect(calls).toHaveLength(1); + expect(calls[0]?.url).toBe(`${DEFAULT_BASE_URL}/v1/usage`); + expect(calls[0]?.headers.authorization).toBe("Bearer sk-test"); + }); + + it("honors a custom baseUrl from params", async () => { + const calls: Array<{ url: string; headers: Record }> = []; + await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + baseUrl: "https://custom.umans.example", + }, + { fetch: fetchRecorder(calls, umansPayload()) }, + ); + expect(calls[0]?.url).toBe("https://custom.umans.example/v1/usage"); + }); + + it("surfaces priority.low as a provider note", async () => { + const payload = umansPayload({ + usage: { + requests_in_window: 250, + remaining_requests: 0, + concurrent_sessions: 1, + tokens_in: 0, + tokens_out: 0, + priority: { low: true, boxed_until: "2026-06-27T12:00:00Z", reason: "burst" }, + }, + }); + const report = await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fakeFetch(payload) }, + ); + expect(report?.notes).toContain("Requests deprioritized after a rate-limit burst."); + }); + + it("returns null on a non-ok HTTP response", async () => { + const report = await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fakeFetch({ message: "unauthorized" }, 401) }, + ); + expect(report).toBeNull(); + }); + + it("returns null when supports() is called for a different provider or credential type", () => { + expect(umansUsageProvider.supports?.({ provider: "zai", credential: { type: "api_key", apiKey: "x" } })).toBe(false); + expect(umansUsageProvider.supports?.({ provider: "umans", credential: { type: "oauth", accessToken: "x" } })).toBe(false); + expect(umansUsageProvider.supports?.({ provider: "umans", credential: { type: "api_key", apiKey: "x" } })).toBe(true); + }); + + it("includes plan display name and account identity in metadata", async () => { + const report = await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test", accountId: "acct-42", email: "dev@example.com" }, + }, + { fetch: fakeFetch(umansPayload({ plan: { display_name: "Code Pro" } })) }, + ); + expect(report?.metadata?.plan).toBe("Code Pro"); + expect(report?.metadata?.accountId).toBe("acct-42"); + expect(report?.metadata?.email).toBe("dev@example.com"); + }); +}); From 209c93d6ed8f14cf5ccd4be1fc4a05b07def2f62 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Sat, 27 Jun 2026 11:54:16 -0700 Subject: [PATCH 002/132] feat(ai): add umans usage provider Fetches GET /v1/usage from api.code.umans.ai and maps the rolling 5h request window + instantaneous concurrency limit onto the shared UsageReport shape. Not yet wired into DEFAULT_USAGE_PROVIDERS. --- packages/ai/src/usage/umans.ts | 178 +++++++++++++++++++++++++++++++++ 1 file changed, 178 insertions(+) create mode 100644 packages/ai/src/usage/umans.ts diff --git a/packages/ai/src/usage/umans.ts b/packages/ai/src/usage/umans.ts new file mode 100644 index 000000000..bf8007b15 --- /dev/null +++ b/packages/ai/src/usage/umans.ts @@ -0,0 +1,178 @@ +import type { + UsageAmount, + UsageFetchContext, + UsageFetchParams, + UsageLimit, + UsageProvider, + UsageReport, + UsageStatus, + UsageWindow, +} from "../usage"; +import { isRecord } from "../utils"; + +const UMANS_PROVIDER = "umans"; +const DEFAULT_ENDPOINT = "https://api.code.umans.ai"; +const USAGE_PATH = "/v1/usage"; +const FIVE_HOURS_MS = 5 * 60 * 60 * 1000; + +/** Umans `GET /v1/usage` response (subset; extras ignored). */ +interface UmansUsagePayload { + plan?: { display_name?: string }; + limits?: { + requests?: { limit?: number; hard_cap?: number | null; window_seconds?: number }; + concurrency?: { limit?: number; hard_cap?: number | null }; + }; + usage?: { + requests_in_window?: number; + remaining_requests?: number; + concurrent_sessions?: number; + tokens_in?: number; + tokens_out?: number; + priority?: { low?: boolean }; + }; +} + +function normalizeBaseUrl(baseUrl?: string): string { + if (!baseUrl?.trim()) return DEFAULT_ENDPOINT; + const trimmed = baseUrl.trim(); + // Strip a trailing slash so `${origin}/v1/usage` is stable. + return trimmed.endsWith("/") ? trimmed.slice(0, -1) : trimmed; +} + +function toFiniteNumber(value: unknown): number | undefined { + if (typeof value !== "number" || !Number.isFinite(value)) return undefined; + return value; +} + +function resolveStatus(usedFraction: number | undefined): UsageStatus | undefined { + if (usedFraction === undefined) return undefined; + if (usedFraction >= 1) return "exhausted"; + if (usedFraction >= 0.9) return "warning"; + return "ok"; +} + +function buildAmount(args: { + used: number | undefined; + limit: number | undefined; + remaining: number | undefined; + unit: UsageAmount["unit"]; +}): UsageAmount { + const used = args.used; + const limit = args.limit; + const usedFraction = + used !== undefined && limit !== undefined && limit > 0 ? Math.min(used / limit, 1) : undefined; + const remainingFraction = usedFraction !== undefined ? Math.max(1 - usedFraction, 0) : undefined; + return { + used, + limit, + remaining: args.remaining, + usedFraction, + remainingFraction, + unit: args.unit, + }; +} + +function buildRequestsLimit(payload: UmansUsagePayload, provider: string): UsageLimit | null { + const limit = toFiniteNumber(payload.limits?.requests?.limit); + const windowSeconds = toFiniteNumber(payload.limits?.requests?.window_seconds); + const used = toFiniteNumber(payload.usage?.requests_in_window); + const remaining = toFiniteNumber(payload.usage?.remaining_requests); + if (limit === undefined && used === undefined) return null; + const amount = buildAmount({ used, limit, remaining, unit: "requests" }); + // Rolling window: each request ages out `window_seconds` after it fired, so + // there is no single reset timestamp. Surface the window size + label only. + // `window.id` is `"5h"` to match the status-line usage segment's window-id + // contract (it only recognizes `"5h"`/`"7d"`); `label` stays human-readable. + const window: UsageWindow = { + id: "5h", + label: "rolling 5h", + durationMs: windowSeconds ? windowSeconds * 1000 : FIVE_HOURS_MS, + }; + return { + id: "umans:requests", + label: "Requests (rolling 5h)", + scope: { provider, windowId: window.id, shared: true }, + window, + amount, + status: resolveStatus(amount.usedFraction), + }; +} + +function buildConcurrencyLimit(payload: UmansUsagePayload, provider: string): UsageLimit | null { + const limit = toFiniteNumber(payload.limits?.concurrency?.limit); + const used = toFiniteNumber(payload.usage?.concurrent_sessions); + if (limit === undefined && used === undefined) return null; + const amount = buildAmount({ used, limit, remaining: undefined, unit: "requests" }); + return { + id: "umans:concurrency", + label: "Concurrency", + // Concurrency is instantaneous, not windowed. + scope: { provider, windowId: "concurrency" }, + amount, + status: resolveStatus(amount.usedFraction), + }; +} + +async function fetchUmansUsage(params: UsageFetchParams, ctx: UsageFetchContext): Promise { + if (params.provider !== UMANS_PROVIDER) return null; + const credential = params.credential; + if (credential.type !== "api_key" || !credential.apiKey) return null; + + const baseUrl = normalizeBaseUrl(params.baseUrl); + const url = `${baseUrl}${USAGE_PATH}`; + const headers: Record = { + authorization: `Bearer ${credential.apiKey}`, + accept: "application/json", + }; + + let payload: UmansUsagePayload | null = null; + try { + const response = await ctx.fetch(url, { headers, signal: params.signal }); + if (!response.ok) { + ctx.logger?.warn("Umans usage fetch failed", { status: response.status, statusText: response.statusText }); + return null; + } + const json = (await response.json()) as unknown; + if (!isRecord(json)) { + ctx.logger?.warn("Umans usage response was not a JSON object"); + return null; + } + payload = json as unknown as UmansUsagePayload; + } catch (error) { + ctx.logger?.warn("Umans usage fetch error", { error: String(error) }); + return null; + } + + const limits: UsageLimit[] = []; + const requests = buildRequestsLimit(payload, params.provider); + if (requests) limits.push(requests); + const concurrency = buildConcurrencyLimit(payload, params.provider); + if (concurrency) limits.push(concurrency); + if (limits.length === 0) return null; + + const notes: string[] = []; + if (payload.usage?.priority?.low === true) { + notes.push("Requests deprioritized after a rate-limit burst."); + } + + return { + provider: params.provider, + fetchedAt: Date.now(), + limits, + notes: notes.length > 0 ? notes : undefined, + metadata: { + plan: payload.plan?.display_name, + accountId: credential.accountId, + email: credential.email, + endpoint: url, + }, + raw: payload as Record, + }; +} + +export const umansUsageProvider: UsageProvider = { + id: UMANS_PROVIDER, + fetchUsage: fetchUmansUsage, + supports: params => params.provider === UMANS_PROVIDER && params.credential.type === "api_key", + validatesCredentials: true, +}; From 72077f0ba23bad6170de582bfec699e715af1988 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Sat, 27 Jun 2026 11:54:19 -0700 Subject: [PATCH 003/132] feat(ai): register umans in DEFAULT_USAGE_PROVIDERS Surfaces the umans provider in /usage, omp usage, and the TUI footer by adding umansUsageProvider to the default usage-provider map. --- packages/ai/src/auth-storage.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/ai/src/auth-storage.ts b/packages/ai/src/auth-storage.ts index fc71c9f70..d4e4cd792 100644 --- a/packages/ai/src/auth-storage.ts +++ b/packages/ai/src/auth-storage.ts @@ -56,6 +56,7 @@ import { } from "./usage/openai-codex-reset"; import { opencodeGoUsageProvider } from "./usage/opencode-go"; import { zaiUsageProvider } from "./usage/zai"; +import { umansUsageProvider } from "./usage/umans"; const USAGE_RANKING_METRIC_EPSILON = 1e-9; @@ -515,6 +516,7 @@ const DEFAULT_USAGE_PROVIDERS: UsageProvider[] = [ ollamaCloudUsageProvider, claudeUsageProvider, zaiUsageProvider, + umansUsageProvider, opencodeGoUsageProvider, githubCopilotUsageProvider, ]; From 04e6c338f3d1c01c36c2f5aa30e4c85e9b01d4c9 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Sat, 27 Jun 2026 17:49:02 -0700 Subject: [PATCH 004/132] chore(ai): biome formatting, import order, CHANGELOG for umans usage provider --- packages/ai/CHANGELOG.md | 3 +++ packages/ai/src/auth-storage.ts | 2 +- packages/ai/src/usage/umans.ts | 3 +-- packages/ai/test/umans-usage.test.ts | 18 ++++++++++++++---- 4 files changed, 19 insertions(+), 7 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 9cbda6ba9..75bee85f3 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -18,6 +18,9 @@ ### Added - Added support for Baseten as an AI provider +### Added + +- Umans usage provider: fetches `GET /v1/usage` and surfaces the rolling 5h request window + concurrency limits in `/usage`, `omp usage`, and the TUI status bar. ### Changed diff --git a/packages/ai/src/auth-storage.ts b/packages/ai/src/auth-storage.ts index d4e4cd792..142316437 100644 --- a/packages/ai/src/auth-storage.ts +++ b/packages/ai/src/auth-storage.ts @@ -55,8 +55,8 @@ import { listCodexResetCredits, } from "./usage/openai-codex-reset"; import { opencodeGoUsageProvider } from "./usage/opencode-go"; -import { zaiUsageProvider } from "./usage/zai"; import { umansUsageProvider } from "./usage/umans"; +import { zaiUsageProvider } from "./usage/zai"; const USAGE_RANKING_METRIC_EPSILON = 1e-9; diff --git a/packages/ai/src/usage/umans.ts b/packages/ai/src/usage/umans.ts index bf8007b15..3e986a678 100644 --- a/packages/ai/src/usage/umans.ts +++ b/packages/ai/src/usage/umans.ts @@ -59,8 +59,7 @@ function buildAmount(args: { }): UsageAmount { const used = args.used; const limit = args.limit; - const usedFraction = - used !== undefined && limit !== undefined && limit > 0 ? Math.min(used / limit, 1) : undefined; + const usedFraction = used !== undefined && limit !== undefined && limit > 0 ? Math.min(used / limit, 1) : undefined; const remainingFraction = usedFraction !== undefined ? Math.max(1 - usedFraction, 0) : undefined; return { used, diff --git a/packages/ai/test/umans-usage.test.ts b/packages/ai/test/umans-usage.test.ts index e4010e206..0414af34e 100644 --- a/packages/ai/test/umans-usage.test.ts +++ b/packages/ai/test/umans-usage.test.ts @@ -32,7 +32,11 @@ function fakeFetch(payload: unknown, status = 200): FetchImpl { return fn as unknown as typeof fetch; } -function fetchRecorder(calls: Array<{ url: string; headers: Record }>, payload: unknown, status = 200): FetchImpl { +function fetchRecorder( + calls: Array<{ url: string; headers: Record }>, + payload: unknown, + status = 200, +): FetchImpl { const fn = async (input: string | URL | Request, init?: RequestInit) => { calls.push({ url: String(input), @@ -145,9 +149,15 @@ describe("umans usage provider", () => { }); it("returns null when supports() is called for a different provider or credential type", () => { - expect(umansUsageProvider.supports?.({ provider: "zai", credential: { type: "api_key", apiKey: "x" } })).toBe(false); - expect(umansUsageProvider.supports?.({ provider: "umans", credential: { type: "oauth", accessToken: "x" } })).toBe(false); - expect(umansUsageProvider.supports?.({ provider: "umans", credential: { type: "api_key", apiKey: "x" } })).toBe(true); + expect(umansUsageProvider.supports?.({ provider: "zai", credential: { type: "api_key", apiKey: "x" } })).toBe( + false, + ); + expect( + umansUsageProvider.supports?.({ provider: "umans", credential: { type: "oauth", accessToken: "x" } }), + ).toBe(false); + expect(umansUsageProvider.supports?.({ provider: "umans", credential: { type: "api_key", apiKey: "x" } })).toBe( + true, + ); }); it("includes plan display name and account identity in metadata", async () => { From ce13bda4bf2477168d5fa80c6d2c959e5613ab14 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Fri, 3 Jul 2026 16:42:11 -0700 Subject: [PATCH 005/132] fix(ai): move umans usage provider changelog entry to Unreleased The rebase onto main (16.3.4 release) moved the entry into the released [16.3.4] section. Release sections are immutable; move it under [Unreleased] where it belongs. --- packages/ai/CHANGELOG.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 75bee85f3..60ce4c279 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Added + +- Umans usage provider: fetches `GET /v1/usage` and surfaces the rolling 5h request window + concurrency limits in `/usage`, `omp usage`, and the TUI status bar. + ## [16.3.5] - 2026-07-04 ### Added @@ -18,9 +22,6 @@ ### Added - Added support for Baseten as an AI provider -### Added - -- Umans usage provider: fetches `GET /v1/usage` and surfaces the rolling 5h request window + concurrency limits in `/usage`, `omp usage`, and the TUI status bar. ### Changed From 61ae1801b17a921b15b9a3d273ce630891c4e5e3 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Fri, 3 Jul 2026 16:42:42 -0700 Subject: [PATCH 006/132] fix(ai): normalize umans usage base URL to origin Provider base URLs from modelRegistry.getProviderBaseUrl() can include a /v1 path (e.g. https://api.code.umans.ai/v1). The old trailing-slash strip produced https://api.code.umans.ai/v1/v1/usage. Switch to new URL().origin (same as zai.ts) to strip any path. Add a test case for the /v1 suffix. --- packages/ai/src/usage/umans.ts | 9 +++++++-- packages/ai/test/umans-usage.test.ts | 13 +++++++++++++ 2 files changed, 20 insertions(+), 2 deletions(-) diff --git a/packages/ai/src/usage/umans.ts b/packages/ai/src/usage/umans.ts index 3e986a678..a71a486dd 100644 --- a/packages/ai/src/usage/umans.ts +++ b/packages/ai/src/usage/umans.ts @@ -35,8 +35,13 @@ interface UmansUsagePayload { function normalizeBaseUrl(baseUrl?: string): string { if (!baseUrl?.trim()) return DEFAULT_ENDPOINT; const trimmed = baseUrl.trim(); - // Strip a trailing slash so `${origin}/v1/usage` is stable. - return trimmed.endsWith("/") ? trimmed.slice(0, -1) : trimmed; + // Normalize to origin so provider base URLs that include a path (e.g. + // `https://api.code.umans.ai/v1`) don't double the `/v1` in the usage path. + try { + return new URL(trimmed).origin; + } catch { + return DEFAULT_ENDPOINT; + } } function toFiniteNumber(value: unknown): number | undefined { diff --git a/packages/ai/test/umans-usage.test.ts b/packages/ai/test/umans-usage.test.ts index 0414af34e..0faa0898b 100644 --- a/packages/ai/test/umans-usage.test.ts +++ b/packages/ai/test/umans-usage.test.ts @@ -116,6 +116,19 @@ describe("umans usage provider", () => { expect(calls[0]?.url).toBe("https://custom.umans.example/v1/usage"); }); + it("strips a /v1 path from a custom baseUrl", async () => { + const calls: Array<{ url: string; headers: Record }> = []; + await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + baseUrl: "https://api.code.umans.ai/v1", + }, + { fetch: fetchRecorder(calls, umansPayload()) }, + ); + expect(calls[0]?.url).toBe("https://api.code.umans.ai/v1/usage"); + }); + it("surfaces priority.low as a provider note", async () => { const payload = umansPayload({ usage: { From 5bd598a57308c1c6016936111b5bae53575bbed1 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Fri, 3 Jul 2026 17:05:05 -0700 Subject: [PATCH 007/132] fix(ai): throw on 401/403 in umans usage probe for credential health When /v1/usage returns 401/403 for a bad Umans key, fetchUsage now throws ProviderHttpError instead of returning null. AuthStorage. checkCredentials treats a thrown probe as ok:false (credential failed) rather than ok:null (unknown), so omp auth-gateway check flags invalid keys instead of showing them as unverifiable no-data. Transient non-auth failures (500, etc.) still return null. --- packages/ai/src/usage/umans.ts | 12 ++++++++++++ packages/ai/test/umans-usage.test.ts | 28 ++++++++++++++++++++++++++-- 2 files changed, 38 insertions(+), 2 deletions(-) diff --git a/packages/ai/src/usage/umans.ts b/packages/ai/src/usage/umans.ts index a71a486dd..774bff895 100644 --- a/packages/ai/src/usage/umans.ts +++ b/packages/ai/src/usage/umans.ts @@ -1,3 +1,4 @@ +import { ProviderHttpError } from "../error"; import type { UsageAmount, UsageFetchContext, @@ -133,6 +134,15 @@ async function fetchUmansUsage(params: UsageFetchParams, ctx: UsageFetchContext) try { const response = await ctx.fetch(url, { headers, signal: params.signal }); if (!response.ok) { + // Auth failures (401/403) must throw so checkCredentials flags the bad + // key as ok:false rather than ok:null (unknown). Other non-ok statuses + // are transient — return null so the probe reports "no data". + if (response.status === 401 || response.status === 403) { + throw new ProviderHttpError( + `Umans usage endpoint returned ${response.status} ${response.statusText}`.trim(), + response.status, + ); + } ctx.logger?.warn("Umans usage fetch failed", { status: response.status, statusText: response.statusText }); return null; } @@ -143,6 +153,8 @@ async function fetchUmansUsage(params: UsageFetchParams, ctx: UsageFetchContext) } payload = json as unknown as UmansUsagePayload; } catch (error) { + // Re-throw auth errors so the credential-health probe can surface them. + if (error instanceof ProviderHttpError) throw error; ctx.logger?.warn("Umans usage fetch error", { error: String(error) }); return null; } diff --git a/packages/ai/test/umans-usage.test.ts b/packages/ai/test/umans-usage.test.ts index 0faa0898b..7cddfcf3b 100644 --- a/packages/ai/test/umans-usage.test.ts +++ b/packages/ai/test/umans-usage.test.ts @@ -150,13 +150,37 @@ describe("umans usage provider", () => { expect(report?.notes).toContain("Requests deprioritized after a rate-limit burst."); }); - it("returns null on a non-ok HTTP response", async () => { + it("throws on a 401 auth failure so checkCredentials flags the bad key", async () => { + await expect( + umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fakeFetch({ message: "unauthorized" }, 401) }, + ), + ).rejects.toThrow(/401/); + }); + + it("throws on a 403 auth failure so checkCredentials flags the bad key", async () => { + await expect( + umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + }, + { fetch: fakeFetch({ message: "forbidden" }, 403) }, + ), + ).rejects.toThrow(/403/); + }); + + it("returns null on a transient non-auth HTTP failure (500)", async () => { const report = await umansUsageProvider.fetchUsage( { provider: "umans", credential: { type: "api_key", apiKey: "sk-test" }, }, - { fetch: fakeFetch({ message: "unauthorized" }, 401) }, + { fetch: fakeFetch({ message: "internal server error" }, 500) }, ); expect(report).toBeNull(); }); From e5a764e439c7f930f23e5bbd82c11a26ac436833 Mon Sep 17 00:00:00 2001 From: Henning Post Date: Fri, 3 Jul 2026 23:47:00 -0700 Subject: [PATCH 008/132] fix(ai): preserve path-mounted base URLs in umans usage provider MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit normalizeBaseUrl collapsed to origin via new URL().origin, stripping all path segments — a path-mounted gateway like https://gateway.example/team/umans/v1 lost its /team/umans prefix and the usage probe hit the proxy root instead. Strip only the trailing /v1 (with optional surrounding slashes), preserving any preceding path. Add a test for the path-mounted gateway case. --- packages/ai/src/usage/umans.ts | 12 +++++------- packages/ai/test/umans-usage.test.ts | 15 ++++++++++++++- 2 files changed, 19 insertions(+), 8 deletions(-) diff --git a/packages/ai/src/usage/umans.ts b/packages/ai/src/usage/umans.ts index 774bff895..c32675d61 100644 --- a/packages/ai/src/usage/umans.ts +++ b/packages/ai/src/usage/umans.ts @@ -36,13 +36,11 @@ interface UmansUsagePayload { function normalizeBaseUrl(baseUrl?: string): string { if (!baseUrl?.trim()) return DEFAULT_ENDPOINT; const trimmed = baseUrl.trim(); - // Normalize to origin so provider base URLs that include a path (e.g. - // `https://api.code.umans.ai/v1`) don't double the `/v1` in the usage path. - try { - return new URL(trimmed).origin; - } catch { - return DEFAULT_ENDPOINT; - } + // Strip a trailing `/v1` (with optional surrounding slashes) so the usage + // path doesn't double it, but preserve any preceding path prefix (e.g. a + // path-mounted gateway like `https://gateway.example/team/umans/v1`). + const withoutTrailingSlash = trimmed.replace(/\/+$/, ""); + return withoutTrailingSlash.replace(/\/v1$/i, "") || DEFAULT_ENDPOINT; } function toFiniteNumber(value: unknown): number | undefined { diff --git a/packages/ai/test/umans-usage.test.ts b/packages/ai/test/umans-usage.test.ts index 7cddfcf3b..e75b1738d 100644 --- a/packages/ai/test/umans-usage.test.ts +++ b/packages/ai/test/umans-usage.test.ts @@ -116,7 +116,7 @@ describe("umans usage provider", () => { expect(calls[0]?.url).toBe("https://custom.umans.example/v1/usage"); }); - it("strips a /v1 path from a custom baseUrl", async () => { + it("strips a trailing /v1 from a custom baseUrl", async () => { const calls: Array<{ url: string; headers: Record }> = []; await umansUsageProvider.fetchUsage( { @@ -129,6 +129,19 @@ describe("umans usage provider", () => { expect(calls[0]?.url).toBe("https://api.code.umans.ai/v1/usage"); }); + it("preserves a path-mounted gateway prefix while stripping /v1", async () => { + const calls: Array<{ url: string; headers: Record }> = []; + await umansUsageProvider.fetchUsage( + { + provider: "umans", + credential: { type: "api_key", apiKey: "sk-test" }, + baseUrl: "https://gateway.example/team/umans/v1", + }, + { fetch: fetchRecorder(calls, umansPayload()) }, + ); + expect(calls[0]?.url).toBe("https://gateway.example/team/umans/v1/usage"); + }); + it("surfaces priority.low as a provider note", async () => { const payload = umansPayload({ usage: { From 29625f08c2443d8f6588f96abf4e1b01e519feec Mon Sep 17 00:00:00 2001 From: "Anthony \"Asterisk\" Ambuehl" Date: Tue, 21 Jul 2026 10:39:43 -0700 Subject: [PATCH 009/132] feat(mcp): add mcp_notification extension event + multi-listener API MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Convert MCPManager's dangling single-slot setOnNotification callback into a multi-listener API and expose server-initiated MCP notifications as an extension event so extensions can bridge push-capable MCP servers (e.g. peer messaging, ticket nudges) into session behavior. API changes: - Removed: MCPManager.setOnNotification(handler) — single-slot, zero callers - Added: MCPManager.addNotificationListener(listener): () => void Multi-listener with per-listener error isolation, returns unsub. - Added: 'mcp_notification' extension event Payload: { server: string; method: string; params: unknown } Wired in sdk.ts: one listener bridges to extensionRunner.emitMcpNotification, captured under postmortem for teardown. Tests: 3 new (multi-listener fanout, error isolation, unsubscribe), fixture pattern matches neighboring mcp tests. bun check passes (biome + tsgo). Docs: extensions.md (new MCP notifications subsection with bridging example), mcp-runtime-lifecycle.md (Server-initiated notifications section), CHANGELOG. --- docs/extensions.md | 17 + docs/mcp-runtime-lifecycle.md | 13 + packages/coding-agent/CHANGELOG.md | 8 + .../src/extensibility/extensions/runner.ts | 65 ++++ .../src/extensibility/extensions/types.ts | 27 ++ packages/coding-agent/src/mcp/manager.ts | 138 ++++++- packages/coding-agent/src/sdk.ts | 67 +++- .../test/extensions-runner.test.ts | 163 +++++++++ .../test/fixtures/notifications-mcp.ts | 86 +++++ ...mcp-manager-notification-listeners.test.ts | 336 ++++++++++++++++++ .../test/mcp-server-tool-ownership.test.ts | 8 +- 11 files changed, 896 insertions(+), 32 deletions(-) create mode 100755 packages/coding-agent/test/fixtures/notifications-mcp.ts create mode 100644 packages/coding-agent/test/mcp-manager-notification-listeners.test.ts diff --git a/docs/extensions.md b/docs/extensions.md index 9fd9e3b32..76f19ce27 100644 --- a/docs/extensions.md +++ b/docs/extensions.md @@ -268,6 +268,23 @@ Cancelable pre-events: - `goal_updated` - `credential_disabled` +### MCP notifications + +- `mcp_notification` — fired for every JSON-RPC notification received from a connected MCP server, AFTER the manager's own handling of known list/update methods (`notifications/tools/list_changed`, `notifications/resources/list_changed`, `notifications/resources/updated`, `notifications/prompts/list_changed`). Unknown or server-custom methods are also delivered. Payload: `{ server: string; method: string; params: unknown }`. Multiple extensions may subscribe; a handler that throws does not prevent other handlers from firing. Notifications received before any listener attaches are buffered (bounded FIFO, cap 100, drop-oldest) and drained into the first subscriber — so startup-time frames aren't lost even if the extension binds after MCP discovery. + +Bridging a push-capable MCP into a session steer: + +```ts +pi.on("mcp_notification", event => { + if (event.server !== "peer-bus") return; + if (event.method !== "notifications/peer_message") return; + const params = event.params as { from: string; text: string }; + pi.sendUserMessage(`[from ${params.from}] ${params.text}`, { deliverAs: "steer" }); +}); +``` + +The runtime handles the JSON-RPC transport and its own list/update refresh first; the handler runs afterwards and can inject a mid-turn steer via `pi.sendMessage` / `pi.sendUserMessage`. + ### User command interception - `user_bash` (override with `{ result }`) diff --git a/docs/mcp-runtime-lifecycle.md b/docs/mcp-runtime-lifecycle.md index 7896bf8ac..45707571c 100644 --- a/docs/mcp-runtime-lifecycle.md +++ b/docs/mcp-runtime-lifecycle.md @@ -158,6 +158,19 @@ Both return structured tool output and convert remaining transport/tool errors i There is also a follow-up path for late connections: after waiting for a specific server, if status becomes `connected`, it re-runs `session.refreshMCPTools(...)` so newly available tools are rebound in-session. + +## Server-initiated notifications + +MCP servers may push JSON-RPC notification frames at any point after `initialize` completes. The transport surfaces them via `onNotification`; the manager fans them out in two paths: + +1. **Internal refresh** for known methods: + - `notifications/tools/list_changed` → `refreshServerTools` + - `notifications/resources/list_changed` → `refreshServerResources` + - `notifications/resources/updated` → `#onResourcesChanged` (only for currently subscribed URIs) + - `notifications/prompts/list_changed` → `refreshServerPrompts` +2. **Listener fanout**: every notification (including the known ones AND server-custom methods) is delivered to registered listeners AFTER the internal refresh runs. Registered via `MCPManager.addNotificationListener(listener)`, which returns an unsubscribe function. Multiple listeners are supported; each is invoked with independent error isolation — a synchronous throw in one listener does not prevent others from firing (thrown errors are logged at `debug`). + +`sdk.ts` registers one listener that bridges to the extension runner's `mcp_notification` event, so extensions receive every server-initiated frame with `{ server, method, params }`. The listener is captured with `postmortem` so it is released on session teardown. ## Health, reconnect, and partial failure behavior Current runtime behavior is connection-event driven: diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index e150972d1..ed2bc494c 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,14 @@ ## [Unreleased] +### Added + +- Added `mcp_notification` extension event and multi-listener `MCPManager.addNotificationListener` API. The runtime already received MCP server-initiated JSON-RPC notifications at the transport layer but had no path to forward them to extensions; every notification (including server-custom methods) is now delivered as `{ server, method, params }` after the manager's own list/update handling. For known list-change methods (`notifications/tools/list_changed`, `notifications/resources/list_changed`, `notifications/prompts/list_changed`) the internal refresh promise is awaited before fanout, so a listener acting on `tools/list_changed` sees fresh `getTools()`. Notifications received before any listener attaches are buffered (bounded FIFO, cap 100, drop-oldest — matches `IrcBus`'s `MAILBOX_CAP`) and drained into the first subscriber, so startup-time frames aren't lost even if the extension binds after MCP discovery. Extensions can use this to bridge push-capable MCP servers (e.g. peer messaging) into session behavior by injecting a mid-turn steer via `pi.sendMessage` / `pi.sendUserMessage`. + +### Removed + +- Removed the dangling `MCPManager.setOnNotification` single-slot setter, which had no callers in the runtime. Replaced by `MCPManager.addNotificationListener` — multi-listener, per-listener error isolation, returns an unsubscribe function. + ## [17.1.2] - 2026-07-24 ### Added diff --git a/packages/coding-agent/src/extensibility/extensions/runner.ts b/packages/coding-agent/src/extensibility/extensions/runner.ts index 43a37f1c3..9c69a147c 100644 --- a/packages/coding-agent/src/extensibility/extensions/runner.ts +++ b/packages/coding-agent/src/extensibility/extensions/runner.ts @@ -39,6 +39,7 @@ import type { ExtensionUIContext, InputEvent, InputEventResult, + McpNotificationEvent, MessageRenderer, RegisteredCommand, RegisteredTool, @@ -132,6 +133,14 @@ async function raceHandlerWithTimeout( const MAX_PENDING_CREDENTIAL_DISABLED = 32; +/** + * Buffer cap for `mcp_notification` events received before {@link ExtensionRunner.initialize} + * has run. Sized to match the manager-side buffer in `MCPManager.NOTIFICATION_BUFFER_CAP` so + * the two layers can't drop different amounts of the same burst — the pipe drains, or it + * spills, but it does so consistently at both ends. Drop-oldest under pressure. + */ +const MAX_PENDING_MCP_NOTIFICATIONS = 100; + /** * Events handled by the generic emit() method. * Events with dedicated emitXxx() methods are excluded for stronger type safety. @@ -259,6 +268,19 @@ export class ExtensionRunner { */ #pendingCredentialDisabled: CredentialDisabledEvent[] = []; + /** + * Buffer for `mcp_notification` events received via {@link emitMcpNotification} before + * {@link initialize} has run. Two-layer race: `MCPManager` also buffers frames until + * its first `addNotificationListener` subscriber attaches, but the sdk.ts bridge is + * registered inside `createAgentSession` — BEFORE the mode controller calls + * `ExtensionRunner.initialize()`. Without this second buffer, the manager's drain + * arrives at the bridge → the bridge calls `emitMcpNotification` → the runner drops + * the frame because `#initialized === false`, and the frame evaporates a second time. + * Bounded at {@link MAX_PENDING_MCP_NOTIFICATIONS}; oldest entries are dropped under + * pressure. Drained in {@link initialize} once the runtime/UI context is wired. + */ + #pendingMcpNotifications: Array> = []; + /** * Timers scheduled by extensions through the sanctioned `ctx.setInterval` / * `ctx.setTimeout` helpers. Callbacks run with the same isolation as handler @@ -345,6 +367,23 @@ export class ExtensionRunner { }); } }); + + // Drain events buffered by emitMcpNotification() before initialize ran, using the + // same deferred-microtask ordering as the credential-disabled drain above so any + // onError listener registered synchronously after initialize() still catches + // handler errors during flush. + const pendingMcp = this.#pendingMcpNotifications.splice(0); + queueMicrotask(() => { + for (const event of pendingMcp) { + this.emit({ type: "mcp_notification", ...event }).catch((error: unknown) => { + logger.warn("mcp_notification handler threw during initialize flush", { + server: event.server, + method: event.method, + error: error instanceof Error ? error.message : String(error), + }); + }); + } + }); } /** @@ -372,6 +411,32 @@ export class ExtensionRunner { await this.emit({ type: "credential_disabled", ...event }); } + /** + * Forward an MCP server notification to extension handlers. + * + * If {@link initialize} has not yet run, the notification is buffered and replayed + * once initialize wires the runtime/UI context. Matches the credential-disabled + * deferral above: the sdk.ts bridge registers `MCPManager.addNotificationListener` + * inside `createAgentSession` — BEFORE the mode controller calls `initialize()` on + * this runner — so notification frames drained by the manager (either fresh + * arrivals or replay from its own startup buffer) can reach us pre-init. Without + * this buffer they would evaporate for a second time here. + * + * Bounded at {@link MAX_PENDING_MCP_NOTIFICATIONS}; oldest entries drop under + * pressure. Never throws; per-handler errors are routed through {@link onError} + * via {@link emit}'s normal isolation. + */ + async emitMcpNotification(event: Omit): Promise { + if (!this.#initialized) { + if (this.#pendingMcpNotifications.length >= MAX_PENDING_MCP_NOTIFICATIONS) { + this.#pendingMcpNotifications.shift(); + } + this.#pendingMcpNotifications.push(event); + return; + } + await this.emit({ type: "mcp_notification", ...event }); + } + async emitSessionStop(event: Omit): Promise { return await this.emit({ type: "session_stop", ...event }); } diff --git a/packages/coding-agent/src/extensibility/extensions/types.ts b/packages/coding-agent/src/extensibility/extensions/types.ts index 0d7493b68..3a6ce857a 100644 --- a/packages/coding-agent/src/extensibility/extensions/types.ts +++ b/packages/coding-agent/src/extensibility/extensions/types.ts @@ -713,6 +713,31 @@ export interface CredentialDisabledEvent { disabledCause: string; } +// ============================================================================ +// MCP Events +// ============================================================================ + +/** + * Fired for every JSON-RPC notification received from a connected MCP server, + * AFTER the runtime's own handling of known list/update methods. Unknown or + * server-custom methods are delivered too — extensions can bridge them into + * session behavior by inspecting `method`/`params` and injecting a follow-up + * via `pi.sendMessage(..., { deliverAs })` or `pi.sendUserMessage(...)`. + */ +export interface McpNotificationEvent { + type: "mcp_notification"; + /** + * Server name as declared in the MCP config (raw, unsanitized). Note this + * differs from the sanitized prefix used in `mcp___` + * tool names — filter by this raw name, not by tool-name prefix matching. + */ + server: string; + /** JSON-RPC method (e.g. `notifications/tools/list_changed`, or server-custom). */ + method: string; + /** JSON-RPC params, opaque to the runtime. */ + params: unknown; +} + // ============================================================================ // User Bash Events // ============================================================================ @@ -941,6 +966,7 @@ export type ExtensionEvent = | TodoReminderEvent | GoalUpdatedEvent | CredentialDisabledEvent + | McpNotificationEvent | UserBashEvent | UserPythonEvent | InputEvent @@ -1134,6 +1160,7 @@ export interface ExtensionAPI { on(event: "tool_result", handler: ExtensionHandler): void; on(event: "user_bash", handler: ExtensionHandler): void; on(event: "user_python", handler: ExtensionHandler): void; + on(event: "mcp_notification", handler: ExtensionHandler): void; // ========================================================================= // Tool Registration diff --git a/packages/coding-agent/src/mcp/manager.ts b/packages/coding-agent/src/mcp/manager.ts index 12a8e1ad9..16fd00321 100644 --- a/packages/coding-agent/src/mcp/manager.ts +++ b/packages/coding-agent/src/mcp/manager.ts @@ -94,6 +94,14 @@ const STARTUP_TIMEOUT_MS = 250; const RECONNECT_BURST_WINDOW_MS = 30_000; const RECONNECT_BURST_LIMIT = 5; +/** + * Bounded buffer for notifications received before any listener attaches. + * Mirrors {@link IrcBus}'s `MAILBOX_CAP` — drop-oldest on overflow. Drained + * into the first {@link MCPManager.addNotificationListener} subscriber, then + * cleared; subsequent frames deliver directly to attached listeners. + */ +const NOTIFICATION_BUFFER_CAP = 100; + function trackPromise(promise: Promise): TrackedPromise { const tracked: TrackedPromise = { promise, status: "pending" }; promise.then( @@ -194,8 +202,14 @@ export class MCPManager { #sources = new Map(); #authStorage: AuthStorage | null = null; #authHandler?: MCPAuthHandler; - #onNotification?: (serverName: string, method: string, params: unknown) => void; - #onToolsChanged?: (tools: CustomTool[]) => void; + #notificationListeners = new Set<(serverName: string, method: string, params: unknown) => void>(); + /** + * Notifications received before any listener attached, to be drained on + * the first {@link addNotificationListener} call. Bounded by + * {@link NOTIFICATION_BUFFER_CAP}, drop-oldest on overflow. + */ + #pendingNotifications: Array<{ server: string; method: string; params: unknown }> = []; + #onToolsChanged?: (tools: CustomTool[]) => void | Promise; #onResourcesChanged?: (serverName: string, uri: string) => void; #onPromptsChanged?: (serverName: string) => void; #notificationsEnabled = false; @@ -219,16 +233,66 @@ export class MCPManager { ) {} /** - * Set a callback to receive all server notifications. + * Register a listener for server-initiated MCP notifications. + * + * The listener is called for every JSON-RPC notification received from any + * connected server, AFTER the manager's own handling of known methods + * (`notifications/tools/list_changed`, `notifications/resources/list_changed`, + * `notifications/resources/updated`, `notifications/prompts/list_changed`). + * For list-change methods the internal refresh promise is awaited before + * fanout, so listeners observe up-to-date manager and tool state. Unknown + * or server-custom methods are also delivered, letting consumers bridge + * server-initiated events into session-level behavior (e.g. an extension + * injecting a steer via `pi.sendMessage`). + * + * Notifications received before any listener attached are buffered + * (bounded FIFO, cap {@link NOTIFICATION_BUFFER_CAP}, drop-oldest) and + * drained into the first subscriber — matches {@link setOnPromptsChanged}'s + * replay-on-attach and {@link IrcBus}'s mailbox semantics. + * + * Returns an unsubscribe function; call it to remove the listener. + * + * Multiple listeners are allowed; each is invoked with independent error + * isolation — a listener that throws does not prevent other listeners from + * firing. */ - setOnNotification(handler: (serverName: string, method: string, params: unknown) => void): void { - this.#onNotification = handler; + addNotificationListener(listener: (serverName: string, method: string, params: unknown) => void): () => void { + const wasEmpty = this.#notificationListeners.size === 0; + this.#notificationListeners.add(listener); + + // Drain startup-buffered notifications into the first attaching listener. + if (wasEmpty && this.#pendingNotifications.length > 0) { + const pending = this.#pendingNotifications.splice(0); + for (const frame of pending) { + try { + listener(frame.server, frame.method, frame.params); + } catch (error) { + logger.debug("MCP notification listener threw during buffered drain", { + path: `mcp:${frame.server}`, + method: frame.method, + error, + }); + } + } + } + + return () => { + this.#notificationListeners.delete(listener); + }; } /** * Set a callback to fire when any server's tools change. + * + * May return a Promise; if so, {@link refreshServerTools} awaits it so that + * downstream consumers (e.g. `mcp_notification` listeners for + * `notifications/tools/list_changed`) observe not just the manager's + * refreshed tool set but also any session-level rebind driven by the + * handler (`session.refreshMCPTools`). Other callsites (initial connect, + * disconnect, reconnect) invoke the handler synchronously — their downstream + * chains don't need to serialize on the rebind. */ - setOnToolsChanged(handler: (tools: CustomTool[]) => void): void { + setOnToolsChanged(handler: (tools: CustomTool[]) => void | Promise): void { this.#onToolsChanged = handler; } @@ -488,7 +552,7 @@ export class MCPManager { this.reconnectServer(name, options); const customTools = MCPTool.fromTools(connection, serverTools, reconnect); this.#replaceServerTools(name, customTools); - this.#onToolsChanged?.(this.#tools); + void this.#onToolsChanged?.(this.#tools); void this.toolCache?.set(name, config, serverTools); onStatus?.({ type: "connected", serverName: name }); @@ -597,7 +661,7 @@ export class MCPManager { sortMCPToolsByName(this.#tools); } - #triggerNotificationRefresh(serverName: string, kind: "tools" | "resources" | "prompts"): void { + #triggerNotificationRefresh(serverName: string, kind: "tools" | "resources" | "prompts"): Promise { const refresh = (() => { switch (kind) { case "tools": @@ -608,22 +672,33 @@ export class MCPManager { return this.refreshServerPrompts(serverName); } })(); - void refresh.catch(error => { + return refresh.catch(error => { logger.debug("Failed MCP notification refresh", { path: `mcp:${serverName}`, kind, error }); }); } - #handleServerNotification(serverName: string, method: string, params: unknown): void { + async #handleServerNotification(serverName: string, method: string, params: unknown): Promise { logger.debug("MCP notification received", { path: `mcp:${serverName}`, method }); + // Only trigger refresh if the connection is already stored — during the + // initial connect handshake, notifications may arrive before + // `#connections.set()` completes, and `refreshServer*` would no-op + // anyway. Skipping the await in that case preserves arrival order + // across concurrently-dispatched notifications (an awaited refresh, + // even a no-op, yields a microtask that lets later frames overtake). + const connectionKnown = this.#connections.has(serverName); + let refreshPromise: Promise | undefined; switch (method) { case MCPNotificationMethods.TOOLS_LIST_CHANGED: - this.#triggerNotificationRefresh(serverName, "tools"); + if (connectionKnown) refreshPromise = this.#triggerNotificationRefresh(serverName, "tools"); break; case MCPNotificationMethods.RESOURCES_LIST_CHANGED: - this.#triggerNotificationRefresh(serverName, "resources"); + if (connectionKnown) refreshPromise = this.#triggerNotificationRefresh(serverName, "resources"); break; case MCPNotificationMethods.RESOURCES_UPDATED: { - const uri = (params as { uri?: string })?.uri; + const uri = + params && typeof params === "object" && "uri" in params && typeof params.uri === "string" + ? params.uri + : undefined; const subscribed = this.#subscribedResources.get(serverName); if (uri && subscribed?.has(uri)) { this.#onResourcesChanged?.(serverName, uri); @@ -631,13 +706,40 @@ export class MCPManager { break; } case MCPNotificationMethods.PROMPTS_LIST_CHANGED: - this.#triggerNotificationRefresh(serverName, "prompts"); + if (connectionKnown) refreshPromise = this.#triggerNotificationRefresh(serverName, "prompts"); break; default: break; } - this.#onNotification?.(serverName, method, params); + // Await internal refresh so listeners see the manager's post-refresh + // state (satisfies the documented "AFTER the manager's own handling" + // contract on `addNotificationListener` — otherwise an extension acting + // on `tools/list_changed` could hit stale `getTools()`). + if (refreshPromise) { + await refreshPromise; + } + + // Buffer for late-attaching subscribers when no listener exists yet. + if (this.#notificationListeners.size === 0) { + this.#pendingNotifications.push({ server: serverName, method, params }); + if (this.#pendingNotifications.length > NOTIFICATION_BUFFER_CAP) { + this.#pendingNotifications.shift(); + } + return; + } + + for (const listener of this.#notificationListeners) { + try { + listener(serverName, method, params); + } catch (error) { + logger.debug("MCP notification listener threw", { + path: `mcp:${serverName}`, + method, + error, + }); + } + } } /** Handle server-to-client JSON-RPC requests (e.g. ping, roots/list). */ @@ -781,7 +883,7 @@ export class MCPManager { // Remove tools from this server and notify consumers const hadTools = this.#tools.some(t => t.mcpServerName === name); this.#tools = this.#tools.filter(t => t.mcpServerName !== name); - if (hadTools) this.#onToolsChanged?.(this.#tools); + if (hadTools) void this.#onToolsChanged?.(this.#tools); // Notify prompt consumers so stale commands are cleared if (connection?.prompts?.length) this.#onPromptsChanged?.(name); @@ -1022,7 +1124,7 @@ export class MCPManager { const customTools = MCPTool.fromTools(connection, serverTools, reconnect); void this.toolCache?.set(name, config, serverTools); this.#replaceServerTools(name, customTools); - this.#onToolsChanged?.(this.#tools); + void this.#onToolsChanged?.(this.#tools); void this.#loadServerResourcesAndPrompts(name, connection); return connection; } catch (error) { @@ -1081,7 +1183,7 @@ export class MCPManager { // Replace tools from this server this.#replaceServerTools(name, customTools); - this.#onToolsChanged?.(this.#tools); + await this.#onToolsChanged?.(this.#tools); } /** diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index a08aa5683..4a8f613cc 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -3203,6 +3203,15 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} throw new Error(`Agent "${resolvedAgentId}" was replaced during session initialization.`); } hasRegistered = true; + // MCP notification bridge cleanup — assigned when the bridge is wired below, + // invoked from the dispose wrapper AND registered as a postmortem so both + // explicit-dispose (SDK embedders that reuse the process across sessions) and + // process-exit paths tear the listener down. Nulled after use so the closure + // graph (`extensionRunner`, `session`) can be GC'd instead of retained by the + // process-global postmortem list. + let unsubscribeMcpNotifications: (() => void) | undefined; + let unregisterMcpPostmortem: (() => void) | undefined; + { const originalDispose = session.dispose.bind(session); session.dispose = async () => { @@ -3233,6 +3242,12 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} } finally { unregisterUnlessParked(); unsubscribeCredentialDisabled?.(); + unsubscribeMcpNotifications?.(); + unregisterMcpPostmortem?.(); + // Drop refs so the process-global postmortem list doesn't retain + // the bridge closure past explicit dispose. + unsubscribeMcpNotifications = undefined; + unregisterMcpPostmortem = undefined; } }; } @@ -3405,19 +3420,24 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} } } - // Wire MCP manager callbacks to session for reactive tool updates. - // Skip when reusing a parent's manager — the parent owns the callbacks. + // MCP manager wiring has two ownership models: + // * Single-slot callbacks (tools/prompts/resources changed) — exactly one + // owner per manager. When reusing a parent's manager (subagent path, + // see task/executor.ts), the parent already owns these slots so we + // MUST NOT overwrite them. Guarded by `!options.mcpManager`. + // * Notification listener — multi-listener by design. Every session with + // an MCP manager (fresh OR reused) needs its own bridge to its own + // `extensionRunner` so extensions loaded in that session receive frames. + // Guarded only by `mcpManager` (see the second `if` below). if (mcpManager && !options.mcpManager) { - mcpManager.setOnToolsChanged(tools => { - void (async () => { - try { - await session.refreshMCPTools(tools); - } catch (error) { - logger.warn("MCP tool refresh failed", { - error: error instanceof Error ? error.message : String(error), - }); - } - })(); + mcpManager.setOnToolsChanged(async tools => { + try { + await session.refreshMCPTools(tools); + } catch (error) { + logger.warn("MCP tool refresh failed", { + error: error instanceof Error ? error.message : String(error), + }); + } }); // Wire prompt refresh → rebuild MCP prompt slash commands mcpManager.setOnPromptsChanged(serverName => { @@ -3450,6 +3470,29 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} }); } + if (mcpManager) { + // Bridge server-initiated notifications to this session's extension + // handlers. Multi-listener registration: fresh-manager and reused-manager + // sessions both install their own listener here, so a subagent's + // extensions get frames even though the parent owns the single-slot + // tool/prompt/resource callbacks above. MCPManager fires known + // list/update refreshes internally, then invokes all registered + // listeners with (server, method, params) for every frame (including + // server-custom methods). Two-layer buffering protects the startup + // race: MCPManager buffers frames received before the first + // `addNotificationListener` subscriber (drains here); ExtensionRunner + // buffers frames received before `initialize()` and drains them on + // init. Both drop-oldest under pressure at cap 100. + unsubscribeMcpNotifications = mcpManager.addNotificationListener((server, method, params) => { + void extensionRunner.emitMcpNotification({ server, method, params }); + }); + // postmortem.register returns a cancel function; capture it so explicit + // session.dispose can remove this from the global list (see finally above). + unregisterMcpPostmortem = postmortem.register("mcp-notification-listener-cleanup", () => + unsubscribeMcpNotifications?.(), + ); + } + startDeferredMCPDiscovery?.(session); return { diff --git a/packages/coding-agent/test/extensions-runner.test.ts b/packages/coding-agent/test/extensions-runner.test.ts index aaaf82136..e51f7fead 100644 --- a/packages/coding-agent/test/extensions-runner.test.ts +++ b/packages/coding-agent/test/extensions-runner.test.ts @@ -2114,6 +2114,169 @@ describe("ExtensionRunner", () => { }); }); + describe("mcp_notification", () => { + it("delivers mcp_notification events to subscribed extensions with the typed payload", async () => { + const eventsPath = path.join(tempDir.path(), "mcp-notification-events.jsonl"); + const extCode = ` + import * as fs from "node:fs"; + + export default function(pi) { + pi.on("mcp_notification", async (event) => { + fs.appendFileSync( + ${JSON.stringify(eventsPath)}, + JSON.stringify({ + type: event.type, + server: event.server, + method: event.method, + params: event.params, + }) + "\\n", + ); + }); + } + `; + fs.writeFileSync(path.join(extensionsDir, "mcp-notification.ts"), extCode); + + const result = await loadTestExtensions(); + const runner = new ExtensionRunner( + result.extensions, + result.runtime, + tempDir.path(), + sessionManager, + modelRegistry, + ); + runner.initialize( + { + sendMessage: () => {}, + sendUserMessage: () => {}, + appendEntry: () => {}, + setLabel: () => {}, + getActiveTools: () => [], + getAllTools: () => [], + setActiveTools: async () => {}, + getCommands: () => [], + setModel: async () => false, + getThinkingLevel: () => undefined, + setThinkingLevel: () => {}, + getSessionName: () => sessionManager.getSessionName(), + setSessionName: async () => {}, + }, + { + getModel: () => undefined, + isIdle: () => true, + abort: () => {}, + hasPendingMessages: () => false, + shutdown: () => {}, + getContextUsage: () => undefined, + compact: async () => {}, + getSystemPrompt: () => [], + }, + ); + + await runner.emitMcpNotification({ + server: "peers", + method: "notifications/peer_message", + params: { from: "alice", text: "hi" }, + }); + + const events = fs + .readFileSync(eventsPath, "utf8") + .trim() + .split("\n") + .map(line => JSON.parse(line)); + expect(events).toEqual([ + { + type: "mcp_notification", + server: "peers", + method: "notifications/peer_message", + params: { from: "alice", text: "hi" }, + }, + ]); + }); + + it("buffers pre-initialize events and drains them on initialize (caps at 100, drops oldest)", async () => { + // Guard against regression of the two-layer startup race that Codex flagged + // on PR #6535 (commit ffa058aa8): the sdk.ts bridge wires + // mcpManager.addNotificationListener inside createAgentSession BEFORE the + // mode controller calls ExtensionRunner.initialize(). Frames the manager + // drains from its own buffer arrive here pre-init. Prior behavior silently + // dropped them; the fix buffers and drains on initialize (same shape as + // emitCredentialDisabled). + const eventsPath = path.join(tempDir.path(), "mcp-notification-cap.jsonl"); + const extCode = ` + import * as fs from "node:fs"; + + export default function(pi) { + pi.on("mcp_notification", async (event) => { + fs.appendFileSync( + ${JSON.stringify(eventsPath)}, + JSON.stringify({ server: event.server, method: event.method }) + "\\n", + ); + }); + } + `; + fs.writeFileSync(path.join(extensionsDir, "mcp-notification-cap.ts"), extCode); + + const result = await loadTestExtensions(); + const runner = new ExtensionRunner( + result.extensions, + result.runtime, + tempDir.path(), + sessionManager, + modelRegistry, + ); + + // Push 101 events while uninitialized — the 1st should be dropped, next 100 buffered. + for (let i = 0; i < 101; i++) { + await runner.emitMcpNotification({ + server: "peers", + method: `notifications/test/${i}`, + params: null, + }); + } + + runner.initialize( + { + sendMessage: () => {}, + sendUserMessage: () => {}, + appendEntry: () => {}, + setLabel: () => {}, + getActiveTools: () => [], + getAllTools: () => [], + setActiveTools: async () => {}, + getCommands: () => [], + setModel: async () => false, + getThinkingLevel: () => undefined, + setThinkingLevel: () => {}, + getSessionName: () => sessionManager.getSessionName(), + setSessionName: async () => {}, + }, + { + getModel: () => undefined, + isIdle: () => true, + abort: () => {}, + hasPendingMessages: () => false, + shutdown: () => {}, + getContextUsage: () => undefined, + compact: async () => {}, + getSystemPrompt: () => [], + }, + ); + + // Drain microtasks so the fire-and-forget emit() calls inside initialize() complete. + for (let i = 0; i < 5; i++) await Promise.resolve(); + + const events = fs + .readFileSync(eventsPath, "utf8") + .trim() + .split("\n") + .map(line => JSON.parse(line)); + expect(events).toHaveLength(100); + // Drop-oldest policy: test/0 was evicted, test/1 survives as the head. + expect(events[0]?.method).toBe("notifications/test/1"); + expect(events[99]?.method).toBe("notifications/test/100"); + }); + }); + describe("managed timers (ctx.setInterval / ctx.setTimeout)", () => { it("contains a throwing interval callback instead of letting it escape as uncaughtException", () => { vi.useFakeTimers(); diff --git a/packages/coding-agent/test/fixtures/notifications-mcp.ts b/packages/coding-agent/test/fixtures/notifications-mcp.ts new file mode 100755 index 000000000..c8aa1c571 --- /dev/null +++ b/packages/coding-agent/test/fixtures/notifications-mcp.ts @@ -0,0 +1,86 @@ +#!/usr/bin/env bun +/** + * Test fixture: minimal stdio MCP server used by + * `mcp-manager-notification-listeners.test.ts` to exercise the notification + * listener API. On receiving the client's `notifications/initialized`, emits + * TWO server-to-client notification frames: + * + * 1. `notifications/tools/list_changed` — a known method that MCPManager + * handles internally (triggers a `tools/list` refresh). Delivered to + * listeners AFTER the internal handling; verifies fanout is not gated + * on the method being unknown. + * 2. A server-custom method (`notifications/custom/test-event`) with a + * distinctive payload — verifies that arbitrary methods are delivered + * to listeners verbatim, which is the extension use case (bridging a + * peer-messaging MCP's custom push into a session steer). + * + * The custom method + payload constants are exported so the test can assert + * on frame equality without duplicating literals. + */ +import * as readline from "node:readline"; + +export const CUSTOM_NOTIFICATION_METHOD = "notifications/custom/test-event"; +export const CUSTOM_NOTIFICATION_PAYLOAD: Readonly> = Object.freeze({ + hello: "world", + n: 42, +}); + +type JsonRpcMessage = { + jsonrpc: "2.0"; + id?: string | number; + method?: string; + params?: Record; +}; + +function buildResult(method: string): Record { + switch (method) { + case "initialize": + return { + protocolVersion: "2025-03-26", + serverInfo: { name: "notifications-fixture", version: "1.0.0" }, + capabilities: { tools: { listChanged: true } }, + }; + case "tools/list": + return { tools: [] }; + default: + return {}; + } +} + +function writeNotification(method: string, params?: Record): void { + const frame: JsonRpcMessage = { jsonrpc: "2.0", method }; + if (params) frame.params = params; + process.stdout.write(`${JSON.stringify(frame)}\n`); +} + +function startServer(): void { + const rl = readline.createInterface({ input: process.stdin }); + rl.on("line", line => { + void (async () => { + const trimmed = line.trim(); + if (trimmed.length === 0) return; + let msg: JsonRpcMessage; + try { + msg = JSON.parse(trimmed) as JsonRpcMessage; + } catch { + return; + } + const isNotification = msg.id === undefined || msg.id === null; + if (isNotification) { + if (msg.method === "notifications/initialized") { + writeNotification("notifications/tools/list_changed"); + writeNotification(CUSTOM_NOTIFICATION_METHOD, { ...CUSTOM_NOTIFICATION_PAYLOAD }); + } + return; + } + if (!msg.method) return; + const response = { jsonrpc: "2.0" as const, id: msg.id, result: buildResult(msg.method) }; + process.stdout.write(`${JSON.stringify(response)}\n`); + })(); + }); + rl.on("close", () => process.exit(0)); +} + +if (import.meta.main) { + startServer(); +} diff --git a/packages/coding-agent/test/mcp-manager-notification-listeners.test.ts b/packages/coding-agent/test/mcp-manager-notification-listeners.test.ts new file mode 100644 index 000000000..8f2678c14 --- /dev/null +++ b/packages/coding-agent/test/mcp-manager-notification-listeners.test.ts @@ -0,0 +1,336 @@ +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { MCPManager } from "@oh-my-pi/pi-coding-agent/mcp/manager"; +import type { MCPServerConfig } from "@oh-my-pi/pi-coding-agent/mcp/types"; +import { removeSyncWithRetries } from "@oh-my-pi/pi-utils"; +import { CUSTOM_NOTIFICATION_METHOD, CUSTOM_NOTIFICATION_PAYLOAD } from "./fixtures/notifications-mcp"; + +const FIXTURE_PATH = path.join(import.meta.dir, "fixtures", "notifications-mcp.ts"); +const BUN_EXEC = process.execPath; + +type CapturedFrame = { server: string; method: string; params: unknown }; + +function serverConfig(): MCPServerConfig { + return { type: "stdio", command: BUN_EXEC, args: [FIXTURE_PATH] }; +} + +/** + * Bind a capturing listener that ALSO exposes a promise-based await for a + * specific `(server, method)` pair. Lets each test wait for the exact frame it + * cares about without polling — the frame arriving is the real signal, not a + * timer. bun:test's per-test timeout backstops a genuine hang. + */ +function makeFrameCollector(): { + listener: (server: string, method: string, params: unknown) => void; + frames: CapturedFrame[]; + awaitFrame: (server: string, method: string) => Promise; +} { + const frames: CapturedFrame[] = []; + const waiters = new Map void>>(); + const key = (server: string, method: string) => `${server}\x00${method}`; + return { + frames, + listener: (server, method, params) => { + const frame: CapturedFrame = { server, method, params }; + frames.push(frame); + const bucket = waiters.get(key(server, method)); + if (bucket && bucket.length > 0) { + waiters.delete(key(server, method)); + for (const resolve of bucket) resolve(frame); + } + }, + awaitFrame: (server, method) => { + const existing = frames.find(f => f.server === server && f.method === method); + if (existing) return Promise.resolve(existing); + const { promise, resolve } = Promise.withResolvers(); + const bucket = waiters.get(key(server, method)) ?? []; + bucket.push(resolve); + waiters.set(key(server, method), bucket); + return promise; + }, + }; +} + +describe("MCPManager notification listeners", () => { + let workDir: string; + + beforeEach(() => { + workDir = fs.mkdtempSync(path.join(os.tmpdir(), "omp-mcp-notif-")); + }); + + afterEach(() => { + removeSyncWithRetries(workDir); + }); + + it("delivers known and server-custom notifications to a registered listener", async () => { + const manager = new MCPManager(workDir); + const collector = makeFrameCollector(); + const unsubscribe = manager.addNotificationListener(collector.listener); + expect(typeof unsubscribe).toBe("function"); + + try { + const result = await manager.connectServers({ alpha: serverConfig() }, {}); + expect(result.connectedServers).toContain("alpha"); + + // Await both the known list_changed and the server-custom frame + // independently. Arrival order across the two isn't guaranteed + // because list_changed triggers an internal `tools/list` refresh + // that MCPManager awaits before fanout (so listeners see fresh + // `getTools()`), which lets a non-refresh frame like the custom + // one overtake in the same batch. Both MUST still be delivered; + // that's the actual contract this test exercises. + const [listChanged, custom] = await Promise.all([ + collector.awaitFrame("alpha", "notifications/tools/list_changed"), + collector.awaitFrame("alpha", CUSTOM_NOTIFICATION_METHOD), + ]); + expect(custom.params).toEqual(CUSTOM_NOTIFICATION_PAYLOAD); + expect(listChanged.method).toBe("notifications/tools/list_changed"); + } finally { + await manager.disconnectAll(); + } + }); + + it("dispatches to every listener and isolates handler errors", async () => { + const manager = new MCPManager(workDir); + const collectorA = makeFrameCollector(); + const collectorC = makeFrameCollector(); + + manager.addNotificationListener(collectorA.listener); + manager.addNotificationListener(() => { + // Middle listener throws synchronously. Must NOT prevent the third + // listener from being invoked for the same frame. + throw new Error("listener-B intentionally throws"); + }); + manager.addNotificationListener(collectorC.listener); + + try { + await manager.connectServers({ alpha: serverConfig() }, {}); + const [frameA, frameC] = await Promise.all([ + collectorA.awaitFrame("alpha", CUSTOM_NOTIFICATION_METHOD), + collectorC.awaitFrame("alpha", CUSTOM_NOTIFICATION_METHOD), + ]); + expect(frameA.params).toEqual(CUSTOM_NOTIFICATION_PAYLOAD); + expect(frameC.params).toEqual(CUSTOM_NOTIFICATION_PAYLOAD); + } finally { + await manager.disconnectAll(); + } + }); + + it("unsubscribe stops delivery to that listener without affecting others", async () => { + const manager = new MCPManager(workDir); + const early = makeFrameCollector(); + const late = makeFrameCollector(); + + const unsubscribeEarly = manager.addNotificationListener(early.listener); + + try { + await manager.connectServers({ alpha: serverConfig() }, {}); + await early.awaitFrame("alpha", CUSTOM_NOTIFICATION_METHOD); + const earlyCountBeforeUnsub = early.frames.length; + + unsubscribeEarly(); + + // Add a fresh listener AFTER unsub, then connect a second server so + // the fixture emits a new batch. The unsubscribed listener must not + // grow; the fresh listener must observe the fresh batch. Dispatch + // per frame is synchronous over the listener set, so once `late` + // resolves for a `beta` frame, any concurrent delivery to `early` + // would have already happened had the unsubscribe not landed. + manager.addNotificationListener(late.listener); + await manager.connectServers({ beta: serverConfig() }, {}); + await late.awaitFrame("beta", CUSTOM_NOTIFICATION_METHOD); + + expect(early.frames.length).toBe(earlyCountBeforeUnsub); + } finally { + await manager.disconnectAll(); + } + }); + + it("buffers notifications received before any listener attaches, drains into first subscriber", async () => { + const manager = new MCPManager(workDir); + + try { + // Connect FIRST, before any listener is registered. The fixture + // emits `notifications/tools/list_changed` plus a server-custom + // notification during its `notifications/initialized` handshake, + // so by the time connect resolves at least one notification frame + // has already been received by the manager with no consumer. + await manager.connectServers({ alpha: serverConfig() }, {}); + + // Attach the FIRST listener AFTER frames were already dispatched. + // Drain-on-first-subscribe must deliver the buffered frames now. + const collector = makeFrameCollector(); + manager.addNotificationListener(collector.listener); + // Both the known list_changed and the server-custom frame must + // have been buffered and drained on attach. Await each + // independently — arrival order between refresh-triggering and + // non-refresh frames is not guaranteed (see the delivery test + // for rationale), but both are always delivered. + const [listChanged, custom] = await Promise.all([ + collector.awaitFrame("alpha", "notifications/tools/list_changed"), + collector.awaitFrame("alpha", CUSTOM_NOTIFICATION_METHOD), + ]); + expect(custom.params).toEqual(CUSTOM_NOTIFICATION_PAYLOAD); + expect(listChanged.method).toBe("notifications/tools/list_changed"); + } finally { + await manager.disconnectAll(); + } + }); + + it("refreshServerTools awaits an async setOnToolsChanged handler before resolving", async () => { + // Verifies the second-layer guarantee (issue raised in review): + // #onToolsChanged in the sdk.ts wiring calls session.refreshMCPTools() + // asynchronously. If refreshServerTools does not await the callback, + // a mcp_notification listener acting on tools/list_changed can call + // getAllTools() on the session before its registry has been rebound. + // This test drives the same pattern by registering a slow async + // setOnToolsChanged handler and confirming refreshServerTools blocks + // on it before resolving — proving the fanout in + // #handleServerNotification (which awaits refreshServerTools) sees + // the callback complete. + const manager = new MCPManager(workDir); + const events: string[] = []; + const { promise: hold, resolve: release } = Promise.withResolvers(); + + manager.setOnToolsChanged(async () => { + events.push("cb:start"); + await hold; + events.push("cb:end"); + }); + + try { + await manager.connectServers({ alpha: serverConfig() }, {}); + // connectServers may kick the callback for the initial tool load via + // its background continuation; drain that first so we can assert + // specifically on the refresh path. + await Bun.sleep(30); + release(); + await Bun.sleep(30); + const priorLen = events.length; + + // Trigger a fresh refresh and hold the callback again — this is the + // path #handleServerNotification takes for tools/list_changed. + const { promise: hold2, resolve: release2 } = Promise.withResolvers(); + let cbState: "started" | "ended" | undefined; + manager.setOnToolsChanged(async () => { + cbState = "started"; + await hold2; + cbState = "ended"; + }); + const refreshDone = manager.refreshServerTools("alpha"); + // Give the callback time to enter the await. + await Bun.sleep(20); + expect(cbState).toBe("started"); + // refreshServerTools has NOT resolved yet — proves it's awaiting. + let refreshResolved = false; + void refreshDone.then(() => { + refreshResolved = true; + }); + await Bun.sleep(10); + expect(refreshResolved).toBe(false); + // Release the callback; refresh should now resolve promptly. + release2(); + await refreshDone; + expect(cbState).toBe("ended"); + // Drain events reference — just to silence unused var if refactored. + expect(events.length).toBeGreaterThanOrEqual(priorLen); + } finally { + await manager.disconnectAll(); + } + }); + + it("mcp_notification listener for tools/list_changed fires only AFTER setOnToolsChanged callback resolves", async () => { + // Option B integration test: exercises the full manager-side chain + // #handleServerNotification → refreshServerTools → #onToolsChanged + // → fanout to notification listeners. Uses a slow async + // setOnToolsChanged callback that gates on a controllable promise + // (matches the mcp-dispose-disconnect-bounded.test.ts gate pattern) + // to prove that if the SDK's setOnToolsChanged handler is async and + // its work isn't complete, the mcp_notification listener does NOT + // fire prematurely. Regression guard: catches any future refactor + // that reintroduces fire-and-forget in the callback chain (e.g. the + // original sdk.ts:3411 `void (async () => ...)()` wrapper — which + // was the specific race raised in review of #6535). + const manager = new MCPManager(workDir); + const events: string[] = []; + let cbCall = 0; + const { promise: gate, resolve: releaseGate } = Promise.withResolvers(); + + manager.setOnToolsChanged(async () => { + cbCall++; + // Initial connect fires the callback via a background continuation + // (manager.ts line ~555, fire-and-forget). Return immediately so + // its floating promise resolves cleanly; only gate on the second + // call, which is the one our simulated post-connect notification + // triggers. + if (cbCall === 1) { + events.push("initial-cb:done"); + return; + } + events.push("cb:start"); + await gate; + events.push("cb:end"); + }); + + manager.addNotificationListener((_server, method) => { + if (method === "notifications/tools/list_changed") { + events.push("listener:fire"); + } + }); + + try { + await manager.connectServers({ alpha: serverConfig() }, {}); + // Wait for the initial-connect background continuation to fire the + // callback (call #1, ungated). Once it's recorded, we know the + // callback is idle and the next invocation will be call #2. + for (let i = 0; i < 20 && !events.includes("initial-cb:done"); i++) { + await Bun.sleep(10); + } + expect(events).toContain("initial-cb:done"); + + // Simulate a post-connect `notifications/tools/list_changed` frame + // by driving the transport's onNotification hook directly — that's + // exactly what the transport does when a server pushes a + // notification after startup. Bypasses the fixture's initial + // handshake burst so we exercise the connectionKnown=true path + // (where refresh actually runs and the callback is awaited). + const conn = manager.getConnection("alpha"); + if (!conn?.transport.onNotification) throw new Error("expected transport onNotification to be wired"); + void conn.transport.onNotification("notifications/tools/list_changed", {}); + + // Give the async chain time to enter the callback's await. + for (let i = 0; i < 20 && !events.includes("cb:start"); i++) { + await Bun.sleep(10); + } + expect(events).toContain("cb:start"); + + // CRITICAL: at this point the callback is holding, and the + // listener MUST NOT have fired. If it has, the manager fanned + // out before awaiting the callback — the exact race this test + // guards against. + expect(events).not.toContain("cb:end"); + expect(events).not.toContain("listener:fire"); + + // Release the gate. Now the callback resolves, refresh completes, + // and the listener fires — in that order. + releaseGate(); + + // Wait for the listener to fire (or bail on timeout). + for (let i = 0; i < 50 && !events.includes("listener:fire"); i++) { + await Bun.sleep(10); + } + + const idxCbEnd = events.indexOf("cb:end"); + const idxFire = events.indexOf("listener:fire"); + expect(idxCbEnd).toBeGreaterThanOrEqual(0); + expect(idxFire).toBeGreaterThanOrEqual(0); + expect(idxFire).toBeGreaterThan(idxCbEnd); + } finally { + // Ensure gate is released even on early failure so no promise leaks. + releaseGate(); + await manager.disconnectAll(); + } + }, 15_000); +}); diff --git a/packages/coding-agent/test/mcp-server-tool-ownership.test.ts b/packages/coding-agent/test/mcp-server-tool-ownership.test.ts index fdd7973c5..bc63ddfdd 100644 --- a/packages/coding-agent/test/mcp-server-tool-ownership.test.ts +++ b/packages/coding-agent/test/mcp-server-tool-ownership.test.ts @@ -61,7 +61,9 @@ describe("MCP tool ownership with prefix-colliding server names", () => { expect(names()).toHaveLength(MANY_TOOL_COUNT * 2); const payloads: string[][] = []; - manager.setOnToolsChanged(tools => payloads.push(tools.map(t => t.name))); + manager.setOnToolsChanged(tools => { + payloads.push(tools.map(t => t.name)); + }); // Same code path a reconnect takes: replace the named server's tools. await manager.refreshServerTools(SHORT_SERVER); @@ -80,7 +82,9 @@ describe("MCP tool ownership with prefix-colliding server names", () => { it("disconnecting a server with sanitized name characters removes exactly its tools", async () => { await manager.connectServers({ [SHORT_SERVER]: fixtureConfig(), [COLON_SERVER]: fixtureConfig() }, {}); const payloads: string[][] = []; - manager.setOnToolsChanged(tools => payloads.push(tools.map(t => t.name))); + manager.setOnToolsChanged(tools => { + payloads.push(tools.map(t => t.name)); + }); await manager.disconnectServer(COLON_SERVER); From 19d7d14a94279436ad7ccf18606b97ab18dec6a4 Mon Sep 17 00:00:00 2001 From: Will Date: Sat, 25 Jul 2026 19:08:02 -0400 Subject: [PATCH 010/132] feat(ai): add Exa API key login --- README.md | 2 +- docs/environment-variables.md | 2 +- docs/tools/web_search.md | 4 +-- packages/ai/CHANGELOG.md | 1 + packages/ai/src/registry/exa.ts | 19 +++++++++++ packages/ai/src/registry/registry.ts | 2 ++ packages/ai/src/stream.ts | 1 - packages/ai/test/exa-login.test.ts | 32 +++++++++++++++++++ packages/ai/test/provider-registry.test.ts | 3 +- packages/coding-agent/CHANGELOG.md | 1 + packages/coding-agent/src/web/search/types.ts | 6 +++- 11 files changed, 66 insertions(+), 7 deletions(-) create mode 100644 packages/ai/src/registry/exa.ts create mode 100644 packages/ai/test/exa-login.test.ts diff --git a/README.md b/README.md index 3e9fe972b..37f6020f6 100644 --- a/README.md +++ b/README.md @@ -373,7 +373,7 @@ Twenty-five backends. Pin one, or let `auto` walk the chain in order. | `codex` | oauth | | `xai` | `XAI_API_KEY` | | `zai` | `ZAI_API_KEY` | -| `exa` | `EXA_API_KEY` (or mcp) | +| `exa` | `/login` / env / mcp | | `tinyfish` | `TINYFISH_API_KEY` | | `jina` | `JINA_API_KEY` | | `kagi` | `KAGI_API_KEY` | diff --git a/docs/environment-variables.md b/docs/environment-variables.md index f5f6536da..fb0deca35 100644 --- a/docs/environment-variables.md +++ b/docs/environment-variables.md @@ -242,7 +242,7 @@ OAuth host chain: `KIMI_CODE_OAUTH_HOST` → `KIMI_OAUTH_HOST` → `https://auth | Variable | Used by | | --------------------------------------------------- | ------------------------------------------------------------- | -| `EXA_API_KEY` | Exa search provider and Exa MCP tools | +| `EXA_API_KEY` | Exa search/MCP; alternatively use `/login exa` | | `BRAVE_API_KEY` | Brave search provider | | `PERPLEXITY_API_KEY` | Perplexity search provider API-key mode | | `PERPLEXITY_COOKIES` | Perplexity cookie-auth search mode | diff --git a/docs/tools/web_search.md b/docs/tools/web_search.md index f54eaa1bb..f14e74e6e 100644 --- a/docs/tools/web_search.md +++ b/docs/tools/web_search.md @@ -147,7 +147,7 @@ Streaming: none. `WebSearchTool.execute()` forwards its `AbortSignal` into `exec - `limit` and `num_search_results` are collapsed together before dispatch. - Output may include parsed free-text `answer`, `sources`, `requestId`. - **Exa** — `packages/coding-agent/src/web/search/providers/exa.ts` - - Availability: env or `agent.db` credential for `exa` admits Exa to the auto chain; settings must not explicitly disable `exa.enabled` or `exa.enableSearch`. Explicit selection (listing `exa` in `providers.webSearchOrder`, or a forced `provider: exa`) reaches Exa even without a credential and falls back to public MCP. + - Availability: `EXA_API_KEY` or a stored credential for `exa` (including one added through `/login exa`) admits Exa to the auto chain; settings must not explicitly disable `exa.enabled` or `exa.enableSearch`. Explicit selection (listing `exa` in `providers.webSearchOrder`, or a forced `provider: exa`) reaches Exa even without a credential and falls back to public MCP. - Querying: POST `https://api.exa.ai/search` with the resolved Exa API key, otherwise JSON-RPC `tools/call` against `https://mcp.exa.ai/mcp` for remote MCP tool `web_search_exa`. - `limit` and `num_search_results` are collapsed together before dispatch. - Output: synthesized `answer` from up to 3 result summaries, `sources`, `requestId`. @@ -279,4 +279,4 @@ Streaming: none. `WebSearchTool.execute()` forwards its `AbortSignal` into `exec - `recency` is implemented by Brave, Perplexity, Tavily, SearXNG, Kagi, TinyFish, Firecrawl, xAI, DuckDuckGo, Bing, Yahoo, Startpage, Google, and Mojeek (Ecosia ignores it; Public Web passes it through). The model-facing prompt does not name specific providers. - `packages/coding-agent/src/config/settings-schema.ts` uses the shared `SEARCH_PROVIDER_PREFERENCES` / `SEARCH_PROVIDER_OPTIONS` metadata, so the settings selector and setup wizard expose `auto` plus every provider in the auto chain. - The credential-free scrapers close the auto chain, cheap plain-fetch engines first (`duckduckgo`, `bing`, `yahoo`, `startpage`) and browser-backed ones after (`google`, `ecosia`, `mojeek`); `public` is listed last and never auto-selected. -- Exa uses `authStorage.getApiKey("exa")`, then `EXA_API_KEY`, then unauthenticated `https://mcp.exa.ai/mcp` fallback. +- `/login exa` stores the pasted key in AuthStorage; Exa resolves credentials in order from `authStorage.getApiKey("exa")`, then `EXA_API_KEY`, then the unauthenticated `https://mcp.exa.ai/mcp` fallback. diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 98cad82ab..880b20846 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -4,6 +4,7 @@ ### Added +- Added interactive Exa API-key login through `/login exa`, opening the official API-key dashboard and saving pasted keys to the credential store ([#1798](https://github.com/can1357/oh-my-pi/issues/1798)). - OAuth logins now stamp `authorizedAt` (epoch ms of the interactive login) on the stored credential, and every refresh-persist path preserves it. Anthropic expires the whole OAuth grant family ~30 days after authorization regardless of refresh-token rotation (observed as `invalid_grant: "Refresh token expired"` on the latest rotated token, exactly 30 days after login, across four production accounts), so the login anchor is what makes re-login deadlines computable. Exported `ANTHROPIC_OAUTH_GRANT_TTL_MS` alongside the anthropic OAuth flow. - Added `GET /v1/credentials/disabled` to the auth broker and `AuthBrokerClient.listDisabledCredentials`: disabled-credential tombstones (`DisabledCredentialSummary` — identity, verbatim disable cause, disable timestamp; never token material) so auto-disabled accounts stay visible to clients instead of silently vanishing from the snapshot. `AuthStorage.listDisabledCredentials` serves the same data locally from SQLite; clients of brokers predating the endpoint get an empty list (404 mapped, no error). - Added `AuthStorage.revalidateCredentials()` and the optional `AuthCredentialStore.refreshSnapshot` hook: remote broker stores re-fetch `GET /v1/snapshot` on demand so callers pairing live per-credential data with stored identities (`omp usage`) never render against the up-to-an-hour-stale disk-cached snapshot; local SQLite stores are always current and only reload. diff --git a/packages/ai/src/registry/exa.ts b/packages/ai/src/registry/exa.ts new file mode 100644 index 000000000..aaa5fb70c --- /dev/null +++ b/packages/ai/src/registry/exa.ts @@ -0,0 +1,19 @@ +import { createApiKeyLogin } from "./api-key-login"; +import type { OAuthLoginCallbacks } from "./oauth/types"; +import type { ProviderDefinition } from "./types"; + +export const loginExa = createApiKeyLogin({ + providerLabel: "Exa", + authUrl: "https://dashboard.exa.ai/api-keys", + instructions: "Create or copy your API key from the Exa dashboard.", + promptMessage: "Paste your Exa API key", + placeholder: "API key", + validation: null, +}); + +export const exaProvider = { + id: "exa", + name: "Exa", + envKeys: "EXA_API_KEY", + login: (cb: OAuthLoginCallbacks) => loginExa(cb), +} as const satisfies ProviderDefinition; diff --git a/packages/ai/src/registry/registry.ts b/packages/ai/src/registry/registry.ts index c6d02aa7d..7c7546f72 100644 --- a/packages/ai/src/registry/registry.ts +++ b/packages/ai/src/registry/registry.ts @@ -12,6 +12,7 @@ import { coreWeaveProvider } from "./coreweave"; import { cursorProvider } from "./cursor"; import { deepseekProvider } from "./deepseek"; import { devinProvider } from "./devin"; +import { exaProvider } from "./exa"; import { firepassProvider } from "./firepass"; import { fireworksProvider } from "./fireworks"; import { githubCopilotProvider } from "./github-copilot"; @@ -134,6 +135,7 @@ const ALL = [ opencodeGoProvider, tavilyProvider, kagiProvider, + exaProvider, parallelProvider, ollamaProvider, ollamaCloudProvider, diff --git a/packages/ai/src/stream.ts b/packages/ai/src/stream.ts index 8a0099f35..cfb9b2e3b 100644 --- a/packages/ai/src/stream.ts +++ b/packages/ai/src/stream.ts @@ -691,7 +691,6 @@ type KeyResolver = string | (() => string | undefined); const LEGACY_ENV_KEYS: Record = { // Non-provider / search-tool keys and API-name keys not modeled as registry provider defs. "azure-openai-responses": "AZURE_OPENAI_API_KEY", - exa: "EXA_API_KEY", jina: "JINA_API_KEY", brave: "BRAVE_API_KEY", tinyfish: "TINYFISH_API_KEY", diff --git a/packages/ai/test/exa-login.test.ts b/packages/ai/test/exa-login.test.ts new file mode 100644 index 000000000..39770d744 --- /dev/null +++ b/packages/ai/test/exa-login.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it } from "bun:test"; +import { loginExa } from "@oh-my-pi/pi-ai/registry/exa"; + +describe("exa login", () => { + it("opens Exa API-key settings and returns a trimmed key without validation requests", async () => { + let authUrl: string | undefined; + let authInstructions: string | undefined; + let promptMessage: string | undefined; + let promptPlaceholder: string | undefined; + + const apiKey = await loginExa({ + onAuth: info => { + authUrl = info.url; + authInstructions = info.instructions; + }, + onPrompt: async prompt => { + promptMessage = prompt.message; + promptPlaceholder = prompt.placeholder; + return " exa-test-key "; + }, + fetch: () => { + throw new Error("Exa login must not make a network request"); + }, + }); + + expect(authUrl).toBe("https://dashboard.exa.ai/api-keys"); + expect(authInstructions).toBe("Create or copy your API key from the Exa dashboard."); + expect(promptMessage).toBe("Paste your Exa API key"); + expect(promptPlaceholder).toBe("API key"); + expect(apiKey).toBe("exa-test-key"); + }); +}); diff --git a/packages/ai/test/provider-registry.test.ts b/packages/ai/test/provider-registry.test.ts index a23b3010f..79a8efb58 100644 --- a/packages/ai/test/provider-registry.test.ts +++ b/packages/ai/test/provider-registry.test.ts @@ -47,7 +47,7 @@ describe("provider registry auth surface", () => { expect(getEnvApiKey("umans")).toBe("umans-env"); Bun.env.LLAMA_CPP_API_KEY = "llama-env"; expect(getEnvApiKey("llama.cpp")).toBe("llama-env"); - // Legacy search-tool key preserved (not a registry provider def). + // Exa is derived from the provider registry's `envKeys` definition. expect(getEnvApiKey("exa")).toBe("exa-env"); }); @@ -65,6 +65,7 @@ describe("provider registry auth surface", () => { const ids = getOAuthProviders().map(provider => provider.id); expect(ids).toContain("zenmux"); expect(ids).toContain("kagi"); + expect(ids).toContain("exa"); expect(ids).toContain("umans"); expect(ids).toContain("llama.cpp"); // openai has no interactive login flow. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index c4777f2be..5307fdb86 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -4,6 +4,7 @@ ### Added +- Added interactive Exa API-key onboarding through `/login exa`, opening the official key dashboard and saving pasted keys for authenticated web search while preserving `EXA_API_KEY` and explicit-selection public MCP fallback behavior ([#1798](https://github.com/can1357/oh-my-pi/issues/1798)). - `omp usage` now surfaces auto-disabled credentials as red `✗` tombstone rows (identity, how long ago, the shortened upstream cause — e.g. `Refresh token expired` — and a re-login hint), including a provider section when no active credential remains. User-driven tombstones (`replaced by newer credential`, `deleted by user`) and API-key rows stay hidden. Requires a broker with `GET /v1/credentials/disabled`; older brokers degrade to no tombstone rows. - `omp usage` warns about Anthropic's ~30-day OAuth grant lifetime: accounts whose interactive login (`authorizedAt`) is within a week of the deadline get a yellow `⚠ re-login within