From 6db649c9eda678542a495aeaca9736d95ac5bee3 Mon Sep 17 00:00:00 2001 From: can1357 Date: Wed, 28 Jan 2026 22:00:12 +0100 Subject: [PATCH] feat(coding-agent): added configurable ask timeout/notifications and centralized process execution - Added configurable ask timeout and notification settings to control ask tool behavior and user notifications. - Added AskSettings interface with timeout (in seconds, default 30) and notification method properties to settings manager. - Added ask timeout and notification configuration options to settings UI with values for timeout (off, 15, 30, 60, 120 seconds) and notification methods (auto, bell, osc99, osc9, off). - Refactored process execution across multiple tools (exec, fetch, grep, read, youtube scraper) to use centralized ptree.execText() API instead of custom implementations. - Refactored ChildProcess class in ptree to use public readonly properties and Promise.withResolvers() for cleaner exit handling. --- .../src/config/settings-manager.ts | 36 ++ .../coding-agent/src/exec/bash-executor.ts | 3 +- packages/coding-agent/src/exec/exec.ts | 22 +- .../src/modes/components/settings-defs.ts | 23 ++ .../coding-agent/src/modes/rpc/rpc-client.ts | 2 +- packages/coding-agent/src/ssh/ssh-executor.ts | 4 +- packages/coding-agent/src/tools/ask.ts | 44 ++- packages/coding-agent/src/tools/fetch.ts | 70 +--- packages/coding-agent/src/tools/grep.ts | 138 ++----- packages/coding-agent/src/tools/read.ts | 25 +- .../coding-agent/src/web/scrapers/utils.ts | 17 +- .../coding-agent/src/web/scrapers/youtube.ts | 68 +--- packages/pi-utils/src/ptree.ts | 371 ++++++++++++------ 13 files changed, 444 insertions(+), 379 deletions(-) diff --git a/packages/coding-agent/src/config/settings-manager.ts b/packages/coding-agent/src/config/settings-manager.ts index a5d1aa1c2..7c44bc784 100644 --- a/packages/coding-agent/src/config/settings-manager.ts +++ b/packages/coding-agent/src/config/settings-manager.ts @@ -72,6 +72,13 @@ export interface NotificationSettings { onComplete?: NotificationMethod; // default: "auto" } +export interface AskSettings { + /** Timeout in seconds for ask tool selections (0 or null to disable, default: 30) */ + timeout?: number | null; + /** Notification method when ask tool is waiting for input (default: "auto") */ + notification?: NotificationMethod; +} + export interface ExaSettings { enabled?: boolean; // default: true (master toggle for all Exa tools) enableSearch?: boolean; // default: true (search, deep, code, crawl) @@ -242,6 +249,7 @@ export interface Settings { showHardwareCursor?: boolean; // Show terminal cursor while still positioning it for IME normativeRewrite?: boolean; // default: false (rewrite tool call arguments to normalized format in session history) readLineNumbers?: boolean; // default: false (prepend line numbers to read tool output by default) + ask?: AskSettings; } export const DEFAULT_BASH_INTERCEPTOR_RULES: BashInterceptorRule[] = [ @@ -307,6 +315,7 @@ const DEFAULT_SETTINGS: Settings = { terminal: { showImages: true }, images: { autoResize: true }, notifications: { onComplete: "auto" }, + ask: { timeout: 30, notification: "auto" }, exa: { enabled: true, enableSearch: true, @@ -1117,6 +1126,33 @@ export class SettingsManager { await this.save(); } + /** Get ask tool timeout in milliseconds (0 or null = disabled) */ + getAskTimeout(): number | null { + const timeout = this.settings.ask?.timeout; + if (timeout === null || timeout === 0) return null; + return (timeout ?? 30) * 1000; + } + + async setAskTimeout(seconds: number | null): Promise { + if (!this.globalSettings.ask) { + this.globalSettings.ask = {}; + } + this.globalSettings.ask.timeout = seconds; + await this.save(); + } + + getAskNotification(): NotificationMethod { + return this.settings.ask?.notification ?? "auto"; + } + + async setAskNotification(method: NotificationMethod): Promise { + if (!this.globalSettings.ask) { + this.globalSettings.ask = {}; + } + this.globalSettings.ask.notification = method; + await this.save(); + } + getImageAutoResize(): boolean { return this.settings.images?.autoResize ?? true; } diff --git a/packages/coding-agent/src/exec/bash-executor.ts b/packages/coding-agent/src/exec/bash-executor.ts index 1fbf14629..4d0d936a8 100644 --- a/packages/coding-agent/src/exec/bash-executor.ts +++ b/packages/coding-agent/src/exec/bash-executor.ts @@ -65,9 +65,8 @@ export async function executeBash(command: string, options?: BashExecutorOptions // Wait for process exit try { - await child.exited; return { - exitCode: child.exitCode ?? 0, + exitCode: await child.exited, cancelled: false, ...(await sink.dump()), }; diff --git a/packages/coding-agent/src/exec/exec.ts b/packages/coding-agent/src/exec/exec.ts index d8994900a..cc62069dc 100644 --- a/packages/coding-agent/src/exec/exec.ts +++ b/packages/coding-agent/src/exec/exec.ts @@ -35,22 +35,20 @@ export async function execCommand( cwd: string, options?: ExecOptions, ): Promise { - using proc = ptree.spawnAttached([command, ...args], { + const result = await ptree.execText([command, ...args], { + mode: "attached", cwd, signal: options?.signal, timeout: options?.timeout, + allowNonZero: true, + allowAbort: true, + stderr: "full", }); - // Read streams before awaiting exit to avoid data loss if streams close - const [stdoutText, stderrText] = await Promise.all([proc.stdout.text(), proc.stderr.text()]); - try { - await proc.exited; - } catch { - // ChildProcess rejects on non-zero exit; we handle it below - } + return { - stdout: stdoutText, - stderr: stderrText, - code: proc.exitCode ?? 0, - killed: proc.exitReason instanceof ptree.AbortError, + stdout: result.stdout, + stderr: result.stderr, + code: result.exitCode ?? 0, + killed: Boolean(result.exitError?.aborted), }; } diff --git a/packages/coding-agent/src/modes/components/settings-defs.ts b/packages/coding-agent/src/modes/components/settings-defs.ts index b2593efd6..672ce8677 100644 --- a/packages/coding-agent/src/modes/components/settings-defs.ts +++ b/packages/coding-agent/src/modes/components/settings-defs.ts @@ -189,6 +189,29 @@ export const SETTINGS_DEFS: SettingDef[] = [ get: sm => sm.getNotificationOnComplete(), set: (sm, v) => sm.setNotificationOnComplete(v as NotificationMethod), }, + { + id: "askTimeout", + tab: "behavior", + type: "enum", + label: "Ask tool timeout", + description: "Auto-select recommended option after timeout (disabled in plan mode)", + values: ["off", "15", "30", "60", "120"], + get: sm => { + const timeout = sm.getAskTimeout(); + return timeout === null ? "off" : String(timeout / 1000); + }, + set: (sm, v) => sm.setAskTimeout(v === "off" ? null : Number.parseInt(v, 10)), + }, + { + id: "askNotification", + tab: "behavior", + type: "enum", + label: "Ask notification", + description: "Notify when ask tool is waiting for input", + values: ["auto", "bell", "osc99", "osc9", "off"], + get: sm => sm.getAskNotification(), + set: (sm, v) => sm.setAskNotification(v as NotificationMethod), + }, { id: "startupQuiet", tab: "behavior", diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index 90916a039..5c28fc42a 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -140,7 +140,7 @@ export class RpcClient { await Bun.sleep(100); try { - const exitCode = await Promise.race([this.process.exited, Bun.sleep(50).then(() => null)]); + const exitCode = await Promise.race([this.process.exited, Bun.sleep(500).then(() => null)]); if (exitCode !== null) { throw new Error( `Agent process exited immediately with code ${exitCode}. Stderr: ${this.process.peekStderr()}`, diff --git a/packages/coding-agent/src/ssh/ssh-executor.ts b/packages/coding-agent/src/ssh/ssh-executor.ts index 6a5149e4a..05190e2ea 100644 --- a/packages/coding-agent/src/ssh/ssh-executor.ts +++ b/packages/coding-agent/src/ssh/ssh-executor.ts @@ -92,10 +92,8 @@ export async function executeSSH( ); try { - await child.exited; - const exitCode = child.exitCode ?? 0; return { - exitCode, + exitCode: await child.exited, cancelled: false, ...(await sink.dump()), }; diff --git a/packages/coding-agent/src/tools/ask.ts b/packages/coding-agent/src/tools/ask.ts index 57b91a6af..1baeb1595 100644 --- a/packages/coding-agent/src/tools/ask.ts +++ b/packages/coding-agent/src/tools/ask.ts @@ -12,7 +12,7 @@ * - Users will always be able to select "Other" to provide custom text input * - Use multi: true to allow multiple answers to be selected for a question * - Use recommended: to mark the default option; "(Recommended)" suffix is added automatically - * - Questions time out after 30 seconds and auto-select the recommended option + * - Questions may time out and auto-select the recommended option (configurable, disabled in plan mode) */ import type { AgentTool, AgentToolContext, AgentToolResult, AgentToolUpdateCallback } from "@oh-my-pi/pi-agent-core"; import type { Component } from "@oh-my-pi/pi-tui"; @@ -23,6 +23,7 @@ import type { RenderResultOptions } from "../extensibility/custom-tools/types"; import { type Theme, theme } from "../modes/theme/theme"; import askDescription from "../prompts/tools/ask.md" with { type: "text" }; import { renderStatusLine } from "../tui"; +import { detectNotificationProtocol, isNotificationSuppressed, sendNotification } from "../utils/terminal-notify"; import type { ToolSession } from "."; import { ToolUIKit } from "./render-utils"; @@ -77,7 +78,8 @@ export interface AskToolDetails { const OTHER_OPTION = "Other (type your own)"; const RECOMMENDED_SUFFIX = " (Recommended)"; -const ASK_TIMEOUT_MS = 30000; +/** Default timeout in milliseconds (used when settings unavailable) */ +const DEFAULT_ASK_TIMEOUT_MS = 30000; function getDoneOptionLabel(): string { return `${theme.status.success} Done selecting`; @@ -119,13 +121,20 @@ interface UIContext { input(prompt: string): Promise; } +interface AskQuestionOptions { + /** Timeout in milliseconds, null/undefined to disable */ + timeout?: number | null; +} + async function askSingleQuestion( ui: UIContext, question: string, optionLabels: string[], multi: boolean, recommended?: number, + options?: AskQuestionOptions, ): Promise { + const timeout = options?.timeout ?? undefined; const doneLabel = getDoneOptionLabel(); let selectedOptions: string[] = []; let customInput: string | undefined; @@ -152,11 +161,11 @@ async function askSingleQuestion( const selectionStart = Date.now(); const choice = await ui.select(`${prefix}${question}`, opts, { initialIndex: cursorIndex, - timeout: ASK_TIMEOUT_MS, + timeout: timeout ?? undefined, outline: true, }); const elapsed = Date.now() - selectionStart; - const timedOut = elapsed >= ASK_TIMEOUT_MS; + const timedOut = timeout != null && elapsed >= timeout; if (choice === undefined || choice === doneLabel) break; @@ -198,7 +207,7 @@ async function askSingleQuestion( } else { const displayLabels = addRecommendedSuffix(optionLabels, recommended); const choice = await ui.select(question, [...displayLabels, OTHER_OPTION], { - timeout: ASK_TIMEOUT_MS, + timeout: timeout ?? undefined, initialIndex: recommended, outline: true, }); @@ -254,8 +263,10 @@ export class AskTool implements AgentTool { public readonly label = "Ask"; public readonly description: string; public readonly parameters = askSchema; + private readonly session: ToolSession; - constructor(_session: ToolSession) { + constructor(session: ToolSession) { + this.session = session; this.description = renderPromptTemplate(askDescription); } @@ -263,6 +274,17 @@ export class AskTool implements AgentTool { return session.hasUI ? new AskTool(session) : null; } + /** Send terminal notification when ask tool is waiting for input */ + private sendAskNotification(): void { + if (isNotificationSuppressed()) return; + + const method = this.session.settingsManager?.getAskNotification() ?? "auto"; + if (method === "off") return; + + const protocol = method === "auto" ? detectNotificationProtocol() : method; + sendNotification(protocol, "Waiting for input"); + } + public async execute( _toolCallId: string, params: AskParams, @@ -280,6 +302,14 @@ export class AskTool implements AgentTool { const { ui } = context; + // Determine timeout based on settings and plan mode + const planModeEnabled = this.session.getPlanModeState?.()?.enabled ?? false; + const settingsTimeout = this.session.settingsManager?.getAskTimeout() ?? DEFAULT_ASK_TIMEOUT_MS; + const timeout = planModeEnabled ? null : settingsTimeout; + + // Send notification if waiting and not suppressed + this.sendAskNotification(); + // Multi-part questions mode if (params.questions && params.questions.length > 0) { const results: QuestionResult[] = []; @@ -292,6 +322,7 @@ export class AskTool implements AgentTool { optionLabels, q.multi ?? false, q.recommended, + { timeout }, ); results.push({ @@ -330,6 +361,7 @@ export class AskTool implements AgentTool { optionLabels, multi, params.recommended, + { timeout }, ); const details: AskToolDetails = { diff --git a/packages/coding-agent/src/tools/fetch.ts b/packages/coding-agent/src/tools/fetch.ts index 9eb372e9a..0ee0765bb 100644 --- a/packages/coding-agent/src/tools/fetch.ts +++ b/packages/coding-agent/src/tools/fetch.ts @@ -75,65 +75,6 @@ const CONVERTIBLE_EXTENSIONS = new Set([ // Utilities // ============================================================================= -/** - * Execute a command and return stdout - */ - -type WritableLike = { - write: (chunk: string | Uint8Array) => unknown; - flush?: () => unknown; - end?: () => unknown; -}; - -const textEncoder = new TextEncoder(); - -async function writeStdin(handle: unknown, input: string | Buffer): Promise { - if (!handle || typeof handle === "number") return; - if (typeof (handle as WritableStream).getWriter === "function") { - const writer = (handle as WritableStream).getWriter(); - try { - const chunk = typeof input === "string" ? textEncoder.encode(input) : new Uint8Array(input); - await writer.write(chunk); - } finally { - await writer.close(); - } - return; - } - - const sink = handle as WritableLike; - sink.write(input); - if (sink.flush) sink.flush(); - if (sink.end) sink.end(); -} - -async function exec( - cmd: string, - args: string[], - options?: { timeout?: number; input?: string | Buffer }, -): Promise<{ stdout: string; stderr: string; ok: boolean }> { - using proc = ptree.spawnGroup([cmd, ...args], { - stdin: options?.input ? "pipe" : null, - timeout: options?.timeout ? options.timeout * 1000 : undefined, - }); - - if (options?.input) { - await writeStdin(proc.stdin, options.input); - } - - const [stdout, stderr] = await Promise.all([proc.stdout.text(), proc.stderr.text()]); - try { - await proc.exited; - } catch { - // Handle non-zero exit or timeout - } - - return { - stdout, - stderr, - ok: proc.exitCode === 0, - }; -} - /** * Check if a command exists (cross-platform) */ @@ -456,13 +397,20 @@ async function renderHtmlToText( try { await Bun.write(tmpFile, html); + const execOptions = { + mode: "group" as const, + timeout: timeout * 1000, + allowNonZero: true, + allowAbort: true, + stderr: "full" as const, + }; // Try lynx first (can't auto-install, system package) const lynx = hasCommand("lynx"); if (lynx) { const normalizedPath = tmpFile.replace(/\\/g, "/"); const fileUrl = normalizedPath.startsWith("/") ? `file://${normalizedPath}` : `file:///${normalizedPath}`; - const result = await exec("lynx", ["-dump", "-nolist", "-width", "120", fileUrl], { timeout }); + const result = await ptree.execText(["lynx", "-dump", "-nolist", "-width", "120", fileUrl], execOptions); if (result.ok) { return { content: result.stdout, ok: true, method: "lynx" }; } @@ -471,7 +419,7 @@ async function renderHtmlToText( // Fall back to html2text (auto-install via uv/pip) const html2text = await ensureTool("html2text", true); if (html2text) { - const result = await exec(html2text, [tmpFile], { timeout }); + const result = await ptree.execText([html2text, tmpFile], execOptions); if (result.ok) { return { content: result.stdout, ok: true, method: "html2text" }; } diff --git a/packages/coding-agent/src/tools/grep.ts b/packages/coding-agent/src/tools/grep.ts index d9cbf26d4..d0aa7bcfb 100644 --- a/packages/coding-agent/src/tools/grep.ts +++ b/packages/coding-agent/src/tools/grep.ts @@ -4,8 +4,7 @@ import { StringEnum } from "@oh-my-pi/pi-ai"; import type { Component } from "@oh-my-pi/pi-tui"; import { Text } from "@oh-my-pi/pi-tui"; import { ptree } from "@oh-my-pi/pi-utils"; -import { Type } from "@sinclair/typebox"; -import { $ } from "bun"; +import { type Static, Type } from "@sinclair/typebox"; import { renderPromptTemplate } from "../config/prompt-templates"; import type { RenderResultOptions } from "../extensibility/custom-tools/types"; import type { Theme } from "../modes/theme/theme"; @@ -34,10 +33,7 @@ const grepSchema = Type.Object({ ), i: Type.Optional(Type.Boolean({ description: "Case-insensitive search (default: false)" })), n: Type.Optional(Type.Boolean({ description: "Show line numbers (default: true)" })), - a: Type.Optional(Type.Number({ description: "Lines to show after each match (default: 0)" })), - b: Type.Optional(Type.Number({ description: "Lines to show before each match (default: 0)" })), - c: Type.Optional(Type.Number({ description: "Lines of context (before and after) (default: 0)" })), - context: Type.Optional(Type.Number({ description: "Lines of context (alias for c)" })), + context: Type.Optional(Type.Number({ description: "Lines of context (default: 5)" })), multiline: Type.Optional(Type.Boolean({ description: "Enable multiline matching (default: false)" })), limit: Type.Optional(Type.Number({ description: "Limit output to first N matches (default: 100 in content mode)" })), offset: Type.Optional(Type.Number({ description: "Skip first N entries before applying limit (default: 0)" })), @@ -78,42 +74,28 @@ export async function runRg( args: string[], options?: { signal?: AbortSignal; timeoutMs?: number }, ): Promise { - using child = ptree.spawnAttached([rgPath, ...args], { signal: options?.signal, timeout: options?.timeoutMs }); const timeoutSeconds = options?.timeoutMs ? Math.max(1, Math.round(options.timeoutMs / 1000)) : undefined; const timeoutMessage = timeoutSeconds ? `rg timed out after ${timeoutSeconds}s` : "rg timed out"; - let stdout: string; - try { - stdout = await child.nothrow().text(); - } catch (err) { - if (err instanceof ptree.TimeoutError) { - throw new ToolError(timeoutMessage); - } - if (err instanceof ptree.Exception && err.aborted) { - throw new ToolAbortError(); - } - throw err; - } + const result = await ptree.execText([rgPath, ...args], { + signal: options?.signal, + timeout: options?.timeoutMs, + allowNonZero: true, + allowAbort: true, + stderr: "buffer", + }); - let exitError: unknown; - try { - await child.exited; - } catch (err) { - exitError = err; - if (err instanceof ptree.TimeoutError) { - throw new ToolError(timeoutMessage); - } - if (err instanceof ptree.Exception && err.aborted) { - throw new ToolAbortError(); - } + if (result.exitError instanceof ptree.TimeoutError) { + throw new ToolError(timeoutMessage); + } + if (result.exitError?.aborted) { + throw new ToolAbortError(); } - - const exitCode = child.exitCode ?? (exitError instanceof ptree.Exception ? exitError.exitCode : null); return { - stdout, - stderr: child.peekStderr(), - exitCode, + stdout: result.stdout, + stderr: result.stderr, + exitCode: result.exitCode, }; } @@ -138,22 +120,7 @@ export interface GrepToolOptions { operations?: GrepOperations; } -interface GrepParams { - pattern: string; - path?: string; - glob?: string; - type?: string; - output_mode?: "content" | "files_with_matches" | "count"; - i?: boolean; - n?: boolean; - a?: number; - b?: number; - c?: number; - context?: number; - multiline?: boolean; - limit?: number; - offset?: number; -} +type GrepParams = Static; export class GrepTool implements AgentTool { public readonly name = "grep"; @@ -172,23 +139,25 @@ export class GrepTool implements AgentTool { /** * Validates a pattern against ripgrep's regex engine. - * Uses a quick dry-run against /dev/null to check for parse errors. + * Uses a quick dry-run against the null device to check for parse errors. */ private async validateRegexPattern(pattern: string, rgPath?: string): Promise<{ valid: boolean; error?: string }> { if (!rgPath) { return { valid: true }; // Can't validate, assume valid } - // Run ripgrep against /dev/null with the pattern - this validates regex syntax + // Run ripgrep against the null device with the pattern - this validates regex syntax // without searching any files - const result = await $`${rgPath} --no-config --quiet -- ${pattern} /dev/null`.quiet().nothrow(); - const stderr = result.stderr?.toString() ?? ""; - const exitCode = result.exitCode ?? 0; + const nullDevice = process.platform === "win32" ? "NUL" : "/dev/null"; + const result = await ptree.execText([rgPath, "--no-config", "--quiet", "--", pattern, nullDevice], { + allowNonZero: true, + allowAbort: true, + }); // Exit code 1 = no matches (pattern is valid), 0 = matches found // Exit code 2 = error (often regex parse error) - if (exitCode === 2 && stderr.includes("regex parse error")) { - return { valid: false, error: stderr.trim() }; + if (result.exitCode === 2 && result.stderr.includes("regex parse error")) { + return { valid: false, error: result.stderr.trim() }; } return { valid: true }; @@ -201,22 +170,7 @@ export class GrepTool implements AgentTool { _onUpdate?: AgentToolUpdateCallback, toolContext?: AgentToolContext, ): Promise> { - const { - pattern, - path: searchDir, - glob, - type, - output_mode, - i, - n, - a, - b, - c, - context, - multiline, - limit, - offset, - } = params; + const { pattern, path: searchDir, glob, type, output_mode, i, n, context, multiline, limit, offset } = params; return untilAborted(signal, async () => { const normalizedPattern = pattern.trim(); @@ -244,25 +198,12 @@ export class GrepTool implements AgentTool { return normalized; }; - const normalizedAfter = normalizeContext(a, "After context"); - const normalizedBefore = normalizeContext(b, "Before context"); - const hasContextParam = context !== undefined; - const hasCParam = c !== undefined; - if (hasContextParam && hasCParam) { - throw new ToolError("Cannot combine context with c"); - } - const normalizedContext = normalizeContext(hasContextParam ? context : c, "Context"); - if (normalizedContext > 0 && (normalizedAfter > 0 || normalizedBefore > 0)) { - throw new ToolError("Cannot combine context with a or b"); - } - const contextAfterValue = normalizedContext > 0 ? normalizedContext : normalizedAfter; - const contextBeforeValue = normalizedContext > 0 ? normalizedContext : normalizedBefore; + const normalizedContext = normalizeContext(context ?? 5, "Context"); const showLineNumbers = n ?? true; const ignoreCase = i ?? false; const normalizedGlob = glob?.trim() ?? ""; const normalizedType = type?.trim() ?? ""; - const hasContentHints = - limit !== undefined || context !== undefined || c !== undefined || a !== undefined || b !== undefined; + const hasContentHints = limit !== undefined || context !== undefined; // Validate regex patterns early to surface parse errors before running rg const rgPath = await ensureTool("rg", { @@ -327,6 +268,9 @@ export class GrepTool implements AgentTool { const args: string[] = []; + // Ignore user config files for consistent behavior + args.push("--no-config"); + // Base arguments depend on output mode if (effectiveOutputMode === "files_with_matches") { args.push("--files-with-matches", "--color=never"); @@ -538,8 +482,8 @@ export class GrepTool implements AgentTool { } const block: string[] = []; - const start = contextBeforeValue > 0 ? Math.max(1, lineNumber - contextBeforeValue) : lineNumber; - const end = contextAfterValue > 0 ? Math.min(lines.length, lineNumber + contextAfterValue) : lineNumber; + const start = normalizedContext > 0 ? Math.max(1, lineNumber - normalizedContext) : lineNumber; + const end = normalizedContext > 0 ? Math.min(lines.length, lineNumber + normalizedContext) : lineNumber; for (let current = start; current <= end; current++) { const lineText = lines[current - 1] ?? ""; @@ -601,7 +545,6 @@ export class GrepTool implements AgentTool { } }; - const decoder = new TextDecoder(); let buffer = ""; const parseBuffer = async () => { while (buffer.length > 0) { @@ -631,6 +574,7 @@ export class GrepTool implements AgentTool { }; // Process stdout stream with JSONL chunk parsing + const decoder = new TextDecoder(); try { for await (const chunk of child.stdout) { if (killedDueToLimit) { @@ -655,7 +599,7 @@ export class GrepTool implements AgentTool { // Wait for process to exit try { - await child.exited; + await child.exitedCleanly; } catch (err) { if (err instanceof ptree.Exception) { if (err.aborted) { @@ -741,9 +685,6 @@ interface GrepRenderArgs { type?: string; i?: boolean; n?: boolean; - a?: number; - b?: number; - c?: number; context?: number; multiline?: boolean; output_mode?: string; @@ -764,10 +705,7 @@ export const grepToolRenderer = { if (args.output_mode && args.output_mode !== "files_with_matches") meta.push(`mode:${args.output_mode}`); if (args.i) meta.push("case:insensitive"); if (args.n === false) meta.push("no-line-numbers"); - const contextValue = args.context ?? args.c; - if (contextValue !== undefined && contextValue > 0) meta.push(`context:${contextValue}`); - if (args.a !== undefined && args.a > 0) meta.push(`after:${args.a}`); - if (args.b !== undefined && args.b > 0) meta.push(`before:${args.b}`); + if (args.context !== undefined && args.context > 0) meta.push(`context:${args.context}`); if (args.multiline) meta.push("multiline"); if (args.limit !== undefined && args.limit > 0) meta.push(`limit:${args.limit}`); if (args.offset !== undefined && args.offset > 0) meta.push(`offset:${args.offset}`); diff --git a/packages/coding-agent/src/tools/read.ts b/packages/coding-agent/src/tools/read.ts index b10284304..c6d9ec30a 100644 --- a/packages/coding-agent/src/tools/read.ts +++ b/packages/coding-agent/src/tools/read.ts @@ -296,22 +296,23 @@ async function convertWithMarkitdown( return { content: "", ok: false, error: "markitdown not found (uv/pip unavailable)" }; } - using child = ptree.spawnGroup([cmd, filePath], { signal }); - let stdout: string; - try { - stdout = await child.nothrow().text(); - } catch (err) { - if (err instanceof ptree.Exception && err.aborted) { - throw new ToolAbortError(); - } - throw err; + const result = await ptree.execText([cmd, filePath], { + mode: "group", + signal, + allowNonZero: true, + allowAbort: true, + stderr: "buffer", + }); + + if (result.exitError?.aborted) { + throw new ToolAbortError(); } - if (child.exitCode === 0 && stdout.length > 0) { - return { content: stdout, ok: true }; + if (result.exitCode === 0 && result.stdout.length > 0) { + return { content: result.stdout, ok: true }; } - return { content: "", ok: false, error: child.peekStderr().trim() || "Conversion failed" }; + return { content: "", ok: false, error: result.stderr.trim() || "Conversion failed" }; } const readSchema = Type.Object({ diff --git a/packages/coding-agent/src/web/scrapers/utils.ts b/packages/coding-agent/src/web/scrapers/utils.ts index eb1bf950b..a1bb212e8 100644 --- a/packages/coding-agent/src/web/scrapers/utils.ts +++ b/packages/coding-agent/src/web/scrapers/utils.ts @@ -49,16 +49,21 @@ export async function convertWithMarkitdown( try { await Bun.write(tmpFile, content); - using child = await ptree.spawnGroup([markitdown, tmpFile], { timeout }); - const [stdout, stderr, exitCode] = await Promise.all([child.stdout.text(), child.stderr.text(), child.exited]); - if (exitCode !== 0) { + const result = await ptree.execText([markitdown, tmpFile], { + mode: "group", + timeout, + allowNonZero: true, + stderr: "full", + }); + if (!result.ok) { return { - content: stdout, + content: result.stdout, ok: false, - error: stderr.length > 0 ? stderr : `markitdown failed (exit ${exitCode})`, + error: + result.stderr.length > 0 ? result.stderr : `markitdown failed (exit ${result.exitCode ?? "unknown"})`, }; } - return { content: stdout, ok: true }; + return { content: result.stdout, ok: true }; } finally { try { await fs.rm(tmpFile, { force: true }); diff --git a/packages/coding-agent/src/web/scrapers/youtube.ts b/packages/coding-agent/src/web/scrapers/youtube.ts index c0615597a..c18f93d85 100644 --- a/packages/coding-agent/src/web/scrapers/youtube.ts +++ b/packages/coding-agent/src/web/scrapers/youtube.ts @@ -8,34 +8,6 @@ import { ensureTool } from "../../utils/tools-manager"; import type { RenderResult, SpecialHandler } from "./types"; import { finalizeOutput } from "./types"; -/** - * Execute a command and return stdout - */ -async function exec( - cmd: string, - args: string[], - options?: { timeout?: number; input?: string | Buffer; signal?: AbortSignal }, -): Promise<{ stdout: string; stderr: string; ok: boolean; exitCode: number | null }> { - using proc = ptree.spawnGroup([cmd, ...args], { - signal: options?.signal, - timeout: options?.timeout, - stdin: options?.input ? Buffer.from(options.input) : undefined, - }); - - const [stdout, stderr, exitResult] = await Promise.all([ - proc.stdout.text(), - proc.stderr.text(), - proc.exited.then(() => proc.exitCode ?? 0), - ]); - - return { - stdout, - stderr, - ok: exitResult === 0, - exitCode: exitResult, - }; -} - interface YouTubeUrl { videoId: string; playlistId?: string; @@ -163,16 +135,20 @@ export const handleYouTube: SpecialHandler = async ( const fetchedAt = new Date().toISOString(); const notes: string[] = []; const videoUrl = `https://www.youtube.com/watch?v=${yt.videoId}`; + const execOptions = { + mode: "group" as const, + signal, + timeout: timeout * 1000, + allowNonZero: true, + allowAbort: true, + stderr: "full" as const, + }; // Fetch video metadata throwIfAborted(signal); - const metaResult = await exec( - ytdlp, - ["--dump-json", "--no-warnings", "--no-playlist", "--skip-download", videoUrl], - { - timeout: timeout * 1000, - signal, - }, + const metaResult = await ptree.execText( + [ytdlp, "--dump-json", "--no-warnings", "--no-playlist", "--skip-download", videoUrl], + execOptions, ); throwIfAborted(signal); @@ -215,13 +191,9 @@ export const handleYouTube: SpecialHandler = async ( // First, list available subtitles throwIfAborted(signal); - const listResult = await exec( - ytdlp, - ["--list-subs", "--no-warnings", "--no-playlist", "--skip-download", videoUrl], - { - timeout: timeout * 1000, - signal, - }, + const listResult = await ptree.execText( + [ytdlp, "--list-subs", "--no-warnings", "--no-playlist", "--skip-download", videoUrl], + execOptions, ); throwIfAborted(signal); @@ -236,9 +208,9 @@ export const handleYouTube: SpecialHandler = async ( // Try manual subtitles first (English preferred) if (hasManualSubs) { throwIfAborted(signal); - const subResult = await exec( - ytdlp, + const subResult = await ptree.execText( [ + ytdlp, "--write-sub", "--sub-lang", "en,en-US,en-GB", @@ -251,7 +223,7 @@ export const handleYouTube: SpecialHandler = async ( tmpBase, videoUrl, ], - { timeout: timeout * 1000, signal }, + execOptions, ); if (subResult.ok) { @@ -271,9 +243,9 @@ export const handleYouTube: SpecialHandler = async ( // Fall back to auto-generated captions if (!transcript && hasAutoSubs) { throwIfAborted(signal); - const autoResult = await exec( - ytdlp, + const autoResult = await ptree.execText( [ + ytdlp, "--write-auto-sub", "--sub-lang", "en,en-US,en-GB", @@ -286,7 +258,7 @@ export const handleYouTube: SpecialHandler = async ( tmpBase, videoUrl, ], - { timeout: timeout * 1000, signal }, + execOptions, ); if (autoResult.ok) { diff --git a/packages/pi-utils/src/ptree.ts b/packages/pi-utils/src/ptree.ts index afdd27a96..22f76311b 100644 --- a/packages/pi-utils/src/ptree.ts +++ b/packages/pi-utils/src/ptree.ts @@ -30,9 +30,17 @@ class AsyncQueue { this.#items.push(item); } - close(): void { - if (this.#closed) return; + close(options?: { discard?: boolean }): void { + if (this.#closed) { + if (options?.discard) { + this.#items = []; + } + return; + } this.#closed = true; + if (options?.discard) { + this.#items = []; + } while (this.#resolvers.length > 0) { const resolver = this.#resolvers.shift(); if (resolver) { @@ -54,7 +62,7 @@ class AsyncQueue { } } -function createProcessStream(queue: AsyncQueue): ReadableStream { +function createProcessStream(queue: AsyncQueue, onCancel?: () => void): ReadableStream { const stream = new ReadableStream({ pull: async controller => { const result = await queue.next(); @@ -64,6 +72,10 @@ function createProcessStream(queue: AsyncQueue): ReadableStream { + onCancel?.(); + queue.close({ discard: true }); + }, }); return stream; } @@ -78,11 +90,17 @@ async function killChild(child: ChildProcess) { if (!pid || child.killed) return; const waitForExit = (timeout = 1000) => - Promise.race([Bun.sleep(timeout).then(() => false), child.exited.then(() => true)]); + Promise.race([ + Bun.sleep(timeout).then(() => false), + child.proc.exited.then( + () => true, + () => true, + ), + ]); const sendSignal = async (signal?: NodeJS.Signals) => { try { - process.kill(pid, signal); + child.proc.kill(signal); } catch {} if (child.isProcessGroup) { @@ -111,58 +129,81 @@ postmortem.register("managed-children", async () => { await Promise.all(children.map(killChild)); }); -/** - * Register a subprocess for managed cleanup. - * Will attach to exit Promise so removal happens even if child exits "naturally". - */ -function registerManaged(child: ChildProcess): void { - if (child.exitCode !== null) return; - managedChildren.add(child); - child.exited.finally(() => { - managedChildren.delete(child); - }); -} - // A Bun subprocess with stdin=Writable/ignore, stdout/stderr=pipe (for tracking/cleanup). type PipedSubprocess = Subprocess<"pipe" | "ignore" | null, "pipe", "pipe">; +type StreamReadResult = { done: boolean; value: Uint8Array | undefined }; + +/** + * Options for capturing process output as text. + */ +export interface CaptureTextOptions { + /** Allow non-zero exit codes without throwing. */ + allowNonZero?: boolean; + /** Allow abort/timeout without throwing. */ + allowAbort?: boolean; + /** Select stderr source: full stream or bounded buffer. */ + stderr?: "full" | "buffer"; +} + +/** + * Result from captureText/execText. + */ +export interface CaptureTextResult { + stdout: string; + stderr: string; + exitCode: number | null; + ok: boolean; + exitError?: Exception; +} + /** * ChildProcess wraps a managed subprocess, capturing output, errors, and providing * cross-platform kill/detach logic plus AbortSignal integration. */ export class ChildProcess { - #proc: PipedSubprocess; - #detached = false; - #group = false; #nothrow = false; #stderrBuffer = ""; #stdoutQueue = new AsyncQueue(); #stderrQueue = new AsyncQueue(); + #stderrDone!: Promise; + #streamStop = new AbortController(); + #stdoutActive = true; #stdoutStream?: ReadableStream; #stderrStream?: ReadableStream; #exitReason?: Exception; #exitReasonPending?: Exception; #exited: Promise; - #resolveExited: (ex?: PromiseLike | Exception) => void; - constructor(proc: PipedSubprocess, group: boolean) { - this.#group = group; - registerManaged(this); + constructor( + public readonly proc: PipedSubprocess, + public readonly isProcessGroup: boolean, + ) { + const stopStreaming: Promise = new Promise(resolve => { + if (this.#streamStop.signal.aborted) { + resolve({ done: true, value: undefined }); + return; + } + this.#streamStop.signal.addEventListener( + "abort", + () => { + resolve({ done: true, value: undefined }); + }, + { once: true }, + ); + }); - const exitSettled = proc.exited.then( - () => {}, - () => {}, - ); + const { promise: stderrDone, resolve: resolveStderrDone } = Promise.withResolvers(); + this.#stderrDone = stderrDone; - // Capture stdout at all times. Close the passthrough when the process exits. + // Capture stdout while active. Buffering starts enabled and is disabled when the + // stream is cancelled. The underlying process stdout is always drained to prevent + // the process from blocking on a full pipe buffer. void (async () => { const reader = proc.stdout.getReader(); try { - while (true) { - const result = await Promise.race([ - reader.read(), - exitSettled.then(() => ({ done: true, value: undefined as Uint8Array | undefined })), - ]); + while (this.#stdoutActive) { + const result = await Promise.race([reader.read(), stopStreaming]); if (result.done) break; if (!result.value) continue; this.#stdoutQueue.push(result.value); @@ -188,10 +229,7 @@ export class ChildProcess { const reader = proc.stderr.getReader(); try { while (true) { - const result = await Promise.race([ - reader.read(), - exitSettled.then(() => ({ done: true, value: undefined as Uint8Array | undefined })), - ]); + const result = await Promise.race([reader.read(), stopStreaming]); if (result.done) break; if (!result.value) continue; this.#stderrQueue.push(result.value); @@ -214,71 +252,99 @@ export class ChildProcess { reader.releaseLock(); } catch {} this.#stderrQueue.close(); + resolveStderrDone(); } })().catch(() => { this.#stderrQueue.close(); + resolveStderrDone(); }); - const { promise, resolve } = Promise.withResolvers(); - - this.#exited = promise.then((ex?: Exception) => { - if (!ex) return proc.exitCode ?? -1337; // success, no exception - if (proc.killed && this.#exitReasonPending) { - ex = this.#exitReasonPending; // propagate reason if killed - } - this.#exitReason = ex; - return Promise.reject(ex); - }); - this.#resolveExited = resolve; + const { promise, resolve, reject } = Promise.withResolvers(); + this.#exited = promise; // On exit, resolve with a ChildError if nonzero code. - proc.exited.then(exitCode => { - if (exitCode !== 0) { - resolve(new NonZeroExitError(exitCode, this.#stderrBuffer)); - } else { - resolve(undefined); - } - }); + if (this.proc.exitCode === null) { + managedChildren.add(this); + } + proc.exited + .catch(() => null) + .then(async exitCode => { + // If we have an exit reason pending (e.g., kill() was called), use it immediately. + if (this.#exitReasonPending) { + this.#exitReason = this.#exitReasonPending; + reject(this.#exitReasonPending); + return; + } - this.#proc = proc; + // If successful, resolve as 0. + if (exitCode === 0) { + resolve(0); + return; + } + + // Wait for stderr capture to complete before creating error with stderr content. + await this.#stderrDone; + + let ex: Exception; + if (exitCode !== null) { + this.#exitReason = new NonZeroExitError(exitCode, this.#stderrBuffer); + resolve(exitCode); + return; + } else if (this.proc.killed) { + ex = new AbortError(new Error("process killed"), this.#stderrBuffer); + } else { + ex = new NonZeroExitError(-1, this.#stderrBuffer); + } + this.#exitReason = ex; + reject(ex); + }) + .finally(() => { + managedChildren.delete(this); + }); } - get isProcessGroup(): boolean { - return this.#group; - } get pid(): number | undefined { - return this.#proc.pid; + return this.proc.pid; } get exited(): Promise { return this.#exited; } + get exitedCleanly(): Promise { + if (this.#nothrow) return this.exited; + return this.exited.then(code => { + if (code !== 0) { + throw new NonZeroExitError(code, this.#stderrBuffer); + } + return code; + }); + } get exitCode(): number | null { - return this.#proc.exitCode; + return this.proc.exitCode; } get exitReason(): Exception | undefined { return this.#exitReason; } get killed(): boolean { - return this.#proc.killed; + return this.proc.killed; } get stdin(): FileSink | undefined { - return this.#proc.stdin; + return this.proc.stdin; } get stdout(): ReadableStream { if (!this.#stdoutStream) { - this.#stdoutStream = createProcessStream(this.#stdoutQueue); + this.#stdoutStream = createProcessStream(this.#stdoutQueue, () => { + this.#stdoutActive = false; + }); } return this.#stdoutStream; } get stderr(): ReadableStream { if (!this.#stderrStream) { - this.#stderrStream = createProcessStream(this.#stderrQueue); + // stderr cancellation doesn't affect the internal buffer used for error context + this.#stderrStream = createProcessStream(this.#stderrQueue, () => {}); } return this.#stderrStream; } - get proc(): PipedSubprocess { - return this.#proc; - } /** * Peek at the stderr buffer. @@ -288,15 +354,8 @@ export class ChildProcess { return this.#stderrBuffer; } - /** - * Detach this process from management (no cleanup on shutdown). - */ - detach(): void { - if (this.#detached || this.#proc.killed) return; - this.#detached = true; - if (managedChildren.delete(this)) { - this.#proc.unref(); - } + #requestStreamStop(): void { + this.#streamStop.abort(); } /** @@ -312,10 +371,11 @@ export class ChildProcess { * Optionally set an exit reason (for better error propagation on cancellation). */ kill(reason?: Exception) { - if (this.#proc.killed) return; - if (reason) { + if (reason && !this.#exitReasonPending) { this.#exitReasonPending = reason; } + this.#requestStreamStop(); + if (this.proc.killed) return; killChild(this); } @@ -337,14 +397,53 @@ export class ChildProcess { const blob = this.stdout.blob(); if (!this.#nothrow) { - this.#exited.catch((ex: Exception) => { - reject(ex); - }); + this.exitedCleanly.catch(reject); } blob.then(resolve, reject); return promise; } + /** + * Capture stdout/stderr as text with optional exit handling. + */ + async captureText(options?: CaptureTextOptions): Promise { + const stderrMode = options?.stderr ?? "buffer"; + const stdoutPromise = this.stdout.text(); + const stderrPromise = + stderrMode === "full" + ? this.stderr.text() + : (async () => { + await Promise.allSettled([stdoutPromise, this.exited, this.#stderrDone]); + return this.peekStderr(); + })(); + + const [stdout, stderr] = await Promise.all([stdoutPromise, stderrPromise]); + + let exitError: Exception | undefined; + try { + await this.exited; + } catch (err) { + if (err instanceof Exception) { + exitError = err; + } else { + throw err; + } + } + + const exitCode = this.exitCode ?? (exitError && !exitError.aborted ? exitError.exitCode : null); + const ok = exitCode === 0; + + if (exitError) { + const allowAbort = options?.allowAbort ?? false; + const allowNonZero = options?.allowNonZero ?? false; + if ((exitError.aborted && !allowAbort) || (!exitError.aborted && !allowNonZero)) { + throw exitError; + } + } + + return { stdout, stderr, exitCode, ok, exitError }; + } + /** * Attach an AbortSignal to this process. Will kill tree with SIGKILL if aborted. */ @@ -352,15 +451,6 @@ export class ChildProcess { const onAbort = () => { const cause = new AbortError(signal.reason, ""); this.kill(cause); - if (this.#proc.killed) { - queueMicrotask(() => { - try { - this.#resolveExited(cause); - } catch { - // Ignore - } - }); - } }; if (signal.aborted) { return void onAbort(); @@ -368,10 +458,10 @@ export class ChildProcess { signal.addEventListener("abort", onAbort, { once: true }); // Use .finally().catch() to avoid unhandled rejection when #exited rejects this.#exited + .catch(() => {}) .finally(() => { signal.removeEventListener("abort", onAbort); - }) - .catch(() => {}); + }); } /** @@ -379,15 +469,19 @@ export class ChildProcess { */ attachTimeout(timeout: number): void { if (timeout <= 0) return; - const timeoutId = setTimeout(() => { - this.kill(new TimeoutError(timeout, this.#stderrBuffer)); - }, timeout); - // Use .finally().catch() to avoid unhandled rejection when #exited rejects - this.#exited - .finally(() => { - clearTimeout(timeoutId); - }) - .catch(() => {}); + if (this.proc.killed) return; + void (async () => { + const result = await Promise.race([ + Bun.sleep(timeout).then(() => true), + this.proc.exited.then( + () => false, + () => false, + ), + ]); + if (result) { + this.kill(new TimeoutError(timeout, this.#stderrBuffer)); + } + }); } [Symbol.dispose](): void { @@ -462,22 +556,20 @@ type ChildSpawnOptions = Omit< signal?: AbortSignal; }; -/** - * Spawn a subprocess as a managed child process. - * - Always pipes stdout/stderr, launches in new session/process group (detached). - * - Optional AbortSignal integrates with kill-on-abort. - */ -export function spawnGroup(cmd: string[], options?: ChildSpawnOptions): ChildProcess { +function spawnManaged( + cmd: string[], + options: ChildSpawnOptions | undefined, + config: { detached: boolean; processGroup: boolean }, +): ChildProcess { const { timeout, ...rest } = options ?? {}; const child = spawn(cmd, { stdin: "ignore", ...rest, stdout: "pipe", stderr: "pipe", - // Windows: new console/pgroup; Unix: setsid for process group. - detached: true, + ...(config.detached ? { detached: true } : {}), }); - const cproc = new ChildProcess(child, true); + const cproc = new ChildProcess(child, config.processGroup); if (options?.signal) { cproc.attachSignal(options.signal); } @@ -492,20 +584,43 @@ export function spawnGroup(cmd: string[], options?: ChildSpawnOptions): ChildPro * - Always pipes stdout/stderr, launches in new session/process group (detached). * - Optional AbortSignal integrates with kill-on-abort. */ -export function spawnAttached(cmd: string[], options?: ChildSpawnOptions): ChildProcess { - const { timeout, ...rest } = options ?? {}; - const child = spawn(cmd, { - stdin: "ignore", - ...rest, - stdout: "pipe", - stderr: "pipe", - }); - const cproc = new ChildProcess(child, false); - if (options?.signal) { - cproc.attachSignal(options.signal); - } - if (timeout && timeout > 0) { - cproc.attachTimeout(timeout); - } - return cproc; +export function spawnGroup(cmd: string[], options?: ChildSpawnOptions): ChildProcess { + return spawnManaged(cmd, options, { detached: true, processGroup: true }); +} + +/** + * Spawn a subprocess as a managed child process. + * - Always pipes stdout/stderr, inherits the current session (not detached). + * - Optional AbortSignal integrates with kill-on-abort. + */ +export function spawnAttached(cmd: string[], options?: ChildSpawnOptions): ChildProcess { + return spawnManaged(cmd, options, { detached: false, processGroup: false }); +} + +/** + * Options for execText. + */ +export interface ExecTextOptions extends Omit, CaptureTextOptions { + /** Spawn mode (process group or attached). */ + mode?: "group" | "attached"; + /** Input to write to stdin (Buffer or UTF-8 string). */ + input?: string | Buffer | Uint8Array; +} + +function toStdinBuffer(input: string | Buffer | Uint8Array): Buffer { + if (typeof input === "string") { + return Buffer.from(input); + } + return Buffer.isBuffer(input) ? input : Buffer.from(input); +} + +/** + * Spawn a process and capture stdout/stderr as text. + */ +export async function execText(cmd: string[], options?: ExecTextOptions): Promise { + const { mode = "attached", input, stderr, allowAbort, allowNonZero, ...spawnOptions } = options ?? {}; + const stdin = input === undefined ? undefined : toStdinBuffer(input); + const resolvedOptions: ChildSpawnOptions = stdin === undefined ? { ...spawnOptions } : { ...spawnOptions, stdin }; + using child = mode === "group" ? spawnGroup(cmd, resolvedOptions) : spawnAttached(cmd, resolvedOptions); + return await child.captureText({ stderr, allowAbort, allowNonZero }); }