mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
diag(mcp): comprehensive HTTP logging for OAuth + streamable-HTTP transport (#5782)
* diag(mcp): comprehensive HTTP logging for OAuth + streamable-HTTP transport Temporary diagnostic to root-cause the Gauge MCP 'initialize' hang. Wraps both the OAuth guarded fetch and the transport pinned fetch to log every request/response with timing: method, url, status, content-type, mcp-session-id, safe headers, and — for the transport phase only — the response body streamed chunk-by-chunk. So one authorize reveals exactly what Gauge returns for initialize (SSE that stalls vs never-sent result vs fast JSON) and where it stalls. Enabled by default; MCP_HTTP_DIAGNOSTICS=false to silence; disabled under test. Never logs request bodies or the OAuth token-response body (carry tokens); every logged value passes through sanitizeForLogging. Revert once root-caused. * diag(mcp): scope body logging to initialize + redact URLs + wrap unpinned Review fixes (Greptile/Cursor): - Only stream-log the response body whose REQUEST is an MCP initialize JSON-RPC call. The transport fetch also carries in-transport OAuth refresh and every tools/call result, so gating on phase alone leaked token responses and tool output (PII/file contents). initialize gating excludes both. - Redact URL query strings (?code=/?token=/?api_key=) — log origin+path only. - Wrap the transport fetch whether pinned or not (SDK default is globalThis.fetch, so passing it is behavior-neutral) so unpinned/allowlisted servers are covered too. * diag(mcp): sanitize initialize preview + log at warn for visibility Review fixes (Cursor): - Run the initialize body preview through sanitizeForLogging (defense-in-depth against a server stuffing token-like values into serverInfo/capabilities). - Log diagnostics at warn, not info, so they surface even if an environment's LOG_LEVEL drops info (staging runs INFO, but this removes the dependency). * diag(mcp): reuse sanitizeUrlForLog + cancel log reader on abort Review fixes (Cursor): - Replace the local safeUrl helper with the shared sanitizeUrlForLog so URL redaction/truncation stays consistent. - The detached initialize-body log reader now observes init.signal and cancels its tee-branch reader on abort, so it can't keep the response stream / connection alive after the SDK's initialize timeout gives up (the hang we're tracing). * diag(mcp): length-gate initialize check to skip parsing large tool payloads isInitializeRequest now rejects bodies over 4096 chars before any JSON.parse. The initialize message is a small fixed structure, so this keeps large tools/call payloads (potentially multi-MB) off the parse hot path while still detecting initialize.
This commit is contained in:
@@ -381,6 +381,9 @@ describe('McpClient notification handler', () => {
|
||||
{
|
||||
authProvider,
|
||||
requestInit: { headers: { 'X-Sim-Via': 'workflow' } },
|
||||
// The transport fetch is always wrapped for diagnostics (a no-op passthrough
|
||||
// under test); it defaults to globalThis.fetch when the server isn't pinned.
|
||||
fetch: expect.any(Function),
|
||||
}
|
||||
)
|
||||
})
|
||||
|
||||
@@ -12,6 +12,7 @@ import { createLogger } from '@sim/logger'
|
||||
import { getErrorMessage } from '@sim/utils/errors'
|
||||
import { getMaxExecutionTimeout } from '@/lib/core/execution-limits'
|
||||
import { getMcpSafeErrorDiagnostics } from '@/lib/mcp/error-diagnostics'
|
||||
import { withMcpHttpDiagnostics } from '@/lib/mcp/http-diagnostics'
|
||||
import { McpOauthRedirectRequired } from '@/lib/mcp/oauth'
|
||||
import { createPinnedMcpFetch } from '@/lib/mcp/pinned-fetch'
|
||||
import {
|
||||
@@ -101,7 +102,9 @@ export class McpClient {
|
||||
this.transport = new StreamableHTTPClientTransport(new URL(this.config.url), {
|
||||
authProvider: useOauth ? this.authProvider : undefined,
|
||||
requestInit: { headers: this.config.headers },
|
||||
...(pinned ? { fetch: pinned.fetch } : {}),
|
||||
// Wrap whether pinned or not (SDK's default is globalThis.fetch, so passing it is
|
||||
// behavior-neutral) so the diagnostic also covers unpinned/allowlisted servers.
|
||||
fetch: withMcpHttpDiagnostics(pinned?.fetch ?? globalThis.fetch, 'transport'),
|
||||
})
|
||||
|
||||
this.client = new Client(
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
import type { FetchLike } from '@modelcontextprotocol/sdk/shared/transport.js'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { sanitizeForLogging } from '@/lib/core/security/redaction'
|
||||
import { sanitizeUrlForLog } from '@/lib/core/utils/logging'
|
||||
import { getMcpSafeErrorDiagnostics } from '@/lib/mcp/error-diagnostics'
|
||||
|
||||
const logger = createLogger('McpHttpDiag')
|
||||
|
||||
const MAX_LOGGED_CHUNKS = 40
|
||||
const PREVIEW_CHARS = 500
|
||||
|
||||
/**
|
||||
* TEMPORARY diagnostic for the MCP OAuth + streamable-HTTP transport HTTP layer.
|
||||
* Enabled by default; set `MCP_HTTP_DIAGNOSTICS=false` to silence. Remove once the
|
||||
* Gauge `initialize` hang is root-caused.
|
||||
*
|
||||
* Secret-safety (the wrapped transport fetch also carries in-transport OAuth
|
||||
* refresh/registration, and every tool result):
|
||||
* - Request and OAuth response bodies are NEVER logged.
|
||||
* - The only response body streamed is the one whose REQUEST is an MCP `initialize`
|
||||
* JSON-RPC call — which excludes token/refresh responses AND `tools/call` results
|
||||
* (tool output can be PII/file contents/credentials). The `initialize` result is
|
||||
* protocol metadata only (serverInfo, capabilities), no credentials.
|
||||
* - URLs are logged origin+path only; query strings (`?code=`, `?token=`, …) are
|
||||
* redacted. Sensitive headers (authorization/cookie/www-authenticate) are omitted.
|
||||
*/
|
||||
function diagnosticsEnabled(): boolean {
|
||||
if (process.env.NODE_ENV === 'test' || process.env.VITEST) return false
|
||||
return process.env.MCP_HTTP_DIAGNOSTICS !== 'false'
|
||||
}
|
||||
|
||||
function rawUrl(input: string | URL | Request): string {
|
||||
return typeof input === 'string' ? input : input instanceof URL ? input.href : input.url
|
||||
}
|
||||
|
||||
/**
|
||||
* The `initialize` request is a small fixed structure (protocolVersion, capabilities,
|
||||
* clientInfo). This length gate lets us skip parsing large `tools/call` payloads —
|
||||
* which can be multi-MB — entirely, keeping the check off the hot path.
|
||||
*/
|
||||
const MAX_INIT_BODY_CHARS = 4096
|
||||
|
||||
/**
|
||||
* True only when the request body is an MCP `initialize` JSON-RPC message. Used to
|
||||
* scope response-body logging to the initialize handshake and nothing else — token
|
||||
* requests (form-encoded) and tool calls (`method: 'tools/call'`) both fail this.
|
||||
* Large bodies are rejected by length before any parse (see {@link MAX_INIT_BODY_CHARS}).
|
||||
*/
|
||||
function isInitializeRequest(body: unknown): boolean {
|
||||
if (typeof body !== 'string' || body.length > MAX_INIT_BODY_CHARS) return false
|
||||
try {
|
||||
return (JSON.parse(body) as { method?: unknown })?.method === 'initialize'
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wraps a `FetchLike` so every request/response in the MCP OAuth or transport flow is
|
||||
* logged with timing. Only the `initialize` response body is streamed (see secret-safety
|
||||
* note above) — that's the suspected hang.
|
||||
*/
|
||||
export function withMcpHttpDiagnostics(
|
||||
fetchFn: FetchLike,
|
||||
phase: 'oauth' | 'transport'
|
||||
): FetchLike {
|
||||
if (!diagnosticsEnabled()) return fetchFn
|
||||
|
||||
return async (input, init) => {
|
||||
const url = sanitizeUrlForLog(rawUrl(input as string | URL | Request))
|
||||
const method = init?.method ?? 'GET'
|
||||
const reqHeaders = new Headers((init?.headers as HeadersInit | undefined) ?? undefined)
|
||||
const startedAt = Date.now()
|
||||
|
||||
logger.warn('request', {
|
||||
phase,
|
||||
method,
|
||||
url,
|
||||
hasAuth: reqHeaders.has('authorization'),
|
||||
accept: reqHeaders.get('accept') ?? undefined,
|
||||
})
|
||||
|
||||
let res: Response
|
||||
try {
|
||||
res = (await fetchFn(input, init)) as Response
|
||||
} catch (error) {
|
||||
logger.warn('fetch rejected', {
|
||||
phase,
|
||||
method,
|
||||
url,
|
||||
ms: Date.now() - startedAt,
|
||||
error: getMcpSafeErrorDiagnostics(error),
|
||||
})
|
||||
throw error
|
||||
}
|
||||
|
||||
logger.warn('response', {
|
||||
phase,
|
||||
method,
|
||||
url,
|
||||
status: res.status,
|
||||
contentType: res.headers.get('content-type') ?? '',
|
||||
mcpSessionId: res.headers.get('mcp-session-id') ? 'present' : 'absent',
|
||||
headerMs: Date.now() - startedAt,
|
||||
})
|
||||
|
||||
// Stream-log ONLY the initialize handshake response — never OAuth bodies or tool results.
|
||||
if (phase !== 'transport' || !isInitializeRequest(init?.body) || !res.body) return res
|
||||
|
||||
let logBranch: ReadableStream<Uint8Array>
|
||||
let passBranch: ReadableStream<Uint8Array>
|
||||
try {
|
||||
;[logBranch, passBranch] = res.body.tee()
|
||||
} catch {
|
||||
return res
|
||||
}
|
||||
|
||||
// Cancel the detached log reader when the caller aborts (e.g. the SDK's 30s
|
||||
// initialize timeout — exactly the hang we're tracing) so this tee branch can't
|
||||
// keep the response stream / connection alive after the SDK has given up.
|
||||
const signal = init?.signal
|
||||
void (async () => {
|
||||
const reader = logBranch.getReader()
|
||||
const cancelReader = () => void reader.cancel().catch(() => {})
|
||||
if (signal?.aborted) {
|
||||
cancelReader()
|
||||
return
|
||||
}
|
||||
signal?.addEventListener('abort', cancelReader, { once: true })
|
||||
const decoder = new TextDecoder()
|
||||
let chunks = 0
|
||||
let bytes = 0
|
||||
try {
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read()
|
||||
if (done) {
|
||||
logger.warn('initialize body complete', {
|
||||
url,
|
||||
chunks,
|
||||
bytes,
|
||||
ms: Date.now() - startedAt,
|
||||
})
|
||||
break
|
||||
}
|
||||
chunks += 1
|
||||
bytes += value.byteLength
|
||||
if (chunks <= MAX_LOGGED_CHUNKS) {
|
||||
logger.warn('initialize body chunk', {
|
||||
url,
|
||||
chunk: chunks,
|
||||
size: value.byteLength,
|
||||
ms: Date.now() - startedAt,
|
||||
preview: sanitizeForLogging(decoder.decode(value, { stream: true }), PREVIEW_CHARS),
|
||||
})
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
logger.warn('initialize body read error', {
|
||||
url,
|
||||
chunks,
|
||||
bytes,
|
||||
ms: Date.now() - startedAt,
|
||||
error: getMcpSafeErrorDiagnostics(error),
|
||||
})
|
||||
} finally {
|
||||
signal?.removeEventListener('abort', cancelReader)
|
||||
}
|
||||
})()
|
||||
|
||||
return new Response(passBranch, {
|
||||
status: res.status,
|
||||
statusText: res.statusText,
|
||||
headers: res.headers,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
import { auth, type OAuthClientProvider } from '@modelcontextprotocol/sdk/client/auth.js'
|
||||
import { withMcpHttpDiagnostics } from '@/lib/mcp/http-diagnostics'
|
||||
import { createSsrfGuardedMcpFetch } from '@/lib/mcp/pinned-fetch'
|
||||
|
||||
type McpAuthOptions = Parameters<typeof auth>[1]
|
||||
@@ -15,5 +16,8 @@ export function mcpAuthGuarded(
|
||||
provider: OAuthClientProvider,
|
||||
options: McpAuthOptions
|
||||
): ReturnType<typeof auth> {
|
||||
return auth(provider, { ...options, fetchFn: options.fetchFn ?? createSsrfGuardedMcpFetch() })
|
||||
return auth(provider, {
|
||||
...options,
|
||||
fetchFn: withMcpHttpDiagnostics(options.fetchFn ?? createSsrfGuardedMcpFetch(), 'oauth'),
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user