From 39c08f5434c2d185b91ee18bbd0bc434fbd07b75 Mon Sep 17 00:00:00 2001 From: can1357 Date: Wed, 10 Jun 2026 01:28:03 +0200 Subject: [PATCH] fix(coding-agent): stopped web-search query mangling and API-key log leakage removed the rewrite replacing every 202x with the current year (corrupted CVE ids); MCP request logs redact key/token/secret/auth query params; fetch honors declared charsets, surfaces transport causes, retries 429 once abort-aware, flags mid-stream truncation, stops double-downloading binaries; browser tab reopen/registry/single-flight races fixed, queued opens honor abort, init failures release the temp hold; MCP calls get a default timeout and per-line SSE parse guards. --- packages/coding-agent/src/mcp/json-rpc.ts | 40 +++++++- .../src/tools/browser/registry.ts | 5 +- .../src/tools/browser/tab-supervisor.ts | 54 ++++++++++- packages/coding-agent/src/tools/fetch.ts | 35 +++++-- .../coding-agent/src/web/scrapers/types.ts | 95 +++++++++++++++++-- .../coding-agent/src/web/scrapers/youtube.ts | 7 +- packages/coding-agent/src/web/search/index.ts | 2 +- .../coding-agent/test/mcp-json-rpc.test.ts | 26 +++++ 8 files changed, 237 insertions(+), 27 deletions(-) create mode 100644 packages/coding-agent/test/mcp-json-rpc.test.ts diff --git a/packages/coding-agent/src/mcp/json-rpc.ts b/packages/coding-agent/src/mcp/json-rpc.ts index 6acd3d916..272d8d462 100644 --- a/packages/coding-agent/src/mcp/json-rpc.ts +++ b/packages/coding-agent/src/mcp/json-rpc.ts @@ -6,6 +6,28 @@ */ import { logger } from "@oh-my-pi/pi-utils"; +/** Hard ceiling on a single MCP HTTP request when the caller provides no signal. */ +const MCP_DEFAULT_TIMEOUT_MS = 60_000; + +const SENSITIVE_QUERY_PARAM = /key|token|secret|auth/i; + +/** + * Redact credential-bearing query params (e.g. `exaApiKey`) so failed + * requests never write secrets to the persistent log file. + */ +export function redactUrlForLog(url: string): string { + try { + const parsed = new URL(url); + for (const name of parsed.searchParams.keys()) { + if (SENSITIVE_QUERY_PARAM.test(name)) parsed.searchParams.set(name, "[redacted]"); + } + return parsed.toString(); + } catch { + // Unparseable URL — drop the query string entirely rather than risk leaking it. + return url.split("?")[0]; + } +} + /** Parse SSE response format (lines starting with "data: ") */ export function parseSSE(text: string): unknown { const lines = text.split("\n"); @@ -13,8 +35,12 @@ export function parseSSE(text: string): unknown { if (line.startsWith("data: ")) { const data = line.slice(6).trim(); if (data === "[DONE]") continue; - const result = JSON.parse(data) as unknown; - if (result) return result; + try { + const result = JSON.parse(data) as unknown; + if (result) return result; + } catch { + // Non-JSON data line (keep-alive/comment) — skip and keep scanning. + } } } // Fallback: try parsing entire response as JSON @@ -71,12 +97,12 @@ export async function callMCP( Accept: "application/json, text/event-stream", }, body: JSON.stringify(body), - signal: options?.signal, + signal: options?.signal ?? AbortSignal.timeout(MCP_DEFAULT_TIMEOUT_MS), }); if (!response.ok) { const errorMsg = `MCP request failed: ${response.status} ${response.statusText}`; - logger.error(errorMsg, { url, method, params }); + logger.error(errorMsg, { url: redactUrlForLog(url), method, params }); throw new Error(errorMsg); } @@ -84,7 +110,11 @@ export async function callMCP( const result = parseSSE(text) as JsonRpcResponse | null; if (!result) { - logger.error("Failed to parse MCP response", { url, method, responseText: text.slice(0, 500) }); + logger.error("Failed to parse MCP response", { + url: redactUrlForLog(url), + method, + responseText: text.slice(0, 500), + }); throw new Error("Failed to parse MCP response"); } diff --git a/packages/coding-agent/src/tools/browser/registry.ts b/packages/coding-agent/src/tools/browser/registry.ts index c8caff4c7..59f836a59 100644 --- a/packages/coding-agent/src/tools/browser/registry.ts +++ b/packages/coding-agent/src/tools/browser/registry.ts @@ -157,7 +157,10 @@ export function holdBrowser(handle: BrowserHandle): void { export async function releaseBrowser(handle: BrowserHandle, opts: { kill: boolean }): Promise { handle.refCount = Math.max(0, handle.refCount - 1); if (handle.refCount === 0) { - browsers.delete(handle.key); + // Only evict if the registry still points at THIS handle. After a disconnect, + // `acquireBrowser` may have already replaced the entry with a fresh live handle + // under the same key; deleting blindly would orphan that new browser. + if (browsers.get(handle.key) === handle) browsers.delete(handle.key); await disposeBrowserHandle(handle, opts); } } diff --git a/packages/coding-agent/src/tools/browser/tab-supervisor.ts b/packages/coding-agent/src/tools/browser/tab-supervisor.ts index a73e3e45f..b06649b43 100644 --- a/packages/coding-agent/src/tools/browser/tab-supervisor.ts +++ b/packages/coding-agent/src/tools/browser/tab-supervisor.ts @@ -84,21 +84,51 @@ export interface ReleaseTabOptions { } const tabs = new Map(); +// Per-name acquisition chain: serializes concurrent `acquireTab` calls for the +// same tab name so the existence check and `tabs.set` (separated by several +// awaits) cannot interleave and leak a worker + browser refCount. +const acquireChains = new Map>(); const GRACE_MS = 750; export function getTab(name: string): TabSession | undefined { return tabs.get(name); } -export async function acquireTab( +export function acquireTab(name: string, browser: BrowserHandle, opts: AcquireTabOptions): Promise { + const prior = acquireChains.get(name) ?? Promise.resolve(); + const result = prior.then(() => acquireTabImpl(name, browser, opts)); + const tail = result.then( + () => undefined, + () => undefined, + ); + acquireChains.set(name, tail); + void tail.then(() => { + if (acquireChains.get(name) === tail) acquireChains.delete(name); + }); + return result; +} + +async function acquireTabImpl( name: string, browser: BrowserHandle, opts: AcquireTabOptions, ): Promise { + // Serialized opens can sit behind a slow predecessor in the per-name + // chain; honor an abort at dequeue instead of spawning a worker and + // browser hold nobody is waiting for. + if (opts.signal?.aborted) { + throw new ToolAbortError("Browser tab open aborted"); + } + // Temporary refCount hold so releasing an existing tab on the SAME browser + // below cannot drop it to refCount 0 and dispose the instance we are about + // to reuse (e.g. reopening the sole tab with a different dialogs policy). + let tempHold = false; const existing = tabs.get(name); if (existing) { if (existing.browser === browser && existing.state === "alive") { if (opts.dialogs !== undefined && opts.dialogs !== existing.dialogPolicy) { + holdBrowser(browser); + tempHold = true; await releaseTab(name, { kill: false }); } else { const reuseSteps: string[] = []; @@ -127,12 +157,25 @@ export async function acquireTab( return { tab: tabs.get(name)!, created: false }; } } else { + if (existing.browser === browser) { + holdBrowser(browser); + tempHold = true; + } await releaseTab(name, { kill: false }); } } - const initPayload = await buildInitPayload(browser, opts); - let worker = await spawnTabWorker(); + let initPayload: WorkerInitPayload; + let worker: WorkerHandle; + try { + initPayload = await buildInitPayload(browser, opts); + worker = await spawnTabWorker(); + } catch (error) { + // Failing before the worker took its own hold must release the + // temporary one, or the browser's refCount never reaches 0 again. + if (tempHold || browser.refCount === 0) await releaseBrowser(browser, { kill: false }); + throw error; + } let info: ReadyInfo; try { info = await initializeTabWorker(worker, initPayload, opts.timeoutMs + GRACE_MS); @@ -142,7 +185,7 @@ export async function acquireTab( // the inline worker here so module-resolution failures don't poison every tab open. await worker.terminate().catch(() => undefined); if (worker.mode === "inline") { - if (browser.refCount === 0) await releaseBrowser(browser, { kill: false }); + if (tempHold || browser.refCount === 0) await releaseBrowser(browser, { kill: false }); throw error; } logger.warn("Tab worker init failed; retrying with inline tab worker (no sync-loop guard)", { @@ -153,7 +196,7 @@ export async function acquireTab( info = await initializeTabWorker(worker, initPayload, opts.timeoutMs + GRACE_MS); } catch (inlineError) { await worker.terminate().catch(() => undefined); - if (browser.refCount === 0) await releaseBrowser(browser, { kill: false }); + if (tempHold || browser.refCount === 0) await releaseBrowser(browser, { kill: false }); const finalError = new ToolError( `Failed to start browser tab worker (inline fallback also failed): ${inlineError instanceof Error ? inlineError.message : String(inlineError)}`, ); @@ -163,6 +206,7 @@ export async function acquireTab( } holdBrowser(browser); + if (tempHold) await releaseBrowser(browser, { kill: false }); const tab: TabSession = { name, browser, diff --git a/packages/coding-agent/src/tools/fetch.ts b/packages/coding-agent/src/tools/fetch.ts index 4a97b4ef0..eb6eca3c6 100644 --- a/packages/coding-agent/src/tools/fetch.ts +++ b/packages/coding-agent/src/tools/fetch.ts @@ -23,7 +23,7 @@ import { ensureTool } from "../utils/tools-manager"; import { extractWithParallel, findParallelApiKey, getParallelExtractContent } from "../web/parallel"; import { specialHandlers } from "../web/scrapers"; import type { RenderResult } from "../web/scrapers/types"; -import { finalizeOutput, loadPage, looksLikeHtml, MAX_OUTPUT_CHARS } from "../web/scrapers/types"; +import { finalizeOutput, loadPage, looksLikeHtml, MAX_BYTES, MAX_OUTPUT_CHARS } from "../web/scrapers/types"; import { convertWithMarkit, fetchBinary } from "../web/scrapers/utils"; import { type ArchiveFormat, listArchiveRoot, sniffArchiveFormat } from "./archive-reader"; import { applyListLimit } from "./list-limit"; @@ -191,7 +191,7 @@ export interface ParsedReadUrlTarget { /** Recognize a single selector token (`raw` or one/many line ranges). */ function isUrlSelectorToken(token: string): boolean { - if (token === "raw") return true; + if (token.toLowerCase() === "raw") return true; try { return parseLineRanges(token) !== null; } catch { @@ -213,7 +213,7 @@ export function parseReadUrlTarget(readPath: string): ParsedReadUrlTarget | null let raw = false; let ranges: readonly LineRange[] | undefined; for (const sel of embedded?.sels ?? []) { - if (sel === "raw") { + if (sel.toLowerCase() === "raw") { raw = true; continue; } @@ -805,6 +805,21 @@ function isArchiveHint(mime: string, extensionHint: string): boolean { return ARCHIVE_MIMES.has(mime) || ARCHIVE_EXTENSIONS.has(extensionHint); } +/** + * Content types whose payload renderUrl always re-fetches via fetchBinary. + * Skipping the initial body read for them avoids downloading and + * string-decoding huge binaries (PDFs, archives, images) twice. + */ +function shouldSkipBodyDownload(contentType: string): boolean { + return ( + CONVERTIBLE_MIMES.has(contentType) || + NOTEBOOK_MIMES.has(contentType) || + SQLITE_MIMES.has(contentType) || + ARCHIVE_MIMES.has(contentType) || + SUPPORTED_INLINE_IMAGE_MIME_TYPES.has(contentType) + ); +} + function getArchiveFormatHint(mime: string, extensionHint: string): ArchiveFormat | undefined { if (extensionHint === ".zip" || mime === "application/zip" || mime === "application/x-zip-compressed") { return "zip"; @@ -901,6 +916,7 @@ async function tryRenderBinaryPayload( mime: string, extHint: string, rawContent: string, + bodySkipped: boolean, timeout: number, signal: AbortSignal | undefined, fetchedAt: string, @@ -909,7 +925,7 @@ async function tryRenderBinaryPayload( const hasNotebookHint = isNotebookHint(mime, extHint); const hasSqliteHint = isSqliteHint(mime, extHint); const hasArchiveHint = isArchiveHint(mime, extHint); - const rawLooksBinary = sampleLooksBinary(rawContent); + const rawLooksBinary = bodySkipped || sampleLooksBinary(rawContent); if (!hasNotebookHint && !hasSqliteHint && !hasArchiveHint && !rawLooksBinary) { return null; } @@ -1092,7 +1108,7 @@ async function renderUrl( } // Step 2: Fetch page - const response = await loadPage(url, { timeout, signal }); + const response = await loadPage(url, { timeout, signal, skipBodyForContentType: shouldSkipBodyDownload }); if (signal?.aborted) { throw new ToolAbortError(); } @@ -1105,11 +1121,17 @@ async function renderUrl( content: "", fetchedAt, truncated: false, - notes: [response.status ? `Failed to fetch URL (HTTP ${response.status})` : "Failed to fetch URL"], + notes: [ + response.status ? `Failed to fetch URL (HTTP ${response.status})` : "Failed to fetch URL", + ...(response.error ? [`Cause: ${response.error}`] : []), + ], }; } const { finalUrl, content: rawContent } = response; + if (response.truncated) { + notes.push(`Response body exceeded ${formatBytes(MAX_BYTES)} and was cut mid-stream; content is incomplete`); + } const mime = normalizeMime(response.contentType); const extHint = getExtensionHint(finalUrl); @@ -1276,6 +1298,7 @@ async function renderUrl( mime, extHint, rawContent, + response.bodySkipped === true, timeout, signal, fetchedAt, diff --git a/packages/coding-agent/src/web/scrapers/types.ts b/packages/coding-agent/src/web/scrapers/types.ts index 695575d3f..ae985a74a 100644 --- a/packages/coding-agent/src/web/scrapers/types.ts +++ b/packages/coding-agent/src/web/scrapers/types.ts @@ -1,6 +1,7 @@ /** * Shared types and utilities for web-fetch handlers */ +import { scheduler } from "node:timers/promises"; import { ptree } from "@oh-my-pi/pi-utils"; import type TurndownService from "turndown"; @@ -70,6 +71,12 @@ export interface LoadPageOptions { body?: string; maxBytes?: number; signal?: AbortSignal; + /** + * Return true to skip reading the response body for this content type + * (lowercased mime, no params). The caller is expected to re-fetch the + * payload as binary; this avoids streaming + decoding huge binaries twice. + */ + skipBodyForContentType?: (contentType: string) => boolean; } export interface LoadPageResult { @@ -78,6 +85,51 @@ export interface LoadPageResult { finalUrl: string; ok: boolean; status?: number; + /** True when the body was cut mid-stream at maxBytes. */ + truncated?: boolean; + /** Last transport-level error message when ok is false. */ + error?: string; + /** True when the body read was skipped via skipBodyForContentType. */ + bodySkipped?: boolean; +} + +const RETRY_AFTER_MAX_MS = 10_000; + +/** Parse a Retry-After header (seconds or HTTP-date) into a bounded delay. */ +function parseRetryAfterMs(value: string | null): number { + if (!value) return 1_000; + const seconds = Number(value); + if (Number.isFinite(seconds)) return Math.min(Math.max(seconds, 0) * 1000, RETRY_AFTER_MAX_MS); + const date = Date.parse(value); + if (!Number.isNaN(date)) return Math.min(Math.max(date - Date.now(), 0), RETRY_AFTER_MAX_MS); + return 1_000; +} + +function charsetFromContentType(header: string): string | undefined { + return /charset\s*=\s*"?([\w-]+)"?/i.exec(header)?.[1]; +} + +/** + * Decode a response body honoring the declared charset (Content-Type header, + * then a cheap sniff), falling back to UTF-8. + */ +function decodeBody(bytes: Buffer, contentTypeHeader: string): string { + let label = charsetFromContentType(contentTypeHeader); + if (!label) { + // All charsets we can decode are ASCII-compatible in the prefix, so a + // latin1 view of the first 2KB is enough to find a . + label = /]+charset\s*=\s*["']?([\w-]+)/i.exec(bytes.subarray(0, 2048).toString("latin1"))?.[1]; + } + if (label && !/^utf-?8$/i.test(label)) { + try { + // Bun.Encoding's union is narrower than the runtime, which accepts + // WHATWG labels (shift_jis, euc-kr, gbk, big5, …); unknowns throw here. + return new TextDecoder(label as Bun.Encoding).decode(bytes); + } catch { + // Unknown/unsupported label — fall back to UTF-8. + } + } + return bytes.toString("utf-8"); } /** @@ -86,6 +138,8 @@ export interface LoadPageResult { export async function loadPage(url: string, options: LoadPageOptions = {}): Promise { const { timeout = 20, headers = {}, maxBytes = MAX_BYTES, signal, method = "GET", body } = options; + let lastError: string | undefined; + let retried429 = false; for (let attempt = 0; attempt < USER_AGENTS.length; attempt++) { if (signal?.aborted) { throw new ToolAbortError(); @@ -114,9 +168,31 @@ export async function loadPage(url: string, options: LoadPageOptions = {}): Prom const response = await fetch(url, requestInit); - const contentType = response.headers.get("content-type")?.split(";")[0]?.trim().toLowerCase() ?? ""; + const rawContentType = response.headers.get("content-type") ?? ""; + const contentType = rawContentType.split(";")[0]?.trim().toLowerCase() ?? ""; const finalUrl = response.url; + if (response.status === 429 && !retried429) { + // Rate limited: retry once, honoring a bounded Retry-After. The + // wait observes the caller's signal so an Esc during the backoff + // does not stall for up to the full delay. + retried429 = true; + const delayMs = parseRetryAfterMs(response.headers.get("retry-after")); + void response.body?.cancel().catch(() => {}); + try { + await scheduler.wait(delayMs, { signal }); + } catch { + throw new ToolAbortError(); + } + attempt--; // Reuse the same user agent for the retry. + continue; + } + + if (response.ok && options.skipBodyForContentType?.(contentType)) { + void response.body?.cancel().catch(() => {}); + return { content: "", contentType, finalUrl, ok: true, status: response.status, bodySkipped: true }; + } + const reader = response.body?.getReader(); if (!reader) { return { content: "", contentType, finalUrl, ok: false, status: response.status }; @@ -124,6 +200,7 @@ export async function loadPage(url: string, options: LoadPageOptions = {}): Prom const chunks: Uint8Array[] = []; let totalSize = 0; + let truncated = false; while (true) { const { done, value } = await reader.read(); @@ -133,32 +210,34 @@ export async function loadPage(url: string, options: LoadPageOptions = {}): Prom totalSize += value.length; if (totalSize > maxBytes) { - reader.cancel(); + truncated = true; + void reader.cancel().catch(() => {}); break; } } - const content = Buffer.concat(chunks).toString("utf-8"); + const content = decodeBody(Buffer.concat(chunks), rawContentType); if (isBotBlocked(response.status, content) && attempt < USER_AGENTS.length - 1) { continue; } if (!response.ok) { - return { content, contentType, finalUrl, ok: false, status: response.status }; + return { content, contentType, finalUrl, ok: false, status: response.status, truncated }; } - return { content, contentType, finalUrl, ok: true, status: response.status }; - } catch { + return { content, contentType, finalUrl, ok: true, status: response.status, truncated }; + } catch (error) { if (signal?.aborted) { throw new ToolAbortError(); } + lastError = error instanceof Error ? error.message : String(error); if (attempt === USER_AGENTS.length - 1) { - return { content: "", contentType: "", finalUrl: url, ok: false }; + return { content: "", contentType: "", finalUrl: url, ok: false, error: lastError }; } } } - return { content: "", contentType: "", finalUrl: url, ok: false }; + return { content: "", contentType: "", finalUrl: url, ok: false, error: lastError }; } /** Module-level Turndown instance — built lazily on first use. */ diff --git a/packages/coding-agent/src/web/scrapers/youtube.ts b/packages/coding-agent/src/web/scrapers/youtube.ts index 6dec1276b..b19af0bcc 100644 --- a/packages/coding-agent/src/web/scrapers/youtube.ts +++ b/packages/coding-agent/src/web/scrapers/youtube.ts @@ -288,12 +288,17 @@ export const handleYouTube: SpecialHandler = async ( } } } finally { - throwIfAborted(signal); // Cleanup temp files (fire-and-forget with error suppression) Array.fromAsync(new Bun.Glob(`${tmpBase}*`).scan({ absolute: true })) .then(tmpFiles => Promise.all(tmpFiles.map(f => fs.unlink(f).catch(() => {})))) .catch(() => {}); } + // Only a user-initiated abort is fatal; the per-fetch time budget expiring + // just means partial metadata/transcript, which we surface as a note. + throwIfAborted(userSignal); + if (signal?.aborted) { + notes.push("Fetch time budget exhausted; metadata/transcript may be incomplete"); + } // Build markdown output let md = `# ${title}\n\n`; diff --git a/packages/coding-agent/src/web/search/index.ts b/packages/coding-agent/src/web/search/index.ts index e0ca3e94d..24034a73a 100644 --- a/packages/coding-agent/src/web/search/index.ts +++ b/packages/coding-agent/src/web/search/index.ts @@ -150,7 +150,7 @@ async function executeSearch( lastProvider = provider; try { const response = await provider.search({ - query: params.query.replace(/202\d/g, String(new Date().getFullYear())), // LUL + query: params.query, limit: params.limit, recency: params.recency, systemPrompt: webSearchSystemPrompt, diff --git a/packages/coding-agent/test/mcp-json-rpc.test.ts b/packages/coding-agent/test/mcp-json-rpc.test.ts new file mode 100644 index 000000000..a71480793 --- /dev/null +++ b/packages/coding-agent/test/mcp-json-rpc.test.ts @@ -0,0 +1,26 @@ +import { describe, expect, it } from "bun:test"; +import { parseSSE, redactUrlForLog } from "@oh-my-pi/pi-coding-agent/mcp/json-rpc"; + +describe("redactUrlForLog", () => { + it("redacts credential-bearing query params but keeps the rest", () => { + const redacted = redactUrlForLog("https://mcp.exa.ai/mcp?exaApiKey=sk-secret-123&foo=bar"); + expect(redacted).not.toContain("sk-secret-123"); + expect(redacted).toContain("foo=bar"); + expect(redacted).toContain("https://mcp.exa.ai/mcp"); + }); + + it("drops the query string entirely for unparseable URLs", () => { + expect(redactUrlForLog("not a url?apiKey=zzz")).toBe("not a url"); + }); +}); + +describe("parseSSE", () => { + it("skips non-JSON data lines (keep-alives) and returns the first JSON payload", () => { + const text = 'data: ping\n\ndata: {"jsonrpc":"2.0","id":1,"result":{}}\n'; + expect(parseSSE(text)).toEqual({ jsonrpc: "2.0", id: 1, result: {} }); + }); + + it("returns null when nothing parses", () => { + expect(parseSSE("data: ping\nnot json either")).toBeNull(); + }); +});