diff --git a/apps/sim/lib/credentials/token-service-accounts/descriptors.ts b/apps/sim/lib/credentials/token-service-accounts/descriptors.ts index 3b3de5551d0..e60338c3b71 100644 --- a/apps/sim/lib/credentials/token-service-accounts/descriptors.ts +++ b/apps/sim/lib/credentials/token-service-accounts/descriptors.ts @@ -369,8 +369,6 @@ export const TOKEN_SERVICE_ACCOUNT_DESCRIPTORS: Record< }, ], docsUrl: 'https://docs.sim.ai/integrations/managed-agent', - helpText: - 'Use an Anthropic workspace API key with Managed Agents access. The key is scoped to one Anthropic workspace — its agents, environments, vaults, and memory stores.', }, } diff --git a/apps/sim/lib/mcp/oauth/probe.ts b/apps/sim/lib/mcp/oauth/probe.ts index fc23cd7efbb..b36d83f3b16 100644 --- a/apps/sim/lib/mcp/oauth/probe.ts +++ b/apps/sim/lib/mcp/oauth/probe.ts @@ -33,16 +33,15 @@ export async function detectMcpAuthType( return 'headers' } - // When the caller pre-validated the IP we pin to it directly and own the Agent's - // lifetime; the SSRF-guarded fetch (used when we must validate ourselves) manages its - // own per-request Agent teardown internally. + // Pre-validated IP → pin directly (we own the Agent); otherwise the SSRF-guarded fetch + // self-manages its per-request Agent teardown. const pinned = resolvedIP ? createPinnedFetchWithDispatcher(resolvedIP) : undefined const probeFetch: FetchLike = pinned?.fetch ?? createSsrfGuardedMcpFetch() const controller = new AbortController() const timer = setTimeout(() => controller.abort(), PROBE_TIMEOUT_MS) - // Best-effort session cleanup runs after we return the classification; the pinned - // Agent must outlive it, so we tear it down once cleanup settles. + // Session cleanup runs after we return; the pinned Agent must outlive it, so tear it + // down once cleanup settles. let sessionClose: Promise = Promise.resolve() try { @@ -87,8 +86,7 @@ export async function detectMcpAuthType( return 'headers' } finally { clearTimeout(timer) - // Destroy (not close) so a hung request can't make teardown itself hang; wait for - // the best-effort session-close hop first so it isn't aborted mid-flight. + // destroy() after the session-close hop settles, so that cleanup isn't aborted mid-flight. if (pinned) { void sessionClose.finally(() => pinned.dispatcher.destroy().catch(() => {})) } diff --git a/apps/sim/lib/mcp/pinned-fetch.test.ts b/apps/sim/lib/mcp/pinned-fetch.test.ts index 415a2a1b358..fc873cfffdd 100644 --- a/apps/sim/lib/mcp/pinned-fetch.test.ts +++ b/apps/sim/lib/mcp/pinned-fetch.test.ts @@ -23,11 +23,37 @@ vi.mock('@/lib/mcp/domain-check', () => ({ validateMcpServerSsrf: mockValidateMcpServerSsrf, })) -import { createSsrfGuardedMcpFetch } from '@/lib/mcp/pinned-fetch' +import { createPinnedMcpFetch, createSsrfGuardedMcpFetch } from '@/lib/mcp/pinned-fetch' /** The per-request pinned Agent is always built with a DoS-backstop response cap. */ const withResponseCap = expect.objectContaining({ maxResponseSize: expect.any(Number) }) +describe('createPinnedMcpFetch', () => { + beforeEach(() => { + vi.clearAllMocks() + mockDestroy.mockResolvedValue(undefined) + mockCreatePinnedFetchWithDispatcher.mockReturnValue({ + fetch: sentinelFetch, + dispatcher: { destroy: mockDestroy }, + }) + }) + + it('builds the transport on HTTP/1.1 — never opts into allowH2 (undici h2 stalls)', () => { + const { close } = createPinnedMcpFetch('203.0.113.10') + + // Called with the IP only: no `allowH2`, so the Agent stays on undici's h1.1 default. + expect(mockCreatePinnedFetchWithDispatcher).toHaveBeenCalledWith('203.0.113.10') + const options = mockCreatePinnedFetchWithDispatcher.mock.calls[0][1] as + | { allowH2?: boolean } + | undefined + expect(options?.allowH2).toBeUndefined() + + // close() tears down the pooled sockets (incl. the long-lived SSE) on disconnect. + void close() + expect(mockDestroy).toHaveBeenCalledTimes(1) + }) +}) + describe('createSsrfGuardedMcpFetch', () => { beforeEach(() => { vi.clearAllMocks() diff --git a/apps/sim/lib/mcp/pinned-fetch.ts b/apps/sim/lib/mcp/pinned-fetch.ts index 563296f10c3..5cbc4f55758 100644 --- a/apps/sim/lib/mcp/pinned-fetch.ts +++ b/apps/sim/lib/mcp/pinned-fetch.ts @@ -6,54 +6,41 @@ import { McpError } from '@/lib/mcp/types' /** Pinned fetch for the live MCP transport, plus a handle to release its sockets. */ export interface PinnedMcpFetch { - /** Pinned fetch to hand to the MCP transport's `fetch` option. */ fetch: typeof fetch - /** Tears down the underlying HTTP/2 Agent; call when the MCP client disconnects. */ + /** Tears down the Agent's pooled sockets; call when the MCP client disconnects. */ close: () => Promise } /** - * Pinned fetch for the long-lived MCP transport, which reuses one Agent across - * a connection's requests. MCP servers are commonly behind HTTP/2 fronts (CDNs, - * cloud LBs), and undici's Agent is h1.1-only unless opted into h2 via ALPN, so - * the transport enables it. h2 is *not* used for one-shot flows (OAuth discovery, - * auth-type probe), where a per-request Agent would leave idle h2 sessions with - * no reuse benefit. Pinning is unaffected: the pinned lookup forces the socket to - * `resolvedIP` regardless of negotiated protocol. The returned `close` binds the - * Agent's teardown to the transport lifecycle so h2 sessions don't linger past - * disconnect. + * Pinned fetch for the long-lived MCP transport (one Agent reused per connection). + * + * Runs HTTP/1.1: we do not opt into undici's experimental `allowH2`, whose h2 path stalls + * with headers-but-no-body on reused POST sessions (nodejs/undici #2311, #3433, #4143) — + * the streamable-HTTP `initialize` hang behind a CDN. Both official MCP SDKs use h1.1, and + * the transport has no concurrency to gain from h2. `close` tears down pooled sockets + * (incl. the SSE connection) on disconnect; IP-pinning is protocol-independent. */ export function createPinnedMcpFetch(resolvedIP: string): PinnedMcpFetch { - const { fetch: pinnedFetch, dispatcher } = createPinnedFetchWithDispatcher(resolvedIP, { - allowH2: true, - }) + const { fetch: pinnedFetch, dispatcher } = createPinnedFetchWithDispatcher(resolvedIP) return { fetch: pinnedFetch, close: () => dispatcher.destroy() } } /** - * Per-request deadline for guarded MCP OAuth / RFC 7009 revocation HTTP calls. - * - * The MCP SDK issues OAuth discovery, dynamic client registration, and token - * exchange with a bare `fetch` and no `AbortSignal` — only the JSON-RPC message - * layer gets the SDK's request timeout. Combined with undici's 5-minute default - * headers/body timeouts, a slow or unresponsive authorization server leaves the - * request (and the browser the user is waiting on during `/oauth/start`) pending - * for minutes. Bounding each leg turns that into a fast, actionable failure. 30s - * mirrors `MCP_CLIENT_CONSTANTS.DEFAULT_CONNECTION_TIMEOUT` and leaves wide - * headroom over a healthy server, which completes each leg in well under a second. + * Per-request deadline for guarded OAuth / RFC 7009 legs. The MCP SDK sets no timeout on + * these (only its JSON-RPC layer gets one) and undici's default is 5 minutes, so an + * unresponsive auth server would hang the flow — and the browser waiting on `/oauth/start`. + * 30s mirrors `MCP_CLIENT_CONSTANTS.DEFAULT_CONNECTION_TIMEOUT`. */ const OAUTH_FETCH_TIMEOUT_MS = 30_000 /** - * Cap on a guarded OAuth response body. Discovery/registration/token/refresh/revocation - * replies are always well under 1 KB; this ceiling is purely a DoS backstop so a - * malicious authorization server (reached via attacker-controllable metadata URLs) - * can't stream a multi-GB body within the deadline and exhaust memory. undici rejects - * with `UND_ERR_RES_EXCEEDED_MAX_SIZE` once the decoded body exceeds it. + * DoS backstop for a guarded OAuth response body (real ones are <1 KB): stops a malicious + * auth server — reached via attacker-controllable metadata URLs — from streaming a huge + * body within the deadline. undici rejects with `UND_ERR_RES_EXCEEDED_MAX_SIZE` past it. */ const MAX_OAUTH_RESPONSE_BYTES = 1_048_576 -/** True for undici's response-too-large rejection, however it's wrapped by `fetch`. */ +/** True for undici's response-too-large rejection, however `fetch` wraps it. */ function isResponseTooLarge(error: unknown): boolean { const e = error as { code?: string; cause?: { code?: string } } | null return ( @@ -63,14 +50,9 @@ function isResponseTooLarge(error: unknown): boolean { } /** - * Bounds `promise` by the composed deadline/caller `signal`, rejecting with the signal's - * reason if it aborts first. Holds every guarded phase inside the deadline — the - * `dns.lookup`-based SSRF validation (which takes no signal of its own), the HTTP - * request, and the response body read. Either path attaches a rejection handler to - * `promise` (the pre-aborted `.catch`, or `.then(resolve, reject)`), so a settlement - * arriving after the deadline has fired — once we've stopped awaiting it and are tearing - * the request down — can't surface as an unhandled rejection. The abort listener is - * removed once `promise` settles. + * Bounds `promise` by `signal`, rejecting with its reason on abort — used to hold SSRF + * validation, the request, and the body read inside the deadline. Attaches a rejection + * handler on both paths so a settlement arriving after the deadline can't leak as unhandled. */ function withDeadline(promise: Promise, signal: AbortSignal): Promise { if (signal.aborted) { @@ -94,19 +76,15 @@ function withDeadline(promise: Promise, signal: AbortSignal): Promise { } /** - * Reads a guarded OAuth response fully under the deadline, then returns a detached, - * in-memory copy. The MCP SDK reads discovery/registration/token/refresh bodies - * lazily — AFTER `auth()`'s injected fetch resolves — and applies no timeout of its - * own, so returning the live response would let the body read escape every deadline - * (the "Connecting… forever" hang). Buffering here brings the body read inside the - * wall-clock `signal`: undici's `bodyTimeout` only measures idle gaps between chunks - * and cannot bound a slow-drip or stalled body. These bodies are always small JSON. + * Reads a guarded OAuth response under the deadline and returns a detached in-memory copy. + * The SDK reads these bodies lazily — after our fetch resolves, with no timeout of its own + * — so buffering here is what keeps the body read inside `signal`; undici's `bodyTimeout` + * measures only idle gaps between chunks, not a stalled body. Bodies are always small JSON. */ async function bufferUnderDeadline(response: Response, signal: AbortSignal): Promise { const body = await withDeadline(response.arrayBuffer(), signal) const headers = new Headers(response.headers) - // The buffered copy is decoded and detached from the socket; drop framing headers - // that would misdescribe it. + // Detached + decoded: drop framing headers that would misdescribe the in-memory copy. headers.delete('content-encoding') headers.delete('content-length') const nullBody = response.status === 204 || response.status === 205 || response.status === 304 @@ -118,14 +96,10 @@ async function bufferUnderDeadline(response: Response, signal: AbortSignal): Pro } /** - * Hands a streaming guarded response back live instead of buffering it. The only - * streaming reply a guarded call can see is the auth-type probe's `initialize` - * (`text/event-stream`), and the probe classifies from headers alone — buffering it - * would drain the stream or stall on a server that holds it open (misclassifying auth). - * The per-request pinned Agent (if any) is torn down once the stream ends, the caller - * cancels it, or the deadline aborts, so the socket is never stranded; the `tee` keeps - * the returned body fully readable meanwhile. Teardown ownership moves here, out of the - * caller's `finally`. + * Hands a streaming guarded response (the auth-type probe's `initialize`) back live rather + * than buffering it — buffering would stall on a server that holds the stream open. Tears + * the per-request Agent down once the stream ends, the caller cancels, or the deadline + * aborts; the `tee` keeps the returned body readable meanwhile. */ function releaseStreamOnSettle( response: Response, @@ -143,11 +117,10 @@ function releaseStreamOnSettle( if (signal.aborted) onAbort() else signal.addEventListener('abort', onAbort, { once: true }) try { - while (!(await reader.read()).done) { - // Drain to end-of-stream so the Agent can be torn down once the reply completes. - } + // Drain to end-of-stream so the Agent can be torn down once the reply completes. + while (!(await reader.read()).done) {} } catch { - // Aborted or errored — the teardown below still runs. + // Aborted or errored — teardown still runs in `finally`. } finally { signal.removeEventListener('abort', onAbort) void dispatcher.destroy().catch(() => {}) @@ -162,29 +135,14 @@ function releaseStreamOnSettle( /** * Builds a `FetchLike` for one-shot MCP OAuth calls (discovery, dynamic client - * registration, token exchange/refresh, RFC 7009 revocation). It validates every - * outbound URL against the MCP SSRF policy before issuing it, then pins the - * connection to the resolved IP. Unlike the live transport — where the server URL is - * validated once up front — these hops follow URLs taken verbatim from - * attacker-controllable authorization-server metadata (`authorization_servers`, - * `token_endpoint`, `revocation_endpoint`, …), so each is re-validated and - * private/reserved/loopback targets are rejected (honoring `ALLOWED_MCP_DOMAINS` and - * self-hosted localhost rules). - * - * Three correctness guarantees, each of which the MCP SDK does NOT provide itself (it - * sets no timeout on any OAuth leg and reads bodies lazily): - * - **Bounded end-to-end.** The `timeoutMs` deadline (`AbortSignal.timeout`) composed - * with any caller signal covers SSRF validation, the request, AND the response body - * read — the body is buffered here so the SDK's later read can't outlive the - * deadline. Only our own deadline is relabeled to an `McpError`; a caller abort or - * any other failure propagates unchanged. - * - **No leaked sockets.** The per-request pinned Agent is `destroy()`ed on every path, - * releasing the keep-alive socket a one-shot flow would otherwise strand. - * - **Detached response.** The returned `Response` is an in-memory copy, safe to read - * after the underlying socket is gone. + * registration, token exchange/refresh, RFC 7009 revocation). Each hop's URL comes from + * attacker-controllable authorization-server metadata, so every request is re-validated + * against the SSRF policy and pinned to the resolved IP. * - * A stalled DNS resolution still runs to completion in the background, but its result - * is discarded. + * The SDK provides none of the safety itself (no timeout on OAuth legs, lazy body reads), + * so the guard owns it: the `timeoutMs` deadline covers validation + request + the buffered + * body read; the per-request pinned Agent is `destroy()`ed on every path; and the returned + * `Response` is a detached in-memory copy. Only our own deadline is relabeled to `McpError`. * * @param timeoutMs Per-request deadline in ms (defaults to 30s; override for tests). * @throws McpSsrfError if a request URL resolves to a blocked IP address @@ -194,11 +152,9 @@ export function createSsrfGuardedMcpFetch(timeoutMs: number = OAUTH_FETCH_TIMEOU return (async (url, init) => { const target = typeof url === 'string' ? url : url.href const timeoutSignal = AbortSignal.timeout(timeoutMs) - // Compose deadline + caller signal up front so every phase — SSRF validation, the - // HTTP request, and the body read — is bounded by the deadline and caller cancellation. + // Bound every phase — validation, request, body read — by the deadline + caller signal. const signal = init?.signal ? AbortSignal.any([init.signal, timeoutSignal]) : timeoutSignal - // The per-request pinned Agent MUST be torn down: it holds a keep-alive socket the - // one-shot OAuth flow never reuses, so leaving it open leaks a socket per leg. + // Per-request Agent must be torn down (finally): a one-shot leg never reuses its socket. let dispatcher: Agent | undefined try { const resolvedIP = await withDeadline(validateMcpServerSsrf(target), signal) @@ -213,22 +169,19 @@ export function createSsrfGuardedMcpFetch(timeoutMs: number = OAUTH_FETCH_TIMEOU // No pin (self-hosted allowlist) — global fetch over the shared dispatcher. response = await withDeadline(globalThis.fetch(url, { ...init, signal }), signal) } - // A `text/event-stream` reply is the auth-type probe's `initialize` (the only - // streaming case a guarded call sees); hand it back live with its own teardown so - // the caller reads headers without the buffer draining/stalling the stream. Every - // OAuth leg is single-shot JSON and falls through to be buffered and torn down. + // The probe's `initialize` can stream (text/event-stream); hand it back live so the + // buffer doesn't drain/stall it. Every OAuth leg is single-shot JSON and is buffered. const contentType = response.headers.get('content-type') ?? '' if (contentType.includes('text/event-stream')) { const streamed = releaseStreamOnSettle(response, dispatcher, signal) - dispatcher = undefined // teardown ownership handed to releaseStreamOnSettle + dispatcher = undefined // teardown ownership moved to releaseStreamOnSettle return streamed } return await bufferUnderDeadline(response, signal) } catch (error) { const host = URL.canParse(target) ? new URL(target).host : target - // Relabel only when our own deadline is what fired — identified by the - // rejection reason's identity, not init.signal's state (which may abort - // independently just after the deadline). + // Relabel only our own deadline — by reason identity, not signal state (which may + // abort independently just after the deadline). if (timeoutSignal.aborted && error === timeoutSignal.reason) { throw new McpError(`MCP authorization request to ${host} timed out after ${timeoutMs}ms`) } @@ -239,8 +192,7 @@ export function createSsrfGuardedMcpFetch(timeoutMs: number = OAUTH_FETCH_TIMEOU } throw error } finally { - // Destroy (not close) so a hung leg can't make teardown itself hang; this releases - // the pooled keep-alive socket the per-request Agent would otherwise strand. + // destroy() not close() so a hung leg can't stall teardown; frees the pooled socket. await dispatcher?.destroy().catch(() => {}) } }) satisfies FetchLike