From b2f577f49d1fa68da65a236814abbfacbfdbe944 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 4 Jan 2026 04:25:56 +0100 Subject: [PATCH] refactor(coding-agent/core): migrated to async streams with ScopeSignal - Added ScopeSignal class for composable cancellation with AbortSignal and timeout support. - Refactored bash-executor to use async/await with stream pipelines and ScopeSignal utility. - Replaced manual abort handling with untilAborted wrapper across find, ls, notebook, and read tools. - Simplified error handling by removing manual abort signal listener setup and cleanup logic. --- packages/coding-agent/CHANGELOG.md | 4 + .../coding-agent/src/core/bash-executor.ts | 279 +++++++-------- packages/coding-agent/src/core/tools/find.ts | 335 +++++++++--------- packages/coding-agent/src/core/tools/ls.ts | 176 +++++---- .../coding-agent/src/core/tools/lsp/index.ts | 4 +- .../coding-agent/src/core/tools/notebook.ts | 239 +++++-------- packages/coding-agent/src/core/tools/read.ts | 266 ++++++-------- packages/coding-agent/src/core/utils.ts | 136 +++++++ 8 files changed, 700 insertions(+), 739 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 07590aced..4561e7cca 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -13,6 +13,10 @@ - Simplified FileDiagnosticsResult interface with renamed fields: `diagnostics` → `messages`, `hasErrors` → `errored`, `serverName` → `server` - Session title generation now triggers before sending the first message rather than after agent work begins +### Fixed + +- Fixed potential text decoding issues in bash executor by using streaming TextDecoder instead of Buffer.toString() + ## [3.6.1337] - 2026-01-03 ## [3.5.1337] - 2026-01-03 diff --git a/packages/coding-agent/src/core/bash-executor.ts b/packages/coding-agent/src/core/bash-executor.ts index ef69c7a5d..2d6673a4e 100644 --- a/packages/coding-agent/src/core/bash-executor.ts +++ b/packages/coding-agent/src/core/bash-executor.ts @@ -14,6 +14,7 @@ import stripAnsi from "strip-ansi"; import { getShellConfig, killProcessTree, sanitizeBinaryOutput } from "../utils/shell"; import { getOrCreateSnapshot, getSnapshotSourceCommand } from "../utils/shell-snapshot"; import { DEFAULT_MAX_BYTES, truncateTail } from "./tools/truncate"; +import { ScopeSignal } from "./utils"; // ============================================================================ // Types @@ -47,6 +48,70 @@ export interface BashResult { // Implementation // ============================================================================ +function createSanitizer(): TransformStream { + const decoder = new TextDecoder(); + return new TransformStream({ + transform(chunk, controller) { + const text = sanitizeBinaryOutput(stripAnsi(decoder.decode(chunk, { stream: true }))).replace(/\r/g, ""); + controller.enqueue(text); + }, + }); +} + +function createOutputSink( + spillThreshold: number, + maxBuffer: number, + onChunk?: (text: string) => void, +): WritableStream & { + dump: (annotation?: string) => { output: string; truncated: boolean; fullOutputPath?: string }; +} { + const chunks: string[] = []; + let chunkBytes = 0; + let totalBytes = 0; + let fullOutputPath: string | undefined; + let fullOutputStream: WriteStream | undefined; + + const sink = new WritableStream({ + write(text) { + totalBytes += text.length; + + // Spill to temp file if needed + if (totalBytes > spillThreshold && !fullOutputPath) { + fullOutputPath = join(tmpdir(), `omp-${crypto.randomUUID()}.buffer`); + const ts = createWriteStream(fullOutputPath); + chunks.forEach((c) => { + ts.write(c); + }); + fullOutputStream = ts; + } + fullOutputStream?.write(text); + + // Rolling buffer + chunks.push(text); + chunkBytes += text.length; + while (chunkBytes > maxBuffer && chunks.length > 1) { + chunkBytes -= chunks.shift()!.length; + } + + onChunk?.(text); + }, + close() { + fullOutputStream?.end(); + }, + }); + + return Object.assign(sink, { + dump(annotation?: string) { + if (annotation) { + chunks.push(`\n\n${annotation}`); + } + const full = chunks.join(""); + const { content, truncated } = truncateTail(full); + return { output: truncated ? content : full, truncated, fullOutputPath: fullOutputPath }; + }, + }); +} + /** * Execute a bash command with optional streaming and cancellation support. * @@ -72,165 +137,61 @@ export async function executeBash(command: string, options?: BashExecutorOptions const prefixedCommand = prefix ? `${prefix} ${command}` : command; const finalCommand = `${snapshotPrefix}${prefixedCommand}`; - return new Promise((resolve, reject) => { - const child: Subprocess = Bun.spawn([shell, ...args, finalCommand], { - cwd: options?.cwd, - stdin: "ignore", - stdout: "pipe", - stderr: "pipe", - env, - }); + using signal = new ScopeSignal(options); - // Track sanitized output for truncation - const outputChunks: string[] = []; - let outputBytes = 0; - const maxOutputBytes = DEFAULT_MAX_BYTES * 2; - - // Temp file for large output - let tempFilePath: string | undefined; - let tempFileStream: WriteStream | undefined; - let totalBytes = 0; - let timedOut = false; - - // Handle abort signal and timeout - const abortHandler = () => { - killProcessTree(child.pid); - }; - - // Set up timeout if specified - let timeoutHandle: Timer | undefined; - if (options?.timeout && options.timeout > 0) { - timeoutHandle = setTimeout(() => { - timedOut = true; - abortHandler(); - }, options.timeout); - } - - if (options?.signal) { - if (options.signal.aborted) { - // Already aborted, don't even start - child.kill(); - if (timeoutHandle) clearTimeout(timeoutHandle); - resolve({ - output: "", - exitCode: undefined, - cancelled: true, - truncated: false, - }); - return; - } - options.signal.addEventListener("abort", abortHandler, { once: true }); - } - - const handleData = (data: Buffer) => { - totalBytes += data.length; - - // Sanitize once at the source: strip ANSI, replace binary garbage, normalize newlines - const text = sanitizeBinaryOutput(stripAnsi(data.toString())).replace(/\r/g, ""); - - // Start writing to temp file if exceeds threshold - if (totalBytes > DEFAULT_MAX_BYTES && !tempFilePath) { - const randomId = crypto.getRandomValues(new Uint8Array(8)); - const id = Array.from(randomId, (b) => b.toString(16).padStart(2, "0")).join(""); - tempFilePath = join(tmpdir(), `omp-bash-${id}.log`); - tempFileStream = createWriteStream(tempFilePath); - // Write already-buffered chunks to temp file - for (const chunk of outputChunks) { - tempFileStream.write(chunk); - } - } - - if (tempFileStream) { - tempFileStream.write(text); - } - - // Keep rolling buffer of sanitized text - outputChunks.push(text); - outputBytes += text.length; - while (outputBytes > maxOutputBytes && outputChunks.length > 1) { - const removed = outputChunks.shift()!; - outputBytes -= removed.length; - } - - // Stream to callback if provided - if (options?.onChunk) { - options.onChunk(text); - } - }; - - // Read streams asynchronously - (async () => { - try { - const stdoutReader = (child.stdout as ReadableStream).getReader(); - const stderrReader = (child.stderr as ReadableStream).getReader(); - - await Promise.all([ - (async () => { - while (true) { - const { done, value } = await stdoutReader.read(); - if (done) break; - handleData(Buffer.from(value)); - } - })(), - (async () => { - while (true) { - const { done, value } = await stderrReader.read(); - if (done) break; - handleData(Buffer.from(value)); - } - })(), - ]); - - const exitCode = await child.exited; - - // Clean up - if (timeoutHandle) clearTimeout(timeoutHandle); - if (options?.signal) { - options.signal.removeEventListener("abort", abortHandler); - } - if (tempFileStream) { - tempFileStream.end(); - } - - // Combine buffered chunks for truncation (already sanitized) - const fullOutput = outputChunks.join(""); - const truncationResult = truncateTail(fullOutput); - - // Handle timeout - if (timedOut) { - const timeoutSecs = Math.round((options?.timeout || 0) / 1000); - resolve({ - output: `${fullOutput}\n\nCommand timed out after ${timeoutSecs} seconds`, - exitCode: undefined, - cancelled: true, - truncated: truncationResult.truncated, - fullOutputPath: tempFilePath, - }); - return; - } - - // Non-zero exit codes or signal-killed processes are considered cancelled if killed via signal - const cancelled = exitCode === null || (exitCode !== 0 && (options?.signal?.aborted ?? false)); - - resolve({ - output: truncationResult.truncated ? truncationResult.content : fullOutput, - exitCode: cancelled ? undefined : exitCode, - cancelled, - truncated: truncationResult.truncated, - fullOutputPath: tempFilePath, - }); - } catch (err) { - // Clean up - if (timeoutHandle) clearTimeout(timeoutHandle); - if (options?.signal) { - options.signal.removeEventListener("abort", abortHandler); - } - if (tempFileStream) { - tempFileStream.end(); - } - - reject(err); - } - })(); + const child: Subprocess = Bun.spawn([shell, ...args, finalCommand], { + cwd: options?.cwd, + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + env, }); + + signal.catch(() => { + killProcessTree(child.pid); + }); + + const sink = createOutputSink(DEFAULT_MAX_BYTES, DEFAULT_MAX_BYTES * 2, options?.onChunk); + + const writer = sink.getWriter(); + try { + async function pumpStream(readable: ReadableStream) { + const reader = readable.pipeThrough(createSanitizer()).getReader(); + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + await writer.write(value); + } + } finally { + reader.releaseLock(); + } + } + await Promise.all([ + pumpStream(child.stdout as ReadableStream), + pumpStream(child.stderr as ReadableStream), + ]); + } finally { + await writer.close(); + } + + // Non-zero exit codes or signal-killed processes are considered cancelled if killed via signal + const exitCode = await child.exited; + + const cancelled = exitCode === null || (exitCode !== 0 && (options?.signal?.aborted ?? false)); + + if (signal.timedOut()) { + const secs = Math.round(options!.timeout! / 1000); + return { + exitCode: undefined, + cancelled: true, + ...sink.dump(`Command timed out after ${secs} seconds`), + }; + } + + return { + exitCode: cancelled ? undefined : exitCode, + cancelled, + ...sink.dump(), + }; } diff --git a/packages/coding-agent/src/core/tools/find.ts b/packages/coding-agent/src/core/tools/find.ts index 765934f69..e293eaea4 100644 --- a/packages/coding-agent/src/core/tools/find.ts +++ b/packages/coding-agent/src/core/tools/find.ts @@ -4,6 +4,7 @@ import type { AgentTool } from "@oh-my-pi/pi-agent-core"; import { Type } from "@sinclair/typebox"; import { globSync } from "glob"; import { ensureTool } from "../../utils/tools-manager"; +import { untilAborted } from "../utils"; import { resolveToCwd } from "./path-utils"; import { DEFAULT_MAX_BYTES, formatSize, type TruncationResult, truncateHead } from "./truncate"; @@ -67,191 +68,171 @@ export function createFindTool(cwd: string): AgentTool { }, signal?: AbortSignal, ) => { - return new Promise((resolve, reject) => { - if (signal?.aborted) { - reject(new Error("Operation aborted")); - return; + return untilAborted(signal, async () => { + // Ensure fd is available + const fdPath = await ensureTool("fd", true); + if (!fdPath) { + throw new Error("fd is not available and could not be downloaded"); } - const onAbort = () => reject(new Error("Operation aborted")); - signal?.addEventListener("abort", onAbort, { once: true }); + const searchPath = resolveToCwd(searchDir || ".", cwd); + const effectiveLimit = limit ?? DEFAULT_LIMIT; + const effectiveType = type ?? "all"; + const includeHidden = hidden ?? false; + const shouldSortByMtime = sortByMtime ?? false; - (async () => { - try { - // Ensure fd is available - const fdPath = await ensureTool("fd", true); - if (!fdPath) { - reject(new Error("fd is not available and could not be downloaded")); - return; - } + // Build fd arguments + const args: string[] = [ + "--glob", // Use glob pattern + "--color=never", // No ANSI colors + "--max-results", + String(effectiveLimit), + ]; - const searchPath = resolveToCwd(searchDir || ".", cwd); - const effectiveLimit = limit ?? DEFAULT_LIMIT; - const effectiveType = type ?? "all"; - const includeHidden = hidden ?? false; - const shouldSortByMtime = sortByMtime ?? false; + if (includeHidden) { + args.push("--hidden"); + } - // Build fd arguments - const args: string[] = [ - "--glob", // Use glob pattern - "--color=never", // No ANSI colors - "--max-results", - String(effectiveLimit), - ]; + // Add type filter + if (effectiveType === "file") { + args.push("--type", "f"); + } else if (effectiveType === "dir") { + args.push("--type", "d"); + } - if (includeHidden) { - args.push("--hidden"); - } + // Include .gitignore files (root + nested) so fd respects them even outside git repos + const gitignoreFiles = new Set(); + const rootGitignore = path.join(searchPath, ".gitignore"); + if (existsSync(rootGitignore)) { + gitignoreFiles.add(rootGitignore); + } - // Add type filter - if (effectiveType === "file") { - args.push("--type", "f"); - } else if (effectiveType === "dir") { - args.push("--type", "d"); - } - - // Include .gitignore files (root + nested) so fd respects them even outside git repos - const gitignoreFiles = new Set(); - const rootGitignore = path.join(searchPath, ".gitignore"); - if (existsSync(rootGitignore)) { - gitignoreFiles.add(rootGitignore); - } - - try { - const nestedGitignores = globSync("**/.gitignore", { - cwd: searchPath, - dot: true, - absolute: true, - ignore: ["**/node_modules/**", "**/.git/**"], - }); - for (const file of nestedGitignores) { - gitignoreFiles.add(file); - } - } catch { - // Ignore glob errors - } - - for (const gitignorePath of gitignoreFiles) { - args.push("--ignore-file", gitignorePath); - } - - // Pattern and path - args.push(pattern, searchPath); - - // Run fd - const result = Bun.spawnSync([fdPath, ...args], { - stdin: "ignore", - stdout: "pipe", - stderr: "pipe", - }); - - signal?.removeEventListener("abort", onAbort); - - const output = result.stdout.toString().trim(); - - if (result.exitCode !== 0) { - const errorMsg = result.stderr.toString().trim() || `fd exited with code ${result.exitCode}`; - // fd returns non-zero for some errors but may still have partial output - if (!output) { - reject(new Error(errorMsg)); - return; - } - } - - if (!output) { - resolve({ - content: [{ type: "text", text: "No files found matching pattern" }], - details: { fileCount: 0, files: [], truncated: false }, - }); - return; - } - - const lines = output.split("\n"); - const relativized: string[] = []; - const mtimes: number[] = []; - - for (const rawLine of lines) { - const line = rawLine.replace(/\r$/, "").trim(); - if (!line) { - continue; - } - - const hadTrailingSlash = line.endsWith("/") || line.endsWith("\\"); - let relativePath = line; - if (line.startsWith(searchPath)) { - relativePath = line.slice(searchPath.length + 1); // +1 for the / - } else { - relativePath = path.relative(searchPath, line); - } - - if (hadTrailingSlash && !relativePath.endsWith("/")) { - relativePath += "/"; - } - - relativized.push(relativePath); - - // Collect mtime if sorting is requested - if (shouldSortByMtime) { - try { - const fullPath = path.join(searchPath, relativePath); - const stat: Stats = statSync(fullPath); - mtimes.push(stat.mtimeMs); - } catch { - mtimes.push(0); - } - } - } - - // Sort by mtime if requested (most recent first) - if (shouldSortByMtime && relativized.length > 0) { - const indexed = relativized.map((path, idx) => ({ path, mtime: mtimes[idx] || 0 })); - indexed.sort((a, b) => b.mtime - a.mtime); - relativized.length = 0; - relativized.push(...indexed.map((item) => item.path)); - } - - // Check if we hit the result limit - const resultLimitReached = relativized.length >= effectiveLimit; - - // Apply byte truncation (no line limit since we already have result limit) - const rawOutput = relativized.join("\n"); - const truncation = truncateHead(rawOutput, { maxLines: Number.MAX_SAFE_INTEGER }); - - let resultOutput = truncation.content; - const details: FindToolDetails = { - fileCount: relativized.length, - files: relativized.slice(0, 50), - truncated: resultLimitReached || truncation.truncated, - }; - - // Build notices - const notices: string[] = []; - - if (resultLimitReached) { - notices.push( - `${effectiveLimit} results limit reached. Use limit=${effectiveLimit * 2} for more, or refine pattern`, - ); - details.resultLimitReached = effectiveLimit; - } - - if (truncation.truncated) { - notices.push(`${formatSize(DEFAULT_MAX_BYTES)} limit reached`); - details.truncation = truncation; - } - - if (notices.length > 0) { - resultOutput += `\n\n[${notices.join(". ")}]`; - } - - resolve({ - content: [{ type: "text", text: resultOutput }], - details: Object.keys(details).length > 0 ? details : undefined, - }); - } catch (e: any) { - signal?.removeEventListener("abort", onAbort); - reject(e); + try { + const nestedGitignores = globSync("**/.gitignore", { + cwd: searchPath, + dot: true, + absolute: true, + ignore: ["**/node_modules/**", "**/.git/**"], + }); + for (const file of nestedGitignores) { + gitignoreFiles.add(file); } - })(); + } catch { + // Ignore glob errors + } + + for (const gitignorePath of gitignoreFiles) { + args.push("--ignore-file", gitignorePath); + } + + // Pattern and path + args.push(pattern, searchPath); + + // Run fd + const result = Bun.spawnSync([fdPath, ...args], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + + const output = result.stdout.toString().trim(); + + if (result.exitCode !== 0) { + const errorMsg = result.stderr.toString().trim() || `fd exited with code ${result.exitCode}`; + // fd returns non-zero for some errors but may still have partial output + if (!output) { + throw new Error(errorMsg); + } + } + + if (!output) { + return { + content: [{ type: "text", text: "No files found matching pattern" }], + details: { fileCount: 0, files: [], truncated: false }, + }; + } + + const lines = output.split("\n"); + const relativized: string[] = []; + const mtimes: number[] = []; + + for (const rawLine of lines) { + const line = rawLine.replace(/\r$/, "").trim(); + if (!line) { + continue; + } + + const hadTrailingSlash = line.endsWith("/") || line.endsWith("\\"); + let relativePath = line; + if (line.startsWith(searchPath)) { + relativePath = line.slice(searchPath.length + 1); // +1 for the / + } else { + relativePath = path.relative(searchPath, line); + } + + if (hadTrailingSlash && !relativePath.endsWith("/")) { + relativePath += "/"; + } + + relativized.push(relativePath); + + // Collect mtime if sorting is requested + if (shouldSortByMtime) { + try { + const fullPath = path.join(searchPath, relativePath); + const stat: Stats = statSync(fullPath); + mtimes.push(stat.mtimeMs); + } catch { + mtimes.push(0); + } + } + } + + // Sort by mtime if requested (most recent first) + if (shouldSortByMtime && relativized.length > 0) { + const indexed = relativized.map((path, idx) => ({ path, mtime: mtimes[idx] || 0 })); + indexed.sort((a, b) => b.mtime - a.mtime); + relativized.length = 0; + relativized.push(...indexed.map((item) => item.path)); + } + + // Check if we hit the result limit + const resultLimitReached = relativized.length >= effectiveLimit; + + // Apply byte truncation (no line limit since we already have result limit) + const rawOutput = relativized.join("\n"); + const truncation = truncateHead(rawOutput, { maxLines: Number.MAX_SAFE_INTEGER }); + + let resultOutput = truncation.content; + const details: FindToolDetails = { + fileCount: relativized.length, + files: relativized.slice(0, 50), + truncated: resultLimitReached || truncation.truncated, + }; + + // Build notices + const notices: string[] = []; + + if (resultLimitReached) { + notices.push( + `${effectiveLimit} results limit reached. Use limit=${effectiveLimit * 2} for more, or refine pattern`, + ); + details.resultLimitReached = effectiveLimit; + } + + if (truncation.truncated) { + notices.push(`${formatSize(DEFAULT_MAX_BYTES)} limit reached`); + details.truncation = truncation; + } + + if (notices.length > 0) { + resultOutput += `\n\n[${notices.join(". ")}]`; + } + + return { + content: [{ type: "text", text: resultOutput }], + details: Object.keys(details).length > 0 ? details : undefined, + }; }); }, }; diff --git a/packages/coding-agent/src/core/tools/ls.ts b/packages/coding-agent/src/core/tools/ls.ts index f37217ebd..57a7df9d2 100644 --- a/packages/coding-agent/src/core/tools/ls.ts +++ b/packages/coding-agent/src/core/tools/ls.ts @@ -2,6 +2,7 @@ import { existsSync, readdirSync, statSync } from "node:fs"; import nodePath from "node:path"; import type { AgentTool } from "@oh-my-pi/pi-agent-core"; import { Type } from "@sinclair/typebox"; +import { untilAborted } from "../utils"; import { resolveToCwd } from "./path-utils"; import { DEFAULT_MAX_BYTES, formatSize, type TruncationResult, truncateHead } from "./truncate"; @@ -28,109 +29,90 @@ export function createLsTool(cwd: string): AgentTool { { path, limit }: { path?: string; limit?: number }, signal?: AbortSignal, ) => { - return new Promise((resolve, reject) => { - if (signal?.aborted) { - reject(new Error("Operation aborted")); - return; + return untilAborted(signal, async () => { + const dirPath = resolveToCwd(path || ".", cwd); + const effectiveLimit = limit ?? DEFAULT_LIMIT; + + // Check if path exists + if (!existsSync(dirPath)) { + throw new Error(`Path not found: ${dirPath}`); } - const onAbort = () => reject(new Error("Operation aborted")); - signal?.addEventListener("abort", onAbort, { once: true }); + // Check if path is a directory + const stat = statSync(dirPath); + if (!stat.isDirectory()) { + throw new Error(`Not a directory: ${dirPath}`); + } + // Read directory entries + let entries: string[]; try { - const dirPath = resolveToCwd(path || ".", cwd); - const effectiveLimit = limit ?? DEFAULT_LIMIT; - - // Check if path exists - if (!existsSync(dirPath)) { - reject(new Error(`Path not found: ${dirPath}`)); - return; - } - - // Check if path is a directory - const stat = statSync(dirPath); - if (!stat.isDirectory()) { - reject(new Error(`Not a directory: ${dirPath}`)); - return; - } - - // Read directory entries - let entries: string[]; - try { - entries = readdirSync(dirPath); - } catch (e: any) { - reject(new Error(`Cannot read directory: ${e.message}`)); - return; - } - - // Sort alphabetically (case-insensitive) - entries.sort((a, b) => a.toLowerCase().localeCompare(b.toLowerCase())); - - // Format entries with directory indicators - const results: string[] = []; - let entryLimitReached = false; - - for (const entry of entries) { - if (results.length >= effectiveLimit) { - entryLimitReached = true; - break; - } - - const fullPath = nodePath.join(dirPath, entry); - let suffix = ""; - - try { - const entryStat = statSync(fullPath); - if (entryStat.isDirectory()) { - suffix = "/"; - } - } catch { - // Skip entries we can't stat - continue; - } - - results.push(entry + suffix); - } - - signal?.removeEventListener("abort", onAbort); - - if (results.length === 0) { - resolve({ content: [{ type: "text", text: "(empty directory)" }], details: undefined }); - return; - } - - // Apply byte truncation (no line limit since we already have entry limit) - const rawOutput = results.join("\n"); - const truncation = truncateHead(rawOutput, { maxLines: Number.MAX_SAFE_INTEGER }); - - let output = truncation.content; - const details: LsToolDetails = {}; - - // Build notices - const notices: string[] = []; - - if (entryLimitReached) { - notices.push(`${effectiveLimit} entries limit reached. Use limit=${effectiveLimit * 2} for more`); - details.entryLimitReached = effectiveLimit; - } - - if (truncation.truncated) { - notices.push(`${formatSize(DEFAULT_MAX_BYTES)} limit reached`); - details.truncation = truncation; - } - - if (notices.length > 0) { - output += `\n\n[${notices.join(". ")}]`; - } - - resolve({ - content: [{ type: "text", text: output }], - details: Object.keys(details).length > 0 ? details : undefined, - }); + entries = readdirSync(dirPath); } catch (e: any) { - signal?.removeEventListener("abort", onAbort); - reject(e); + throw new Error(`Cannot read directory: ${e.message}`); } + + // Sort alphabetically (case-insensitive) + entries.sort((a, b) => a.toLowerCase().localeCompare(b.toLowerCase())); + + // Format entries with directory indicators + const results: string[] = []; + let entryLimitReached = false; + + for (const entry of entries) { + if (results.length >= effectiveLimit) { + entryLimitReached = true; + break; + } + + const fullPath = nodePath.join(dirPath, entry); + let suffix = ""; + + try { + const entryStat = statSync(fullPath); + if (entryStat.isDirectory()) { + suffix = "/"; + } + } catch { + // Skip entries we can't stat + continue; + } + + results.push(entry + suffix); + } + + if (results.length === 0) { + return { content: [{ type: "text", text: "(empty directory)" }], details: undefined }; + } + + // Apply byte truncation (no line limit since we already have entry limit) + const rawOutput = results.join("\n"); + const truncation = truncateHead(rawOutput, { maxLines: Number.MAX_SAFE_INTEGER }); + + let output = truncation.content; + const details: LsToolDetails = {}; + + // Build notices + const notices: string[] = []; + + if (entryLimitReached) { + notices.push(`${effectiveLimit} entries limit reached. Use limit=${effectiveLimit * 2} for more`); + details.entryLimitReached = effectiveLimit; + } + + if (truncation.truncated) { + notices.push(`${formatSize(DEFAULT_MAX_BYTES)} limit reached`); + details.truncation = truncation; + } + + if (notices.length > 0) { + output += `\n\n[${notices.join(". ")}]`; + } + + return { + content: [{ type: "text", text: output }], + details: Object.keys(details).length > 0 ? details : undefined, + }; }); }, }; diff --git a/packages/coding-agent/src/core/tools/lsp/index.ts b/packages/coding-agent/src/core/tools/lsp/index.ts index 3a89b99c6..365058fe9 100644 --- a/packages/coding-agent/src/core/tools/lsp/index.ts +++ b/packages/coding-agent/src/core/tools/lsp/index.ts @@ -2,8 +2,10 @@ import * as fs from "node:fs"; import path from "node:path"; import type { AgentTool } from "@oh-my-pi/pi-agent-core"; import type { BunFile } from "bun"; +import { utils } from "packages/coding-agent/src/core"; import type { Theme } from "../../../modes/interactive/theme/theme"; import { logger } from "../../logger"; +import { untilAborted } from "../../utils"; import { resolveToCwd } from "../path-utils"; import { ensureFileOpen, @@ -53,8 +55,6 @@ import { symbolKindToIcon, uriToFile, } from "./utils"; -import { utils } from "packages/coding-agent/src/core"; -import { untilAborted } from "../../utils"; export type { LspServerStatus } from "./client"; export type { LspToolDetails } from "./types"; diff --git a/packages/coding-agent/src/core/tools/notebook.ts b/packages/coding-agent/src/core/tools/notebook.ts index a668c8a60..fe89902e4 100644 --- a/packages/coding-agent/src/core/tools/notebook.ts +++ b/packages/coding-agent/src/core/tools/notebook.ts @@ -1,5 +1,6 @@ import type { AgentTool } from "@oh-my-pi/pi-agent-core"; import { Type } from "@sinclair/typebox"; +import { untilAborted } from "../utils"; import { resolveToCwd } from "./path-utils"; const notebookSchema = Type.Object({ @@ -66,160 +67,104 @@ export function createNotebookTool(cwd: string): AgentTool { const absolutePath = resolveToCwd(notebook_path, cwd); - return new Promise<{ - content: Array<{ type: "text"; text: string }>; - details: NotebookToolDetails | undefined; - }>((resolve, reject) => { - if (signal?.aborted) { - reject(new Error("Operation aborted")); - return; + return untilAborted(signal, async () => { + // Check if file exists + const file = Bun.file(absolutePath); + if (!(await file.exists())) { + throw new Error(`Notebook not found: ${notebook_path}`); } - let aborted = false; - - const onAbort = () => { - aborted = true; - reject(new Error("Operation aborted")); - }; - - if (signal) { - signal.addEventListener("abort", onAbort, { once: true }); + // Read and parse notebook + let notebook: Notebook; + try { + notebook = await file.json(); + } catch { + throw new Error(`Invalid JSON in notebook: ${notebook_path}`); } - (async () => { - try { - // Check if file exists - const file = Bun.file(absolutePath); - if (!(await file.exists())) { - if (signal) signal.removeEventListener("abort", onAbort); - reject(new Error(`Notebook not found: ${notebook_path}`)); - return; - } + // Validate notebook structure + if (!notebook.cells || !Array.isArray(notebook.cells)) { + throw new Error(`Invalid notebook structure (missing cells array): ${notebook_path}`); + } - if (aborted) return; + const cellCount = notebook.cells.length; - // Read and parse notebook - let notebook: Notebook; - try { - notebook = await file.json(); - } catch { - if (signal) signal.removeEventListener("abort", onAbort); - reject(new Error(`Invalid JSON in notebook: ${notebook_path}`)); - return; - } - - if (aborted) return; - - // Validate notebook structure - if (!notebook.cells || !Array.isArray(notebook.cells)) { - if (signal) signal.removeEventListener("abort", onAbort); - reject(new Error(`Invalid notebook structure (missing cells array): ${notebook_path}`)); - return; - } - - const cellCount = notebook.cells.length; - - // Validate cell_index based on action - if (action === "insert") { - if (cell_index < 0 || cell_index > cellCount) { - if (signal) signal.removeEventListener("abort", onAbort); - reject( - new Error( - `Cell index ${cell_index} out of range for insert (0-${cellCount}) in ${notebook_path}`, - ), - ); - return; - } - } else { - if (cell_index < 0 || cell_index >= cellCount) { - if (signal) signal.removeEventListener("abort", onAbort); - reject( - new Error(`Cell index ${cell_index} out of range (0-${cellCount - 1}) in ${notebook_path}`), - ); - return; - } - } - - // Validate content for edit/insert - if ((action === "edit" || action === "insert") && content === undefined) { - if (signal) signal.removeEventListener("abort", onAbort); - reject(new Error(`Content is required for ${action} action`)); - return; - } - - if (aborted) return; - - // Perform the action - let resultMessage: string; - let finalCellType: string | undefined; - - switch (action) { - case "edit": { - const sourceLines = splitIntoLines(content!); - notebook.cells[cell_index].source = sourceLines; - finalCellType = notebook.cells[cell_index].cell_type; - resultMessage = `Replaced cell ${cell_index} (${finalCellType})`; - break; - } - case "insert": { - const sourceLines = splitIntoLines(content!); - const newCellType = (cell_type as "code" | "markdown") || "code"; - const newCell: NotebookCell = { - cell_type: newCellType, - source: sourceLines, - metadata: {}, - }; - if (newCellType === "code") { - newCell.execution_count = null; - newCell.outputs = []; - } - notebook.cells.splice(cell_index, 0, newCell); - finalCellType = newCellType; - resultMessage = `Inserted ${newCellType} cell at position ${cell_index}`; - break; - } - case "delete": { - finalCellType = notebook.cells[cell_index].cell_type; - notebook.cells.splice(cell_index, 1); - resultMessage = `Deleted cell ${cell_index} (${finalCellType})`; - break; - } - default: { - if (signal) signal.removeEventListener("abort", onAbort); - reject(new Error(`Invalid action: ${action}`)); - return; - } - } - - if (aborted) return; - - // Write back with single-space indentation - await Bun.write(absolutePath, JSON.stringify(notebook, null, 1)); - - if (aborted) return; - - if (signal) signal.removeEventListener("abort", onAbort); - - const newCellCount = notebook.cells.length; - resolve({ - content: [ - { - type: "text", - text: `${resultMessage}. Notebook now has ${newCellCount} cells.`, - }, - ], - details: { - action: action as "edit" | "insert" | "delete", - cellIndex: cell_index, - cellType: finalCellType, - totalCells: newCellCount, - }, - }); - } catch (error: any) { - if (signal) signal.removeEventListener("abort", onAbort); - if (!aborted) reject(error); + // Validate cell_index based on action + if (action === "insert") { + if (cell_index < 0 || cell_index > cellCount) { + throw new Error( + `Cell index ${cell_index} out of range for insert (0-${cellCount}) in ${notebook_path}`, + ); } - })(); + } else { + if (cell_index < 0 || cell_index >= cellCount) { + throw new Error(`Cell index ${cell_index} out of range (0-${cellCount - 1}) in ${notebook_path}`); + } + } + + // Validate content for edit/insert + if ((action === "edit" || action === "insert") && content === undefined) { + throw new Error(`Content is required for ${action} action`); + } + + // Perform the action + let resultMessage: string; + let finalCellType: string | undefined; + + switch (action) { + case "edit": { + const sourceLines = splitIntoLines(content!); + notebook.cells[cell_index].source = sourceLines; + finalCellType = notebook.cells[cell_index].cell_type; + resultMessage = `Replaced cell ${cell_index} (${finalCellType})`; + break; + } + case "insert": { + const sourceLines = splitIntoLines(content!); + const newCellType = (cell_type as "code" | "markdown") || "code"; + const newCell: NotebookCell = { + cell_type: newCellType, + source: sourceLines, + metadata: {}, + }; + if (newCellType === "code") { + newCell.execution_count = null; + newCell.outputs = []; + } + notebook.cells.splice(cell_index, 0, newCell); + finalCellType = newCellType; + resultMessage = `Inserted ${newCellType} cell at position ${cell_index}`; + break; + } + case "delete": { + finalCellType = notebook.cells[cell_index].cell_type; + notebook.cells.splice(cell_index, 1); + resultMessage = `Deleted cell ${cell_index} (${finalCellType})`; + break; + } + default: { + throw new Error(`Invalid action: ${action}`); + } + } + + // Write back with single-space indentation + await Bun.write(absolutePath, JSON.stringify(notebook, null, 1)); + + const newCellCount = notebook.cells.length; + return { + content: [ + { + type: "text", + text: `${resultMessage}. Notebook now has ${newCellCount} cells.`, + }, + ], + details: { + action: action as "edit" | "insert" | "delete", + cellIndex: cell_index, + cellType: finalCellType, + totalCells: newCellCount, + }, + }; }); }, }; diff --git a/packages/coding-agent/src/core/tools/read.ts b/packages/coding-agent/src/core/tools/read.ts index d4ebd3a60..c29d84c81 100644 --- a/packages/coding-agent/src/core/tools/read.ts +++ b/packages/coding-agent/src/core/tools/read.ts @@ -6,6 +6,7 @@ import type { AgentTool } from "@oh-my-pi/pi-agent-core"; import type { ImageContent, TextContent } from "@oh-my-pi/pi-ai"; import { Type } from "@sinclair/typebox"; import { detectSupportedImageMimeTypeFromFile } from "../../utils/mime"; +import { untilAborted } from "../utils"; import { resolveReadPath } from "./path-utils"; import { DEFAULT_MAX_BYTES, DEFAULT_MAX_LINES, formatSize, type TruncationResult, truncateHead } from "./truncate"; @@ -69,169 +70,120 @@ Usage: ) => { const absolutePath = resolveReadPath(path, cwd); - return new Promise<{ content: (TextContent | ImageContent)[]; details: ReadToolDetails | undefined }>( - (resolve, reject) => { - // Check if already aborted - if (signal?.aborted) { - reject(new Error("Operation aborted")); - return; - } + return untilAborted(signal, async () => { + // Check if file exists + await access(absolutePath, constants.R_OK); - let aborted = false; + const mimeType = await detectSupportedImageMimeTypeFromFile(absolutePath); + const ext = extname(absolutePath).toLowerCase(); - // Set up abort handler - const onAbort = () => { - aborted = true; - reject(new Error("Operation aborted")); - }; + // Read the file based on type + let content: (TextContent | ImageContent)[]; + let details: ReadToolDetails | undefined; - if (signal) { - signal.addEventListener("abort", onAbort, { once: true }); - } + if (mimeType) { + // Read as image (binary) + const buffer = await readFile(absolutePath); + const base64 = buffer.toString("base64"); - // Perform the read operation - (async () => { - try { - // Check if file exists - await access(absolutePath, constants.R_OK); + content = [ + { type: "text", text: `Read image file [${mimeType}]` }, + { type: "image", data: base64, mimeType }, + ]; + } else if (CONVERTIBLE_EXTENSIONS.has(ext)) { + // Convert document via markitdown + const result = convertWithMarkitdown(absolutePath); + if (result.ok) { + // Apply truncation to converted content + const truncation = truncateHead(result.content); + let outputText = truncation.content; - // Check if aborted before reading - if (aborted) { - return; - } - - const mimeType = await detectSupportedImageMimeTypeFromFile(absolutePath); - const ext = extname(absolutePath).toLowerCase(); - - // Read the file based on type - let content: (TextContent | ImageContent)[]; - let details: ReadToolDetails | undefined; - - if (mimeType) { - // Read as image (binary) - const buffer = await readFile(absolutePath); - const base64 = buffer.toString("base64"); - - content = [ - { type: "text", text: `Read image file [${mimeType}]` }, - { type: "image", data: base64, mimeType }, - ]; - } else if (CONVERTIBLE_EXTENSIONS.has(ext)) { - // Convert document via markitdown - const result = convertWithMarkitdown(absolutePath); - if (result.ok) { - // Apply truncation to converted content - const truncation = truncateHead(result.content); - let outputText = truncation.content; - - if (truncation.truncated) { - outputText += `\n\n[Document converted via markitdown. Output truncated to $formatSize( - DEFAULT_MAX_BYTES, - )]`; - details = { truncation }; - } - - content = [{ type: "text", text: outputText }]; - } else { - // markitdown not available or failed - const errorMsg = - result.error === "markitdown not found" - ? `markitdown not installed. Install with: pip install markitdown` - : result.error || "conversion failed"; - content = [{ type: "text", text: `[Cannot read ${ext} file: ${errorMsg}]` }]; - } - } else { - // Read as text - const textContent = await readFile(absolutePath, "utf-8"); - const allLines = textContent.split("\n"); - const totalFileLines = allLines.length; - - // Apply offset if specified (1-indexed to 0-indexed) - const startLine = offset ? Math.max(0, offset - 1) : 0; - const startLineDisplay = startLine + 1; // For display (1-indexed) - - // Check if offset is out of bounds - if (startLine >= allLines.length) { - throw new Error(`Offset ${offset} is beyond end of file (${allLines.length} lines total)`); - } - - // If limit is specified by user, use it; otherwise we'll let truncateHead decide - let selectedContent: string; - let userLimitedLines: number | undefined; - if (limit !== undefined) { - const endLine = Math.min(startLine + limit, allLines.length); - selectedContent = allLines.slice(startLine, endLine).join("\n"); - userLimitedLines = endLine - startLine; - } else { - selectedContent = allLines.slice(startLine).join("\n"); - } - - // Apply truncation (respects both line and byte limits) - const truncation = truncateHead(selectedContent); - - let outputText: string; - - if (truncation.firstLineExceedsLimit) { - // First line at offset exceeds 30KB - tell model to use bash - const firstLineSize = formatSize(Buffer.byteLength(allLines[startLine], "utf-8")); - outputText = `[Line ${startLineDisplay} is ${firstLineSize}, exceeds ${formatSize( - DEFAULT_MAX_BYTES, - )} limit. Use bash: sed -n '${startLineDisplay}p' ${path} | head -c ${DEFAULT_MAX_BYTES}]`; - details = { truncation }; - } else if (truncation.truncated) { - // Truncation occurred - build actionable notice - const endLineDisplay = startLineDisplay + truncation.outputLines - 1; - const nextOffset = endLineDisplay + 1; - - outputText = truncation.content; - - if (truncation.truncatedBy === "lines") { - outputText += `\n\n[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines}. Use offset=${nextOffset} to continue]`; - } else { - outputText += `\n\n[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines} (${formatSize( - DEFAULT_MAX_BYTES, - )} limit). Use offset=${nextOffset} to continue]`; - } - details = { truncation }; - } else if (userLimitedLines !== undefined && startLine + userLimitedLines < allLines.length) { - // User specified limit, there's more content, but no truncation - const remaining = allLines.length - (startLine + userLimitedLines); - const nextOffset = startLine + userLimitedLines + 1; - - outputText = truncation.content; - outputText += `\n\n[${remaining} more lines in file. Use offset=${nextOffset} to continue]`; - } else { - // No truncation, no user limit exceeded - outputText = truncation.content; - } - - content = [{ type: "text", text: outputText }]; - } - - // Check if aborted after reading - if (aborted) { - return; - } - - // Clean up abort handler - if (signal) { - signal.removeEventListener("abort", onAbort); - } - - resolve({ content, details }); - } catch (error: any) { - // Clean up abort handler - if (signal) { - signal.removeEventListener("abort", onAbort); - } - - if (!aborted) { - reject(error); - } + if (truncation.truncated) { + outputText += `\n\n[Document converted via markitdown. Output truncated to $formatSize( + DEFAULT_MAX_BYTES, + )]`; + details = { truncation }; } - })(); - }, - ); + + content = [{ type: "text", text: outputText }]; + } else { + // markitdown not available or failed + const errorMsg = + result.error === "markitdown not found" + ? `markitdown not installed. Install with: pip install markitdown` + : result.error || "conversion failed"; + content = [{ type: "text", text: `[Cannot read ${ext} file: ${errorMsg}]` }]; + } + } else { + // Read as text + const textContent = await readFile(absolutePath, "utf-8"); + const allLines = textContent.split("\n"); + const totalFileLines = allLines.length; + + // Apply offset if specified (1-indexed to 0-indexed) + const startLine = offset ? Math.max(0, offset - 1) : 0; + const startLineDisplay = startLine + 1; // For display (1-indexed) + + // Check if offset is out of bounds + if (startLine >= allLines.length) { + throw new Error(`Offset ${offset} is beyond end of file (${allLines.length} lines total)`); + } + + // If limit is specified by user, use it; otherwise we'll let truncateHead decide + let selectedContent: string; + let userLimitedLines: number | undefined; + if (limit !== undefined) { + const endLine = Math.min(startLine + limit, allLines.length); + selectedContent = allLines.slice(startLine, endLine).join("\n"); + userLimitedLines = endLine - startLine; + } else { + selectedContent = allLines.slice(startLine).join("\n"); + } + + // Apply truncation (respects both line and byte limits) + const truncation = truncateHead(selectedContent); + + let outputText: string; + + if (truncation.firstLineExceedsLimit) { + // First line at offset exceeds 30KB - tell model to use bash + const firstLineSize = formatSize(Buffer.byteLength(allLines[startLine], "utf-8")); + outputText = `[Line ${startLineDisplay} is ${firstLineSize}, exceeds ${formatSize( + DEFAULT_MAX_BYTES, + )} limit. Use bash: sed -n '${startLineDisplay}p' ${path} | head -c ${DEFAULT_MAX_BYTES}]`; + details = { truncation }; + } else if (truncation.truncated) { + // Truncation occurred - build actionable notice + const endLineDisplay = startLineDisplay + truncation.outputLines - 1; + const nextOffset = endLineDisplay + 1; + + outputText = truncation.content; + + if (truncation.truncatedBy === "lines") { + outputText += `\n\n[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines}. Use offset=${nextOffset} to continue]`; + } else { + outputText += `\n\n[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines} (${formatSize( + DEFAULT_MAX_BYTES, + )} limit). Use offset=${nextOffset} to continue]`; + } + details = { truncation }; + } else if (userLimitedLines !== undefined && startLine + userLimitedLines < allLines.length) { + // User specified limit, there's more content, but no truncation + const remaining = allLines.length - (startLine + userLimitedLines); + const nextOffset = startLine + userLimitedLines + 1; + + outputText = truncation.content; + outputText += `\n\n[${remaining} more lines in file. Use offset=${nextOffset} to continue]`; + } else { + // No truncation, no user limit exceeded + outputText = truncation.content; + } + + content = [{ type: "text", text: outputText }]; + } + + return { content, details }; + }); }, }; } diff --git a/packages/coding-agent/src/core/utils.ts b/packages/coding-agent/src/core/utils.ts index b2ae6ab13..f25d6d613 100644 --- a/packages/coding-agent/src/core/utils.ts +++ b/packages/coding-agent/src/core/utils.ts @@ -49,3 +49,139 @@ export function once(fn: () => T): () => T { return value; }; } + +// ScopeSignal is a cancellation/helper utility similar to AbortController but +// allows composition of an existing AbortSignal and/or a timeout. It exposes a +// simple API for cancellation observation (finally, catch). +interface ScopeSignalOptions { + signal?: AbortSignal; + timeout?: number; +} + +const kTimeoutReason = new Error("Timeout"); +const kDisposedReason = new Error("Disposed"); + +/** + * Type of signal exit (None = disposed, TimedOut = timed out, Aborted = underlying signal aborted) + */ +enum ExitReason { + None = 0, + TimedOut = 1, + Aborted = 2, +} + +/** + * ScopeSignal: composable cancellation for async work–observes an external AbortSignal and/or a timeout. + * + * Use .finally(fn) to register a one-time callback invoked on *any* exit (abort, timeout, or manual dispose). + * Use .catch(fn) to register a one-time callback invoked only on abort/timeout. + * + * Disposing ScopeSignal disables further callbacks. + */ +export class ScopeSignal implements Disposable { + #signal: AbortSignal | undefined; + #timer: NodeJS.Timeout | undefined; + #exit = undefined as ExitReason | undefined; + #onAbort: (() => void) | undefined; + #callbacks?: (() => void)[]; + #reason: unknown | undefined; + + /** + * Provides abort/timeout reason (Error or user-defined). + */ + get reason(): unknown | undefined { + return this.#reason; + } + + /** + * True if exited due to external AbortSignal or timeout. + */ + get aborted(): boolean { + return this.#exit !== undefined && this.#exit > ExitReason.None; + } + + /** + * True if this ScopeSignal timed out (not external abort). + */ + timedOut(): boolean { + return this.#exit === ExitReason.TimedOut; + } + + /** + * Create a new ScopeSignal, optionally observing an AbortSignal and/or auto-aborting after a timeout (ms). + */ + constructor(options?: ScopeSignalOptions) { + const { signal, timeout } = options ?? {}; + + if (signal?.aborted) { + this.#abort(ExitReason.Aborted, signal.reason); // Immediately abort if already-aborted + return; + } + if (timeout && timeout <= 0) { + this.#abort(ExitReason.TimedOut, kTimeoutReason); + return; + } + + // Observe external signal if provided + if (signal) { + const onAbort = () => { + this.#abort(ExitReason.Aborted, signal.reason); + }; + this.#signal = signal; + this.#onAbort = onAbort; + this.#signal.addEventListener("abort", onAbort, { once: true }); + } + + // Set up timeout if provided + if (timeout) { + this.#timer = setTimeout(() => { + this.#abort(ExitReason.TimedOut, kTimeoutReason); + }, timeout); + } + } + + /** + * Register a one-time callback invoked on any exit (abort, timeout, or manual dispose). + * Runs immediately if already exited. + */ + finally(onfinally: () => void): void { + if (this.#exit !== undefined) { + onfinally(); + return; + } + this.#callbacks ??= []; + this.#callbacks.push(onfinally); + } + + /** + * Register a one-time callback invoked only if exited due to abort/timeout (not normal disposal). + */ + catch(oncatch: (reason: unknown) => void): void { + this.finally(() => { + if (this.aborted) { + oncatch(this.reason); + } + }); + } + + /** Internal: cause exit; only first call takes effect. */ + #abort(exit: ExitReason, reason?: unknown): void { + if (this.#exit !== undefined) return; + this.#reason = reason; + clearTimeout(this.#timer); + this.#signal?.removeEventListener("abort", this.#onAbort!); + + this.#exit = exit; + + const callbacks = this.#callbacks; + this.#callbacks = undefined; + callbacks?.forEach((fn) => void fn()); + } + + /** + * Dispose: marks as normally exited (not abort/timeout); disables further callback registration. + */ + [Symbol.dispose](): void { + this.#abort(ExitReason.None, kDisposedReason); + } +}