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.
This commit is contained in:
can1357
2026-06-10 01:28:03 +02:00
parent c902f0a7d9
commit 39c08f5434
8 changed files with 237 additions and 27 deletions
+35 -5
View File
@@ -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<T = unknown>(
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<T = unknown>(
const result = parseSSE(text) as JsonRpcResponse<T> | 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");
}
@@ -157,7 +157,10 @@ export function holdBrowser(handle: BrowserHandle): void {
export async function releaseBrowser(handle: BrowserHandle, opts: { kill: boolean }): Promise<void> {
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);
}
}
@@ -84,21 +84,51 @@ export interface ReleaseTabOptions {
}
const tabs = new Map<string, TabSession>();
// 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<string, Promise<void>>();
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<AcquireTabResult> {
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<AcquireTabResult> {
// 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,
+29 -6
View File
@@ -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,
@@ -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 <meta charset> 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 <meta charset>.
label = /<meta[^>]+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<LoadPageResult> {
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. */
@@ -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`;
@@ -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,
@@ -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();
});
});