diff --git a/package-lock.json b/package-lock.json index 87e25df..918b474 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@aauth/proxy", - "version": "5.6.0", + "version": "5.7.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@aauth/proxy", - "version": "5.6.0", + "version": "5.7.0", "license": "MIT", "dependencies": { "@aauth/call-log": "^0.1.3", diff --git a/package.json b/package.json index 0afb87b..d555326 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@aauth/proxy", - "version": "5.6.0", + "version": "5.7.0", "description": "The user's AAuth agent in MCP form — discovery, identity, interaction relay", "type": "module", "exports": { diff --git a/src/__tests__/connect-resources.test.ts b/src/__tests__/connect-resources.test.ts index aa3a95d..2de733a 100644 --- a/src/__tests__/connect-resources.test.ts +++ b/src/__tests__/connect-resources.test.ts @@ -24,7 +24,8 @@ const mockSignedFetch = vi.fn() vi.mock('@hellocoop/httpsig', () => ({ fetch: mockSignedFetch })) const { McpServer, InMemoryTransport } = await import('@modelcontextprotocol/server') -const { Client } = await import('@modelcontextprotocol/client') +const { serveStdio } = await import('@modelcontextprotocol/server/stdio') +const { Client, UrlElicitationRequiredError } = await import('@modelcontextprotocol/client') const { buildProxyTools } = await import('../tools.js') import type { L1Entry, L1Store } from '../store.js' import type { ProxyConfig } from '../agent.js' @@ -130,6 +131,7 @@ async function connectClient( connectBudgetMs = 30, connectProgressBudgetMs?: number, onInteraction?: (url: string, code: string) => void | Promise, + connectDrainMs?: number, ) { const server = new McpServer({ name: 'test', version: '0.0.0' }) const cfg = makeCfg() @@ -141,6 +143,7 @@ async function connectClient( // A short slice keeps the test fast: only the waiting is cut short. connectBudgetMs, ...(connectProgressBudgetMs !== undefined ? { connectProgressBudgetMs } : {}), + ...(connectDrainMs !== undefined ? { connectDrainMs } : {}), }) const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair() const client = new Client({ name: 'test-client', version: '0.0.0' }) @@ -513,4 +516,268 @@ describe('connect_resources', () => { await close() } }) + + // ── One URL per person server per connect (5.7.0) ── + // + // Prod, 2026-09-28: four Google items, the first two both answered + // `requirement=interaction`. The first call threw a URL elicitation from + // onInteraction before the second item started; the retry started it and + // threw a second elicitation — for a code the person's wallet tab already + // held, since the PS queues every pending interaction per person and drains + // the queue into the open tab. + + /** + * A PS as the wallet behaves: the first `interactions` token requests answer + * `requirement=interaction`, later ones a bare 202 (an open tab is + * reachable). Each pending URL answers per `poll(code, n)` on its nth poll. + */ + function scriptedPS(opts: { interactions: number; poll: (code: string, n: number) => Response }) { + const posted: string[] = [] + const polls = new Map() + let codeSeq = 0 + mockSignedFetch.mockImplementation(async (url: string, init?: { method?: string }) => { + if (url === 'https://ps.example/person') return makeResponse(200, { person_token: 'pt_abc', expires_in: 3600 }) + if (url.endsWith('/connections')) { + if ((init?.method ?? 'GET') === 'POST') { + posted.push(url) + return makeResponse(200, { resource_token: 'rt_conn' }) + } + return makeResponse(200, { connections: [] }) + } + if (url === 'https://ps.example/token') { + codeSeq += 1 + const code = `CODE-${String(codeSeq).padStart(4, '0')}` + return makeResponse(202, {}, { + ...(codeSeq <= opts.interactions ? { 'aauth-requirement': `requirement=interaction; code="${code}"` } : {}), + location: `https://ps.example/pending/${code}`, + }) + } + if (url.startsWith('https://ps.example/pending/')) { + const n = (polls.get(url) ?? 0) + 1 + polls.set(url, n) + return opts.poll(url.slice('https://ps.example/pending/'.length), n) + } + throw new Error(`unexpected signed fetch: ${url}`) + }) + return { posted } + } + + const deferred = (status: 'pending' | 'interacting', position: number) => + makeResponse(202, { status, queue_position: position, queue_depth: 2 }) + + /** + * A client that declares `elicitation.url` and opens every URL it is handed + * (records it), on either protocol era. `modern` serves the connection the + * way the stdio bin does (serveStdio picks the era from the opening + * exchange) and lets the SDK's driver fulfil `input_required` and retry — + * what Claude Code does. + */ + async function elicitingClient( + l1: L1Store, + era: 'modern' | 'legacy', + opts: { drainMs?: number; onInteraction?: (url: string, code: string) => void } = {}, + ) { + const cfg = makeCfg() + const build = async () => { + const server = new McpServer({ name: 'test', version: '0.0.0' }) + await buildProxyTools(server, { + ...(opts.onInteraction ? { onInteraction: opts.onInteraction } : {}), + l1, + registryCache: { read: async () => undefined, write: async () => {} }, + identity: { resolve: async () => ({ kind: 'ready', cfg }), peek: () => cfg }, + connectBudgetMs: 30, + connectProgressBudgetMs: 20_000, + ...(opts.drainMs !== undefined ? { connectDrainMs: opts.drainMs } : {}), + }) + return server + } + const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair() + if (era === 'modern') serveStdio(build, { transport: serverTransport as never }) + else await (await build()).connect(serverTransport) + const client = new Client( + { name: 'test-client', version: '0.0.0' }, + { capabilities: { elicitation: { url: {} } }, versionNegotiation: { mode: era === 'modern' ? 'auto' : 'legacy' } } as never, + ) + const opened: string[] = [] + client.setRequestHandler('elicitation/create' as never, (async (req: { params: { url?: string } }) => { + opened.push(req.params.url ?? '') + return { action: 'accept' } + }) as never) + await client.connect(clientTransport) + return { client, opened, close: async () => void (await client.close()) } + } + + const FOUR = ['gmail.example', 'calendar.example', 'chat.example', 'meet.example'] + + it('2026-07-28: four items, the first two interaction — one input_required with one URL, and the retry lands all four', async () => { + // Gmail and Calendar both need a URL; the person opens Gmail's, the tab + // holds both codes (`interacting`), and the person approves each in turn. + // Chat and Meet start as slots free and the PS reaches the open tab. + const { posted } = scriptedPS({ + interactions: 2, + poll: (code, n) => { + if (code === 'CODE-0001') return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + if (code === 'CODE-0002') return n < 3 ? deferred('interacting', n < 2 ? 2 : 1) : makeResponse(200, {}) + return makeResponse(200, {}) + }, + }) + const handedOver: string[] = [] + const { client, opened, close } = await elicitingClient(memoryL1(FOUR.map((h) => entry(h))), 'modern', { + onInteraction: (_url, code) => void handedOver.push(code), + }) + try { + const result = await client.callTool( + { name: 'connect_resources', arguments: { items: FOUR.map((h) => ({ resource: h, account: 'a@b.co' })) } }, + { onprogress: () => {}, timeout: 20_000 }, + ) + // One URL, the head's — Calendar's code rode in the same tab. + expect(opened).toEqual(['https://ps.example/auth?code=CODE-0001']) + expect(handedOver).toEqual(['CODE-0001']) + // Both live slots were filled before the URL went out: two codes in the + // first call, then one start per item as slots freed. + expect(posted).toEqual(FOUR.map((h) => `https://${h}/connections`)) + const summary = summaryOf(result) + expect(summary.results.map((r) => r.outcome)).toEqual(['connected', 'connected', 'connected', 'connected']) + expect(summary.next).toBeUndefined() + } finally { + await close() + } + }, 20_000) + + it('2025-era: the same list throws one -32042 with one URL; the retry lands all four', async () => { + const { posted } = scriptedPS({ + interactions: 2, + poll: (code, n) => { + if (code === 'CODE-0001' || code === 'CODE-0002') return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + return makeResponse(200, {}) + }, + }) + const { client, close } = await elicitingClient(memoryL1(FOUR.map((h) => entry(h))), 'legacy') + try { + const args = { items: FOUR.map((h) => ({ resource: h, account: 'a@b.co' })) } + const thrown = await client.callTool({ name: 'connect_resources', arguments: args }).catch((e: unknown) => e) + expect(thrown).toBeInstanceOf(UrlElicitationRequiredError) + expect((thrown as InstanceType).elicitations.map((e) => e.url)).toEqual([ + 'https://ps.example/auth?code=CODE-0001', + ]) + expect(posted).toHaveLength(2) + + const result = await client.callTool({ name: 'connect_resources', arguments: args }, { onprogress: () => {}, timeout: 20_000 }) + expect(summaryOf(result).results.map((r) => r.outcome)).toEqual(['connected', 'connected', 'connected', 'connected']) + } finally { + await close() + } + }, 20_000) + + // Two items, both `interaction`; the client has no elicitation, so each URL + // comes back as text and each call is the model's. The first call hands over + // Gmail's URL and covers Calendar's code. + const TWO = FOUR.slice(0, 2) + const twoArgs = { items: TWO.map((h) => ({ resource: h, account: 'a@b.co' })) } + + it('a covered item still `pending` at the head past the drain bound gets its own URL', async () => { + // The person approved Gmail, but no browser ever took Calendar's code (a + // record created without presence, a tab that never connected): the PS + // will not re-advertise it, so the proxy hands its URL over. + scriptedPS({ + interactions: 2, + poll: (code, n) => { + if (code === 'CODE-0001') return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + return deferred('pending', 1) + }, + }) + const { client, close } = await connectClient(memoryL1(TWO.map((h) => entry(h))), 30, 20_000, undefined, 1_500) + try { + const first = textOf(await client.callTool({ name: 'connect_resources', arguments: twoArgs })) + expect(first).toContain('CODE-0001') + expect(first).not.toContain('CODE-0002') + + const started = Date.now() + const second = await client.callTool({ name: 'connect_resources', arguments: twoArgs }, { onprogress: () => {}, timeout: 20_000 }) + const body = textOf(second) + expect(body).toContain('IMPORTANT') + expect(body).toContain('https://ps.example/auth?code=CODE-0002') + expect(summaryOf(second).results.map((r) => r.outcome)).toEqual(['connected', 'still_pending']) + // Not before the bound. + expect(Date.now() - started).toBeGreaterThanOrEqual(1_500) + } finally { + await close() + } + }, 20_000) + + it('a covered item a browser holds (`interacting`) gets no URL of its own, however long it waits', async () => { + scriptedPS({ + interactions: 2, + poll: (code, n) => { + if (code === 'CODE-0001') return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + return deferred('interacting', 1) + }, + }) + const { client, close } = await connectClient(memoryL1(TWO.map((h) => entry(h))), 30, 4_000, undefined, 500) + try { + await client.callTool({ name: 'connect_resources', arguments: twoArgs }) + const second = await client.callTool({ name: 'connect_resources', arguments: twoArgs }, { onprogress: () => {}, timeout: 20_000 }) + expect(textOf(second)).not.toContain('IMPORTANT') + const summary = summaryOf(second) + expect(summary.results.map((r) => r.outcome)).toEqual(['connected', 'still_pending']) + expect(summary.next).toBeTruthy() + } finally { + await close() + } + }, 20_000) + + it('a covered code the PS re-advertises gets its own URL at once', async () => { + // The PS says it cannot reach the person with Calendar's code. That ends + // Calendar's poll the moment it arrives, bound or no bound. + scriptedPS({ + interactions: 2, + poll: (code, n) => + code === 'CODE-0001' + ? n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + : makeResponse(202, { status: 'pending', requirement: 'interaction', code }, { 'aauth-requirement': `requirement=interaction; code="${code}"` }), + }) + const { client, close } = await connectClient(memoryL1(TWO.map((h) => entry(h))), 30, 20_000) + try { + await client.callTool({ name: 'connect_resources', arguments: twoArgs }) + const second = textOf(await client.callTool({ name: 'connect_resources', arguments: twoArgs }, { onprogress: () => {}, timeout: 20_000 })) + expect(second).toContain('IMPORTANT') + expect(second).toContain('https://ps.example/auth?code=CODE-0002') + + // Handed over once: a third call waits on it instead of returning it + // again, though every poll still advertises it. + const third = await client.callTool({ name: 'connect_resources', arguments: twoArgs }) + expect(textOf(third)).not.toContain('IMPORTANT') + expect(summaryOf(third).awaiting?.url).toBe('https://ps.example/auth?code=CODE-0002') + } finally { + await close() + } + }, 30_000) + + it('an item started while a URL is out is covered by it', async () => { + // Three items, all `interaction`: Gmail's URL goes out with Calendar + // covered; when Gmail lands, Chat starts in its slot and is covered by the + // same URL — the wallet tab takes it from the queue. + const { posted } = scriptedPS({ + interactions: 3, + poll: (code, n) => { + if (code === 'CODE-0001') return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + if (code === 'CODE-0002') return n < 3 ? deferred('interacting', 1) : makeResponse(200, {}) + return n < 2 ? deferred('interacting', 1) : makeResponse(200, {}) + }, + }) + const handedOver: string[] = [] + const three = FOUR.slice(0, 3) + const args = { items: three.map((h) => ({ resource: h, account: 'a@b.co' })) } + const { client, close } = await connectClient(memoryL1(three.map((h) => entry(h))), 30, 20_000, (_u, code) => void handedOver.push(code)) + try { + await client.callTool({ name: 'connect_resources', arguments: args }) + const second = await client.callTool({ name: 'connect_resources', arguments: args }, { onprogress: () => {}, timeout: 20_000 }) + expect(textOf(second)).not.toContain('IMPORTANT') + expect(summaryOf(second).results.map((r) => r.outcome)).toEqual(['connected', 'connected', 'connected']) + expect(posted).toHaveLength(3) + expect(handedOver).toEqual(['CODE-0001']) + } finally { + await close() + } + }, 20_000) }) diff --git a/src/__tests__/connect.test.ts b/src/__tests__/connect.test.ts index d25763a..1cf9925 100644 --- a/src/__tests__/connect.test.ts +++ b/src/__tests__/connect.test.ts @@ -184,6 +184,9 @@ describe('pollConnection (B2 slice)', () => { kind: 'still_pending', pollUrl: 'https://ps.example/pending/QQQQ-1111', interaction: { url: 'https://ps.example/auth', code: 'QQQQ-1111', pollUrl: 'https://ps.example/pending/QQQQ-1111' }, + // This poll is the one that advertised it, and the deferred-response status rides along. + advertised: true, + status: 'pending', }) }) @@ -258,9 +261,60 @@ describe('invoke on the approval → interaction fallback (issue #21)', () => { .mockResolvedValueOnce(approval()) .mockResolvedValue(makeResponse(202, { status: 'pending' }, { 'aauth-requirement': 'requirement=interaction; code="C"', location: 'https://ps.example/pending/P' })) const out = await pollConnection(config(), 'https://ps.example/pending/P', 60_000) - expect(out).toEqual({ kind: 'still_pending', pollUrl: 'https://ps.example/pending/P', interaction: { url: 'https://ps.example/auth', code: 'C', pollUrl: 'https://ps.example/pending/P' } }) + expect(out).toEqual({ + kind: 'still_pending', + pollUrl: 'https://ps.example/pending/P', + interaction: { url: 'https://ps.example/auth', code: 'C', pollUrl: 'https://ps.example/pending/P' }, + advertised: true, + status: 'pending', + }) expect(mockSignedFetch).toHaveBeenCalledTimes(2) }) + + it('a poll reports the deferred-response status and queue, and a known code re-advertised is not new', async () => { + // The wallet answers `interacting` once a browser holds the code, with the + // person's queue. A poll that carries no AAuth-Requirement keeps the prior + // interaction but does not mark it advertised. + mockPSWellKnown() + const prior = { url: 'https://ps.example/auth', code: 'C', pollUrl: 'https://ps.example/pending/P' } + mockSignedFetch.mockResolvedValue(makeResponse(202, { status: 'interacting', queue_position: 2, queue_depth: 3 }, { location: 'https://ps.example/pending/P' })) + expect(await pollConnection(config(), prior, 10)).toEqual({ + kind: 'still_pending', + pollUrl: 'https://ps.example/pending/P', + interaction: prior, + status: 'interacting', + queuePosition: 2, + queueDepth: 3, + }) + + // stopOnAdvertise: a code the caller has not handed over ends the slice at once. + mockSignedFetch.mockReset() + mockSignedFetch.mockResolvedValue(makeResponse(202, { status: 'pending' }, { 'aauth-requirement': 'requirement=interaction; code="C"', location: 'https://ps.example/pending/P' })) + const start = Date.now() + const out = await pollConnection(config(), prior, 60_000, undefined, { stopOnAdvertise: true }) + expect(out).toMatchObject({ kind: 'still_pending', advertised: true, interaction: prior }) + expect(mockSignedFetch).toHaveBeenCalledTimes(1) + expect(Date.now() - start).toBeLessThan(1_000) + }) + + it('a re-advertisement without Location is still recognised: the pending is the one polled', async () => { + // Wallet poll.js re-advertises with AAuth-Requirement and Retry-After only. + // Requiring Location here dropped every re-advertised code on the floor. + mockPSWellKnown() + mockSignedFetch.mockResolvedValue( + makeResponse(202, { status: 'pending', requirement: 'interaction', code: 'R', queue_position: 1, queue_depth: 1 }, { 'aauth-requirement': 'requirement=interaction; code="R"' }), + ) + expect(await pollConnection(config(), 'https://ps.example/pending/R', 60_000)).toEqual({ + kind: 'still_pending', + pollUrl: 'https://ps.example/pending/R', + interaction: { url: 'https://ps.example/auth', code: 'R', pollUrl: 'https://ps.example/pending/R' }, + advertised: true, + status: 'pending', + queuePosition: 1, + queueDepth: 1, + }) + expect(mockSignedFetch).toHaveBeenCalledTimes(1) + }) }) describe('listConnections / disconnectAll', () => { diff --git a/src/__tests__/invoke-resume.test.ts b/src/__tests__/invoke-resume.test.ts index f2d202b..a6a6f35 100644 --- a/src/__tests__/invoke-resume.test.ts +++ b/src/__tests__/invoke-resume.test.ts @@ -22,6 +22,7 @@ vi.mock('../resource.js', async (importOriginal) => { }) const { McpServer, InMemoryTransport } = await import('@modelcontextprotocol/server') +const { serveStdio } = await import('@modelcontextprotocol/server/stdio') const { Client } = await import('@modelcontextprotocol/client') const { buildProxyTools } = await import('../tools.js') import type { L1Entry, L1Store } from '../store.js' @@ -221,4 +222,46 @@ describe('invoke resumes a pending authorization', () => { await close() } }) + + it('hands the URL over itself when the host hook returns (5.7.0): one elicitation, then the same code', async () => { + // Before 5.7.0 mcp.aauth.dev threw the elicitation from onInteraction. + // The hook now returns, and invoke elicits natively on a client that + // declared elicitation.url — 2026-07-28 here, as the stdio bin serves it. + const ps = fakePS() + const handedOver: string[] = [] + const cfg = makeCfg() + const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair() + serveStdio(async () => { + const server = new McpServer({ name: 'test', version: '0.0.0' }) + await buildProxyTools(server, { + l1: memoryL1([entry('gmail.example')]), + registryCache: { read: async () => undefined, write: async () => {} }, + identity: { resolve: async () => ({ kind: 'ready', cfg }), peek: () => cfg }, + connectBudgetMs: 30, + onInteraction: (_url, code) => void handedOver.push(code), + }) + return server + }, { transport: serverTransport as never }) + const client = new Client( + { name: 'test-client', version: '0.0.0' }, + { capabilities: { elicitation: { url: {} } }, versionNegotiation: { mode: 'auto' } } as never, + ) + const opened: string[] = [] + client.setRequestHandler('elicitation/create' as never, (async (req: { params: { url?: string } }) => { + opened.push(req.params.url ?? '') + return { action: 'accept' } + }) as never) + await client.connect(clientTransport) + try { + // The driver fulfils the elicitation and retries; the retry resumes the + // same pending (still 202) and says so, without a second elicitation. + const result = textOf(await client.callTool({ name: 'invoke', arguments: { resource: 'gmail.example', op_id: 'whoami' } })) + expect(opened).toEqual(['https://ps.example/auth?code=CODE-1']) + expect(handedOver).toEqual(['CODE-1']) + expect(result).toContain('still in progress') + expect(ps.codes).toBe(1) + } finally { + await client.close() + } + }) }) diff --git a/src/agent.ts b/src/agent.ts index 95111ac..9004cbf 100644 --- a/src/agent.ts +++ b/src/agent.ts @@ -345,9 +345,12 @@ async function terminalChallenge(res: Response, req: ParsedRequirement): Promise // header carries the code only and the agent composes // `{interaction_endpoint}?code=`. A `url=` parameter is still honoured when a // 2.x-era issuer sends one. Neither → not an interaction this agent can drive. -function interactionFrom(res: Response, publishedUrl?: string): Interaction | undefined { +// `polledUrl` is the pending URL a poll answer came from: the PS re-advertises +// a code on the poll without repeating Location (Wallet poll.js), and the +// pending is still the one polled. +function interactionFrom(res: Response, publishedUrl?: string, polledUrl?: string): Interaction | undefined { const parsed = parseRequirement(res.headers.get('aauth-requirement')) - const pollUrl = res.headers.get('location') ?? '' + const pollUrl = res.headers.get('location') ?? polledUrl ?? '' if (parsed?.requirement !== 'interaction' || !parsed.code || !pollUrl) return undefined const url = parsed.url ?? publishedUrl return url ? { url, code: parsed.code, pollUrl } : undefined @@ -382,7 +385,7 @@ async function drivePending( ): Promise<{ kind: 'done'; body: unknown; res: Response } | { kind: 'interaction'; interaction: Interaction } | { kind: 'pending' } | { kind: 'result'; status: number; body: unknown }> { const res = await pollUntilDone(makeAgentPoll(cfg), pollUrl, timeoutMs, undefined, advertisesInteraction) if (res.status === 202) { - const interaction = interactionFrom(res, publishedUrl) + const interaction = interactionFrom(res, publishedUrl, pollUrl) return interaction ? { kind: 'interaction', interaction } : { kind: 'pending' } } if (!res.ok) return { kind: 'result', status: res.status, body: await safeBody(res) } @@ -1412,8 +1415,22 @@ export type ConnectOutcome = * when the person has a URL to open; absent while the PS is reaching them by * its own channels (an open wallet tab, a device) — a later poll may * advertise one. + * + * From `pollConnection`, the PS's own word on the pending: `advertised` when + * the last poll answer carried `requirement=interaction` (the PS could not + * reach the person and wants the URL opened), and the -11 deferred-response + * `status` (`interacting` once a browser holds the code, else `pending`) with + * the person's queue, when the PS sends them. */ - | { kind: 'still_pending'; pollUrl: string; interaction?: Interaction } + | { + kind: 'still_pending' + pollUrl: string + interaction?: Interaction + advertised?: boolean + status?: string + queuePosition?: number + queueDepth?: number + } | { kind: 'error'; status: number; body: unknown } /** @@ -1468,16 +1485,32 @@ export async function connectAtResource(cfg: ProxyConfig, l1: L1Entry, args: Con * advertises, if the PS gave up reaching the person itself), and `error` on a * terminal failure. `pollUrl` may be the URL or a prior `Interaction`. */ -export async function pollConnection(cfg: ProxyConfig, pending: string | Interaction, budgetMs: number, onPoll?: (elapsedMs: number) => void | Promise): Promise { +export async function pollConnection( + cfg: ProxyConfig, + pending: string | Interaction, + budgetMs: number, + onPoll?: (elapsedMs: number) => void | Promise, + opts: { stopOnAdvertise?: boolean } = {}, +): Promise { const pollUrl = typeof pending === 'string' ? pending : pending.pollUrl const prior = typeof pending === 'string' ? undefined : pending - // With no interaction known yet, one the PS starts advertising ends the wait: - // the person needs its URL now, not at the end of the slice. - const res = await pollUntilDone(makeAgentPoll(cfg), pollUrl, budgetMs, onPoll, prior ? undefined : advertisesInteraction) + // An interaction the PS starts advertising ends the wait: the person needs + // its URL now, not at the end of the slice. By default only when no + // interaction is known yet; the caller says otherwise when the URL it knows + // has not reached the person (a PS that re-advertises on every poll would + // otherwise end every slice at once). + const stop = opts.stopOnAdvertise ?? !prior + const res = await pollUntilDone(makeAgentPoll(cfg), pollUrl, budgetMs, onPoll, stop ? advertisesInteraction : undefined) if (res.status === 202) { - const advertised = interactionFrom(res, prior?.url ?? (await psMetadata(cfg.psUrl).catch(() => undefined))?.interaction_endpoint) + const advertised = interactionFrom(res, prior?.url ?? (await psMetadata(cfg.psUrl).catch(() => undefined))?.interaction_endpoint, pollUrl) const interaction = advertised ?? prior - return { kind: 'still_pending', pollUrl, ...(interaction ? { interaction } : {}) } + return { + kind: 'still_pending', + pollUrl, + ...(interaction ? { interaction } : {}), + ...(advertised ? { advertised: true } : {}), + ...deferredFrom(await safeBody(res)), + } } if (res.status >= 200 && res.status < 300) { const access = res.headers.get('aauth-access') @@ -1486,6 +1519,19 @@ export async function pollConnection(cfg: ProxyConfig, pending: string | Interac return { kind: 'error', status: res.status, body: await safeBody(res) } } +// The -11 deferred-response body of a 202 poll answer: `status` and, from a PS +// that queues a person's pendings, where this one sits. Absent fields are left +// out; a body that is not an object yields nothing. +function deferredFrom(body: unknown): { status?: string; queuePosition?: number; queueDepth?: number } { + if (!body || typeof body !== 'object') return {} + const b = body as { status?: unknown; queue_position?: unknown; queue_depth?: unknown } + return { + ...(typeof b.status === 'string' ? { status: b.status } : {}), + ...(typeof b.queue_position === 'number' ? { queuePosition: b.queue_position } : {}), + ...(typeof b.queue_depth === 'number' ? { queueDepth: b.queue_depth } : {}), + } +} + /** * Keep what a settled pending delivered. The PS answers the pending URL with * the token it issued on approval — `person_token` for a person token request, diff --git a/src/tools.ts b/src/tools.ts index b6ecf6e..38c658d 100644 --- a/src/tools.ts +++ b/src/tools.ts @@ -24,12 +24,20 @@ // that do not. Two items are live at a time, finished items are answered from // connectState instead of the resource, and a URL is handed over once, at once. // +// 5.7.0: one URL per person server per connect. Both live items are started +// before any URL goes out (the host's onInteraction no longer ends the call by +// throwing), the head's URL is handed over, and the other live codes are +// covered by it — the PS queues them for the person and the wallet tab works +// through the queue. A covered code gets its own URL only when the PS +// re-advertises it or no browser has taken it from the head of the queue. +// invoke hands its URL over natively too. +// // Transport-agnostic: no fs, no stdio, no child_process. The stdio bin // (server.ts) supplies fs/local-keys deps + a browser-launch onInteraction; // other hosts supply their own backends and surface interaction URLs however // their transport allows. -import { UrlElicitationRequiredError, inputRequired } from '@modelcontextprotocol/server' +import { CLIENT_CAPABILITIES_META_KEY, UrlElicitationRequiredError, inputRequired } from '@modelcontextprotocol/server' import type { ClientCapabilities, Icon, McpServer, RegisteredTool, ServerContext, StandardSchemaWithJSON, ToolAnnotations, ToolCallback } from '@modelcontextprotocol/server' import { renderUnicodeCompact } from 'uqr' import { z } from 'zod' @@ -65,12 +73,18 @@ export interface ProxyDeps { // Optional shared/persistent L3 vocab-doc cache. Defaults (inside resource.ts) // to a process-wide in-memory cache when omitted. docCache?: DocCache - // Called when invoke or connect encounters an interaction (authorization - // URL). May throw to initiate a native protocol-level flow (e.g. MCP URL - // elicitation for cloud hosts). For stdio hosts: open the OS browser and - // return; the tool falls back to returning the URL as text. onComplete is - // called by the host when authorization finishes, resolving any waiters - // registered via authPending. + // Called when invoke or connect_resources hands an authorization URL to the + // person — once per URL handed over, not once per code minted + // (connect_resources covers the codes queued behind it). The tool then hands + // it to the client itself: an MCP URL elicitation when the client declared + // one, else the URL and a QR code as text. For stdio hosts: open the OS + // browser and return. For cloud hosts: arm whatever runs in the background + // and return. onComplete is called by the host when authorization finishes, + // resolving any waiters registered via authPending. + // + // Before 5.7.0 a cloud host threw a URL elicitation from here. A throw is + // still passed through (the URL is recorded as handed over first), but the + // tool's own elicitation is the supported path. onInteraction?: (url: string, code: string, pollUrl: string, onComplete?: () => void | Promise) => void | Promise // Tracks in-flight authorization per resource. Implementations should survive // across MCP session DO instances (e.g. backed by a longer-lived UserStore DO). @@ -100,6 +114,10 @@ export interface ProxyDeps { // progressToken. The call walks the whole list and reports progress as items // land; this only bounds a runaway (each item already times out on its own). connectProgressBudgetMs?: number + // How long an item covered by another item's URL may sit at the head of the + // person's queue with no browser holding it before its own URL is handed + // over. Default DEFAULT_CONNECT_DRAIN_MS; tests shorten it. + connectDrainMs?: number // In-flight connects, per resource host, so a repeat connect_resources call // resumes the same PS pending record instead of starting a new flow. A host // that builds a fresh server per request (the hosted MCP) MUST back this @@ -145,6 +163,16 @@ export interface ConnectFlight { startedAt: number /** The interaction code last handed back to the client, so a resumed call does not hand it back again. */ surfaced?: string + /** + * The code of another item's URL, already handed over, that reaches this + * one too: the PS queues every pending interaction per person, and the + * wallet tab that URL opens works through the queue. This item's own URL is + * handed over only if the PS re-advertises it, or it reaches the head of the + * queue and no browser picks it up (`headAt`). + */ + coveredBy?: string + /** When a poll first found this covered item at the head of the person's queue, not yet held by a browser. */ + headAt?: number } // An item that finished connecting: `connected`, or the resource said it @@ -194,6 +222,20 @@ const MAX_CONNECT_ITEMS = 64 // is already waiting at the PS when the person finishes the one in front of // them, and at most one code counts down behind it. const MAX_LIVE_CONNECTS = 2 +// One URL per person server per connect (5.7.0). The PS queues every pending +// interaction for the person and drains the queue into the tab the first URL +// opened, so a second URL is a second tab for codes the first already holds +// (prod, 2026-09-28: four items, two URL elicitations, the wallet tab already +// held both codes). An item covered by a URL gets its own only when the PS +// re-advertises it, or when it has been at the head of the queue this long +// with no browser holding it (poll `status: 'pending'`, not `'interacting'`). +// The PS sends a waiting tab every queued code when it connects and the tab +// acks them within a second (16:17:19.8 → 16:17:20 in that incident), but a +// poll held open by `Prefer: wait=20` answers with the record as it was when +// the poll arrived — up to 20 s old. Thirty seconds after first seeing the item +// at the head, a `pending` answer is from at least ten seconds after it got +// there: long enough for an open tab to have taken it. +const DEFAULT_CONNECT_DRAIN_MS = 30_000 const BOOTSTRAP_GUIDANCE = `The agent proxy has no AAuth identity on this machine yet. @@ -283,6 +325,7 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi const { l1, registryCache, identity, docCache } = deps const budgetMs = deps.connectBudgetMs ?? DEFAULT_CONNECT_BUDGET_MS const progressBudgetMs = deps.connectProgressBudgetMs ?? DEFAULT_CONNECT_PROGRESS_BUDGET_MS + const drainMs = deps.connectDrainMs ?? DEFAULT_CONNECT_DRAIN_MS // Identity is resolved lazily per call; the provider owns any caching (which // must be per-principal — a shared process-global cache would leak identities @@ -320,8 +363,14 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi const started = Date.now() const fields = toolFields(name, cbArgs.length > 1 ? cbArgs[0] : undefined) try { - const result = (await (handler as (...a: unknown[]) => unknown)(...cbArgs)) as { isError?: boolean } - deps.log?.('tool.call', { ...fields, ok: !result?.isError, duration_ms: Date.now() - started }) + const result = (await (handler as (...a: unknown[]) => unknown)(...cbArgs)) as { isError?: boolean; resultType?: string } + deps.log?.('tool.call', { + ...fields, + ok: !result?.isError, + // A URL elicitation on the 2026-07-28 revision: the client opens it and calls again. + ...(result?.resultType === 'input_required' ? { outcome: 'input_required' } : {}), + duration_ms: Date.now() - started, + }) return result } catch (e) { const error = e instanceof UrlElicitationRequiredError ? 'url_elicitation' : ((e as Error)?.name ?? 'error') @@ -472,14 +521,16 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi // a bare `elicitation:{}`), so a declared `elicitation` key is enough. // Form mode is deliberately not attempted: it is gated on `elicitation.form` // the same way, and handing over a link is not what form mode is for. - function surfaceNatively(ctx: ServerContext, interaction: Interaction, host: string) { + function surfaceNatively(ctx: ServerContext, interaction: Interaction, message: string) { const url = `${interaction.url}?code=${interaction.code}` - const message = `Authorize ${host} — open this URL to connect, then the agent continues.` - // Deprecated accessor, but the supported per-request one: the SDK backfills - // it from the validated envelope on instances that never see an initialize. - let caps: ClientCapabilities | undefined + // On 2026-07-28 the request's own envelope says what the client declared. + // serveStdio does not backfill the instance from it (the stdio bin saw + // `undefined` here and fell back to text); createMcpHandler does. + const envelope = ctx.mcpReq.envelope as Record | undefined + let caps = envelope?.[CLIENT_CAPABILITIES_META_KEY] as ClientCapabilities | undefined + // Deprecated accessor, but the supported per-request one otherwise. try { - caps = server.server.getClientCapabilities() + caps ??= server.server.getClientCapabilities() } catch { return undefined } @@ -578,10 +629,12 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi type Slot = { host: string; item: ConnectItem; row: Record; entry: L1Entry } const rows: Record[] = [] const waiting: Slot[] = [] - // The first item that needs the person at a URL. Only one can be surfaced - // — the PS shows its queue one at a time — and it is the head of that queue. - let toSurface: Interaction | undefined - let surfaceHost: string | undefined + // Items started or polled in this call: only these can need a URL. A + // resumed flight's stored code can die before CONNECT_MAX_MS, so it is + // not handed over until a poll in this call shows it is still pending. + const seen = new Set() + // The deferred-response status of each item's last poll in this call. + const lastStatus = new Map() const rowFor = (entry: L1Entry, account?: string): Record => ({ resource: entry.resource, @@ -665,13 +718,20 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi pollUrl: outcome.pollUrl, ...(outcome.interaction ? { interaction: outcome.interaction } : {}), } + if (outcome.advertised) { + // The PS could not reach the person with this code (its reach + // fallback, or the tab that held it went away) and wants its URL + // opened: no other URL covers it now. + delete next.coveredBy + delete next.headAt + } else if (next.coveredBy && notHeld(outcome.status) && atHead(host, outcome.queuePosition)) { + next.headAt ??= Date.now() + } + seen.add(host) + lastStatus.set(host, outcome.status) await inflight.set(host, next) row.outcome = 'still_pending' row.waiting_on = entry.connection?.upstream_name ?? 'the upstream' - if (next.interaction && !toSurface) { - toSurface = next.interaction - surfaceHost = host - } return } case 'error': { @@ -694,6 +754,82 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi } } + // No browser holds the code. A PS that sends no deferred-response status + // gives no evidence one does, so a covered code there gets its own URL + // after the bound instead of waiting out CONNECT_MAX_MS. + const notHeld = (status: string | undefined): boolean => status === undefined || status === 'pending' + + // Whether a covered item is at the head of the person's queue: from the + // PS's queue_position when it sends one, else no live item of this call + // is ahead of it. + const atHead = (host: string, queuePosition: number | undefined): boolean => { + if (queuePosition !== undefined) return queuePosition <= 1 + const at = waiting.findIndex((w) => w.host === host) + return !waiting.slice(0, Math.max(at, 0)).some((w) => w.row.outcome === 'still_pending') + } + + // The code of a URL already out at this interaction endpoint (one person + // server) that a new code there is covered by: a live item's own handed- + // over URL, or the one covering it. + const coveringCode = async (url: string, except: string): Promise => { + for (const w of waiting) { + if (w.host === except || w.row.outcome !== 'still_pending') continue + const f = await inflight.get(w.host) + if (!f?.interaction || f.interaction.url !== url) continue + if (f.surfaced === f.interaction.code) return f.interaction.code + if (f.coveredBy) return f.coveredBy + } + return undefined + } + + // Whether the person needs this item's own URL now. Not when the client + // has it already, nor while another URL covers it — unless it has sat at + // the head of the queue for drainMs with no browser holding it. + const needsUrl = (host: string, f: InFlight | undefined): f is InFlight & { interaction: Interaction } => { + if (!f?.interaction || !seen.has(host)) return false + if (f.surfaced === f.interaction.code) return false + if (!f.coveredBy) return true + return notHeld(lastStatus.get(host)) && f.headAt !== undefined && Date.now() - f.headAt >= drainMs + } + + // The first live item, in list order, whose URL the person needs. + const nextUrl = async (): Promise<{ host: string; interaction: Interaction } | undefined> => { + for (const w of waiting) { + if (w.row.outcome !== 'still_pending') continue + const f = await inflight.get(w.host) + if (needsUrl(w.host, f)) return { host: w.host, interaction: f.interaction } + } + return undefined + } + + // Record a URL as handed over, and every other live code at the same + // person server as covered by it: the wallet tab it opens is handed the + // rest of the person's queue. + const handOver = async (host: string, interaction: Interaction): Promise => { + const own = await inflight.get(host) + if (own) { + const { coveredBy: _c, headAt: _h, ...rest } = own + await inflight.set(host, { ...rest, surfaced: interaction.code }) + } + for (const w of waiting) { + if (w.host === host || w.row.outcome !== 'still_pending') continue + const f = await inflight.get(w.host) + if (!f?.interaction || f.interaction.url !== interaction.url || f.surfaced === f.interaction.code) continue + const { headAt: _h, ...rest } = f + await inflight.set(w.host, { ...rest, coveredBy: interaction.code }) + } + } + + // A live item whose own URL the client was handed: what the person is on. + const handedLive = async (): Promise<{ resource: string; url: string } | undefined> => { + for (const w of waiting) { + if (w.row.outcome !== 'still_pending') continue + const f = await inflight.get(w.host) + if (f?.interaction && f.surfaced === f.interaction.code) return { resource: w.host, url: `${f.interaction.url}?code=${f.interaction.code}` } + } + return undefined + } + // Start one connect and account for it. Returns true when the item now // holds a live slot (the person, or the PS reaching the person, still has // to act) and has been added to `waiting`; false when it settled on the @@ -713,33 +849,23 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi } if (outcome.kind === 'interaction') { + // Not handed over here: pass 1 starts every live item first, and the + // call hands over one URL at the end (5.7.0). A code started while a + // URL at the same person server is out is covered by it. const { interaction } = outcome - const flight: InFlight = { + const cover = await coveringCode(interaction.url, host) + await inflight.set(host, { pollUrl: interaction.pollUrl, interaction, ...(item.account ? { account: item.account } : {}), ...(item.scopes ? { scopes: item.scopes } : {}), startedAt: Date.now(), - } - await inflight.set(host, flight) - try { - await deps.onInteraction?.(interaction.url, interaction.code, interaction.pollUrl, () => - deps.authPending?.resolve(host), - ) - } catch (err) { - // The host handed the URL over natively (a cloud host throws a URL - // elicitation). Record it, so the retry waits on this code instead - // of handing it over a second time. - await inflight.set(host, { ...flight, surfaced: interaction.code }) - throw err - } + ...(cover ? { coveredBy: cover } : {}), + }) await deps.authPending?.register(host) + seen.add(host) row.outcome = 'still_pending' row.waiting_on = entry.connection?.upstream_name ?? 'the upstream' - if (!toSurface) { - toSurface = interaction - surfaceHost = host - } waiting.push({ host, item, row, entry }) return true } @@ -873,20 +999,14 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi } } - // A URL the person must open that the client has not been handed yet. A - // resumed flight's interaction was handed over by an earlier call. - const unsurfaced = async (): Promise => { - if (!toSurface || !surfaceHost) return false - return (await inflight.get(surfaceHost))?.surfaced !== toSurface.code - } - // Pass 2 — wait on the live items until the list is finished, the // deadline passes, or the person has a URL to open that the client has // not been handed. Each live item is polled for at most a slice, in turn, // so a stuck head does not starve the item behind it and progress goes // out at least once a slice. Each item that lands frees its slot for the - // next held item without a round trip through the model. - while (!(await unsurfaced()) && Date.now() < deadline) { + // next held item without a round trip through the model. A retry after a + // URL was handed over waits here, on the codes that URL covers. + while (!(await nextUrl()) && Date.now() < deadline) { const active = waiting.filter((w) => w.row.outcome === 'still_pending') if (active.length === 0) break for (const { host, item, row, entry } of active) { @@ -902,22 +1022,27 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi row.detail = 'the person did not finish; include this item again to start over' delete row.waiting_on } else { - const polled = await pollConnection(cfg, flight.interaction ?? flight.pollUrl, Math.min(remaining, POLL_SLICE_MS)) + // Stop early on an advertised code unless this item's own URL is + // already out: a covered code the PS advertises needs its URL now. + // A covered code is polled in shorter slices: the drain bound is + // judged from poll answers, and a 25 s slice would push the second + // URL rounds past it. + const handed = !!flight.interaction && flight.surfaced === flight.interaction.code + const slice = flight.coveredBy && !handed ? Math.min(POLL_SLICE_MS, drainMs / 2) : POLL_SLICE_MS + const polled = await pollConnection(cfg, flight.interaction ?? flight.pollUrl, Math.min(remaining, slice), undefined, { + stopOnAdvertise: !handed, + }) await settle(row, entry, item, polled, flight) } if (row.outcome !== 'still_pending') { - // This one landed (or failed): it no longer holds a slot, and it is - // no longer what the person should be looking at. + // This one landed (or failed): it no longer holds a slot. live -= 1 - if (surfaceHost === host) { - toSurface = undefined - surfaceHost = undefined - } await fill() await report(`${host}: ${row.outcome as string}. ${status()}`) } - // A newly started item may need the person at a URL: hand it over now. - if (await unsurfaced()) break + // A newly started item, a re-advertised code, or a covered code no + // browser took may need the person at a URL: hand it over now. + if (await nextUrl()) break } await report(status()) } @@ -936,25 +1061,27 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi : {}), } - // Nothing for the person to open: either everything landed, or the PS is - // reaching them by its own channels (an open wallet tab, a device). - if (pending === 0 || !toSurface || !surfaceHost) return json(summary) - - // A URL an earlier call handed over: the client has shown it already, so - // it rides along for reference rather than as an instruction. - if (!(await unsurfaced())) { - return json({ ...summary, awaiting: { resource: surfaceHost, url: `${toSurface.url}?code=${toSurface.code}` } }) + const surface = pending > 0 ? await nextUrl() : undefined + if (!surface) { + // Nothing new for the person to open: everything landed, the PS is + // reaching them by its own channels (an open wallet tab, a device), or + // the URL is out already — then it rides along for reference rather + // than as an instruction. + const awaiting = pending > 0 ? await handedLive() : undefined + return json(awaiting ? { ...summary, awaiting } : summary) } - // The person must open a URL. Hand it over once: the next call resumes - // the flight and waits instead of returning it again. - const flight = await inflight.get(surfaceHost) - if (flight) await inflight.set(surfaceHost, { ...flight, surfaced: toSurface.code }) - - // Prefer a native prompt; the host may already have taken over (stdio - // opens a browser, a cloud host may throw its own elicitation), in which - // case onInteraction never returned here. - const native = surfaceNatively(ctx, toSurface, surfaceHost) + // The person must open a URL. Hand it over once, and cover the other + // live codes at the same person server with it: the next call waits on + // them instead of handing over another. + const { host: surfaceHost, interaction: toSurface } = surface + await handOver(surfaceHost, toSurface) + // The host's hook: stdio opens a browser; a cloud host arms its + // background poll. A host that throws its own elicitation ends the call + // here, with the URL already recorded as handed over. + await deps.onInteraction?.(toSurface.url, toSurface.code, toSurface.pollUrl, () => deps.authPending?.resolve(surfaceHost)) + + const native = surfaceNatively(ctx, toSurface, `Authorize ${surfaceHost} — open this URL to connect, then the agent continues.`) if (native) return native return text( @@ -1211,18 +1338,23 @@ export async function buildProxyTools(server: McpServer, deps: ProxyDeps): Promi } if (result.kind === 'interaction') { - // Record the flight FIRST: onInteraction may throw (cloud hosts raise - // an MCP URL elicitation) and the retry must find it either way. - await inflight.set(host, { pollUrl: result.interaction.pollUrl, interaction: result.interaction, startedAt: Date.now() }) + // Record the flight FIRST, as handed over: the URL goes to the client + // below (or the host throws its own elicitation), and the retry must + // find the flight either way. + await inflight.set(host, { + pollUrl: result.interaction.pollUrl, + interaction: result.interaction, + startedAt: Date.now(), + surfaced: result.interaction.code, + }) // onComplete resolves the UserStore pending-auth waiter when the poll finishes. const onComplete = () => deps.authPending?.resolve(found.l1.resource) - - // onInteraction may throw (cloud: MCP URL elicitation) or return (stdio/fallback). await deps.onInteraction?.(result.interaction.url, result.interaction.code, result.interaction.pollUrl, onComplete) - - // Only reached if onInteraction returned (fallback path, not elicitation). await deps.authPending?.register(found.l1.resource) + const native = surfaceNatively(ctx, result.interaction, `Authorize ${host} — open this URL, then call invoke again.`) + if (native) return native + return text( `Authorization required for ${found.l1.resource}.\n\n` + `IMPORTANT: You MUST do all of the following in your response:\n` +