fix(coding-agent/core): corrected async stream handling and timeout conversion

- Fixed async stream handling in executors by adding await to stream.dump() calls.
- Improved output sink to handle large outputs and prevent truncation.
- Fixed timeout handling in web fetch tool with proper millisecond conversion.
- Updated process execution in web scrapers to use ptree.cspawn instead of Bun executor.
This commit is contained in:
can1357
2026-01-21 06:40:25 +01:00
parent 08519c1cf5
commit 1db58085e3
19 changed files with 239 additions and 407 deletions
+11
View File
@@ -1,6 +1,17 @@
# Changelog
## [Unreleased]
### Changed
- Updated output sink to properly handle large outputs
- Improved error message formatting in SSH executor
- Updated web fetch timeout bounds and conversion
### Fixed
- Fixed output truncation handling in streaming output
- Fixed timeout handling in web fetch tool
- Fixed async stream dumping in executors
## [6.8.3] - 2026-01-21
+1 -26
View File
@@ -324,6 +324,7 @@ const myExtension: ExtensionFactory = async (pi) => {
```
Key extension events:
- `before_agent_start`: Receives `systemPrompt` and can return full replacement (not just append)
- `user_bash`: Intercept `!`/`!!` commands for custom execution (e.g., remote SSH)
- `session_shutdown`: Cleanup notification before exit
@@ -411,32 +412,6 @@ bun test --testNamePattern="RPC"
bun test/rpc-example.ts
```
### Pluggable Tool Operations
Built-in tools support pluggable operations for remote execution:
- **BashOperations**: Execute commands on remote systems
- **LsOperations**: Remote directory listing
- **GrepOperations**: Remote content search
- **FindOperations**: Remote file search
Example: SSH extension overriding bash execution:
```typescript
pi.on("user_bash", async (event) => {
if (shouldRunRemotely()) {
return {
operations: {
exec: async (cmd, cwd, opts) => {
// Execute via SSH
return { exitCode: 0 };
},
},
};
}
});
```
### Managed Binaries
Tools like `fd` and `rg` are auto-downloaded to `~/.omp/bin/` (migrated from `~/.omp/agent/tools/`).
@@ -22,7 +22,7 @@ import type { Rule } from "../capability/rule";
import { getAgentDbPath } from "../config";
import { theme } from "../modes/interactive/theme/theme";
import ttsrInterruptTemplate from "../prompts/system/ttsr-interrupt.md" with { type: "text" };
import { type BashResult, executeBash as executeBashCommand, executeBashWithOperations } from "./bash-executor";
import { type BashResult, executeBash as executeBashCommand } from "./bash-executor";
import {
type CompactionResult,
calculateContextTokens,
@@ -60,7 +60,6 @@ import type { Skill, SkillWarning } from "./skills";
import { expandSlashCommand, type FileSlashCommand } from "./slash-commands";
import { closeAllConnections } from "./ssh/connection-manager";
import { unmountAll } from "./ssh/sshfs-mount";
import type { BashOperations } from "./tools/bash";
import { normalizeDiff, normalizeToLF, ParseError, previewPatch, stripBom } from "./tools/patch";
import { resolveToCwd } from "./tools/path-utils";
import type { TodoItem } from "./tools/todo-write";
@@ -2546,25 +2545,19 @@ export class AgentSession {
* @param command The bash command to execute
* @param onChunk Optional streaming callback for output
* @param options.excludeFromContext If true, command output won't be sent to LLM (!! prefix)
* @param options.operations Custom BashOperations for remote execution
*/
async executeBash(
command: string,
onChunk?: (chunk: string) => void,
options?: { excludeFromContext?: boolean; operations?: BashOperations },
options?: { excludeFromContext?: boolean },
): Promise<BashResult> {
this._bashAbortController = new AbortController();
try {
const result = options?.operations
? await executeBashWithOperations(command, process.cwd(), options.operations, {
onChunk,
signal: this._bashAbortController.signal,
})
: await executeBashCommand(command, {
onChunk,
signal: this._bashAbortController.signal,
});
const result = await executeBashCommand(command, {
onChunk,
signal: this._bashAbortController.signal,
});
this.recordBashResult(command, result, options);
return result;
@@ -8,7 +8,6 @@ import { cspawn, Exception, ptree } from "@oh-my-pi/pi-utils";
import { getShellConfig } from "../utils/shell";
import { getOrCreateSnapshot, getSnapshotSourceCommand } from "../utils/shell-snapshot";
import { OutputSink } from "./streaming-output";
import type { BashOperations } from "./tools/bash";
export interface BashExecutorOptions {
cwd?: string;
@@ -34,7 +33,7 @@ export async function executeBash(command: string, options?: BashExecutorOptions
const prefixedCommand = prefix ? `${prefix} ${command}` : command;
const finalCommand = `${snapshotPrefix}${prefixedCommand}`;
const stream = new OutputSink({ onChunk: options?.onChunk });
const sink = new OutputSink({ onChunk: options?.onChunk });
const child = cspawn([shell, ...args, finalCommand], {
cwd: options?.cwd,
@@ -45,12 +44,9 @@ export async function executeBash(command: string, options?: BashExecutorOptions
// Pump streams - errors during abort/timeout are expected
// Use preventClose to avoid closing the shared sink when either stream finishes
await Promise.allSettled([
child.stdout.pipeTo(stream.createWritable()),
child.stderr.pipeTo(stream.createWritable()),
])
.then(() => stream.close())
.catch(() => {});
await Promise.allSettled([child.stdout.pipeTo(sink.createInput()), child.stderr.pipeTo(sink.createInput())]).catch(
() => {},
);
// Wait for process exit
try {
@@ -58,7 +54,7 @@ export async function executeBash(command: string, options?: BashExecutorOptions
return {
exitCode: child.exitCode ?? 0,
cancelled: false,
...stream.dump(),
...(await sink.dump()),
};
} catch (err) {
// Exception covers NonZeroExitError, AbortError, TimeoutError
@@ -71,7 +67,7 @@ export async function executeBash(command: string, options?: BashExecutorOptions
return {
exitCode: undefined,
cancelled: true,
...stream.dump(annotation),
...(await sink.dump(annotation)),
};
}
@@ -79,60 +75,7 @@ export async function executeBash(command: string, options?: BashExecutorOptions
return {
exitCode: err.exitCode,
cancelled: false,
...stream.dump(),
};
}
throw err;
}
}
export async function executeBashWithOperations(
command: string,
cwd: string,
operations: BashOperations,
options?: BashExecutorOptions,
): Promise<BashResult> {
const stream = new OutputSink({ onChunk: options?.onChunk });
const writable = stream.createWritable();
const writer = writable.getWriter();
const closeStreams = async () => {
try {
await writer.close();
} catch {}
try {
await writable.close();
} catch {}
try {
await stream.close();
} catch {}
};
try {
const result = await operations.exec(command, cwd, {
onData: (data) => writer.write(data),
signal: options?.signal,
timeout: options?.timeout,
});
await closeStreams();
const cancelled = options?.signal?.aborted ?? false;
return {
exitCode: cancelled ? undefined : (result.exitCode ?? undefined),
cancelled,
...stream.dump(),
};
} catch (err) {
await closeStreams();
if (options?.signal?.aborted) {
return {
exitCode: undefined,
cancelled: true,
...stream.dump(),
...(await sink.dump()),
};
}
@@ -29,7 +29,6 @@ import type {
SessionManager,
} from "../session-manager";
import type { BashToolDetails, FindToolDetails, GrepToolDetails, LsToolDetails, ReadToolDetails } from "../tools";
import type { BashOperations } from "../tools/bash";
import type { EditToolDetails } from "../tools/patch";
export type { ExecOptions, ExecResult } from "../exec";
@@ -551,8 +550,6 @@ export interface InputEventResult {
/** Result from user_bash event handler */
export interface UserBashEventResult {
/** Custom operations to use for execution */
operations?: BashOperations;
/** Full replacement: extension handled execution, use this result */
result?: BashResult;
}
+1 -1
View File
@@ -11,7 +11,7 @@ export {
type PromptOptions,
type SessionStats,
} from "./agent-session";
export { type BashExecutorOptions, type BashResult, executeBash, executeBashWithOperations } from "./bash-executor";
export { type BashExecutorOptions, type BashResult, executeBash } from "./bash-executor";
export type { CompactionResult } from "./compaction/index";
export {
discoverAndLoadExtensions,
@@ -1,4 +1,4 @@
import { logger, sanitizeText } from "@oh-my-pi/pi-utils";
import { logger } from "@oh-my-pi/pi-utils";
import {
checkPythonKernelAvailability,
type KernelDisplayOutput,
@@ -210,30 +210,20 @@ async function executeWithKernel(
code: string,
options: PythonExecutorOptions | undefined,
): Promise<PythonResult> {
const sink = new OutputSink({ onLine: options?.onChunk });
const sink = new OutputSink({ onChunk: options?.onChunk });
const displayOutputs: KernelDisplayOutput[] = [];
try {
const writable = sink.createStringWritable();
const writer = writable.getWriter();
let result: KernelExecuteResult;
try {
result = await kernel.execute(code, {
signal: options?.signal,
timeoutMs: options?.timeout,
onChunk: (text) => {
writer.write(sanitizeText(text));
},
onDisplay: (output) => {
displayOutputs.push(output);
},
});
} catch (err) {
await writer.abort(err);
throw err;
} finally {
await writer.close().catch(() => {});
}
const result = await kernel.execute(code, {
signal: options?.signal,
timeoutMs: options?.timeout,
onChunk: (text) => {
sink.push(text);
},
onDisplay: (output) => {
displayOutputs.push(output);
},
});
if (result.cancelled) {
const secs = options?.timeout ? Math.round(options.timeout / 1000) : undefined;
@@ -244,7 +234,7 @@ async function executeWithKernel(
cancelled: true,
displayOutputs,
stdinRequested: result.stdinRequested,
...sink.dump(annotation),
...(await sink.dump(annotation)),
};
}
@@ -254,7 +244,7 @@ async function executeWithKernel(
cancelled: false,
displayOutputs,
stdinRequested: true,
...sink.dump("Kernel requested stdin; interactive input is not supported."),
...(await sink.dump("Kernel requested stdin; interactive input is not supported.")),
};
}
@@ -264,7 +254,7 @@ async function executeWithKernel(
cancelled: false,
displayOutputs,
stdinRequested: false,
...sink.dump(),
...(await sink.dump()),
};
} catch (err) {
const error = err instanceof Error ? err : new Error(String(err));
@@ -70,16 +70,11 @@ export async function executeSSH(
timeout: options?.timeout,
});
const sink = new OutputSink({ onLine: options?.onChunk });
const sink = new OutputSink({ onChunk: options?.onChunk });
try {
await Promise.allSettled([
child.stdout.pipeTo(sink.createWritable()),
child.stderr.pipeTo(sink.createWritable()),
]);
} finally {
await sink.close();
}
await Promise.allSettled([child.stdout.pipeTo(sink.createInput()), child.stderr.pipeTo(sink.createInput())]).catch(
() => {},
);
try {
await child.exited;
@@ -87,7 +82,7 @@ export async function executeSSH(
return {
exitCode,
cancelled: false,
...sink.dump(),
...(await sink.dump()),
};
} catch (err) {
if (err instanceof ptree.Exception) {
@@ -95,20 +90,20 @@ export async function executeSSH(
return {
exitCode: undefined,
cancelled: true,
...sink.dump(`SSH: ${err.message}`),
...(await sink.dump(`SSH: ${err.message}`)),
};
}
if (err.aborted) {
return {
exitCode: undefined,
cancelled: true,
...sink.dump(`SSH command aborted: ${err.message}`),
...(await sink.dump(`Command aborted: ${err.message}`)),
};
}
return {
exitCode: err.exitCode,
cancelled: false,
...sink.dump(`Unexpected error: ${err.message}`),
...(await sink.dump(`Unexpected error: ${err.message}`)),
};
}
throw err;
@@ -2,7 +2,7 @@ import { tmpdir } from "node:os";
import { join } from "node:path";
import { sanitizeText } from "@oh-my-pi/pi-utils";
import { nanoid } from "nanoid";
import { DEFAULT_MAX_BYTES, DEFAULT_MAX_COLUMN } from "./tools/truncate";
import { DEFAULT_MAX_BYTES } from "./tools/truncate";
export interface OutputResult {
output: string;
@@ -14,7 +14,6 @@ export interface OutputSinkOptions {
allocateFilePath?: () => string;
spillThreshold?: number;
maxColumn?: number;
onLine?: (line: string) => void;
onChunk?: (chunk: string) => void;
}
@@ -29,182 +28,93 @@ function defaultFilePathAllocator(): string {
* When memory limit exceeded, spills ~half to file in one batch operation.
*/
export class OutputSink {
private buffer = "";
private lineEnds: number[] = []; // String index after each \n
#buffer = "";
#file?: {
path: string;
sink: Bun.FileSink;
};
#bytesWritten: number = 0;
private fileSink?: Bun.FileSink;
private filePath?: string;
private readonly allocateFilePath: () => string;
private readonly spillThreshold: number;
private readonly maxColumn: number;
private readonly onLine?: (line: string) => void;
private readonly onChunk?: (chunk: string) => void;
readonly #allocateFilePath: () => string;
readonly #spillThreshold: number;
readonly #onChunk?: (chunk: string) => void;
constructor(options?: OutputSinkOptions) {
const {
allocateFilePath = defaultFilePathAllocator,
spillThreshold = DEFAULT_MAX_BYTES,
maxColumn = DEFAULT_MAX_COLUMN,
onLine,
onChunk,
} = options ?? {};
this.allocateFilePath = allocateFilePath;
this.spillThreshold = spillThreshold;
this.maxColumn = maxColumn;
this.onLine = onLine;
this.onChunk = onChunk;
this.#allocateFilePath = allocateFilePath;
this.#spillThreshold = spillThreshold;
this.#onChunk = onChunk;
}
private pushLine(line: string, term?: string): void {
while (line.length > this.maxColumn) {
this.pushLine(line.slice(0, this.maxColumn), "--\n");
line = line.slice(this.maxColumn);
}
async #pushSanitized(data: string): Promise<void> {
this.#onChunk?.(data);
const dataBytes = Buffer.byteLength(data);
const overflow = dataBytes + this.#bytesWritten > this.#spillThreshold || this.#file != null;
this.buffer += line;
if (term) {
this.buffer += term;
}
const sink = overflow ? await this.#fileSink() : null;
this.lineEnds.push(this.buffer.length);
this.onLine?.(line);
this.#buffer += data;
await sink?.write(data);
if (this.buffer.length > this.spillThreshold) {
this.spillHalf();
if (this.#buffer.length > this.#spillThreshold) {
this.#buffer = this.#buffer.slice(-this.#spillThreshold);
}
}
private pushChunk(line: string): void {
this.onChunk?.(line);
this.pushLine(line);
async #fileSink(): Promise<Bun.FileSink> {
if (!this.#file) {
const filePath = this.#allocateFilePath();
this.#file = {
path: filePath,
sink: Bun.file(filePath).writer(),
};
await this.#file.sink.write(this.#buffer);
}
return this.#file.sink;
}
private getFileSink(): Bun.FileSink {
if (!this.fileSink) {
const filePath = this.allocateFilePath();
this.filePath = filePath;
this.fileSink = Bun.file(filePath).writer();
}
return this.fileSink;
async push(chunk: string): Promise<void> {
chunk = sanitizeText(chunk);
await this.#pushSanitized(chunk);
}
private spillHalf(): void {
const target = this.buffer.length >>> 1;
createInput(): WritableStream<Uint8Array | string> {
let decoder: TextDecoder | undefined;
let finalize = async () => {};
// Binary search: first line ending >= target
let lo = 0;
let hi = this.lineEnds.length;
while (lo < hi) {
const mid = (lo + hi) >>> 1;
if (this.lineEnds[mid] < target) {
lo = mid + 1;
} else {
hi = mid;
}
}
// Clamp: evict at least 1 line, keep at least 1 line
const splitIdx = Math.max(1, Math.min(lo, this.lineEnds.length - 1));
const splitPos = this.lineEnds[splitIdx - 1];
// Write evicted portion to file
this.getFileSink().write(this.buffer.slice(0, splitPos));
// Truncate buffer, shift line positions
this.buffer = this.buffer.slice(splitPos);
const remaining = this.lineEnds.length - splitIdx;
for (let i = 0; i < remaining; i++) {
this.lineEnds[i] = this.lineEnds[i + splitIdx] - splitPos;
}
this.lineEnds.length = remaining;
}
createWritable(): WritableStream<Uint8Array> {
const decoder = new TextDecoder("utf-8", { ignoreBOM: true });
let buf = "";
const flushLines = () => {
let start = 0;
while (true) {
const nl = buf.indexOf("\n", start);
if (nl === -1) break;
this.pushChunk(buf.slice(start, nl + 1));
start = nl + 1;
}
buf = buf.slice(start);
};
const finalize = () => {
buf += sanitizeText(decoder.decode());
flushLines();
buf = buf.trimEnd();
if (buf) {
this.pushChunk(`${buf}\n`);
}
};
return new WritableStream<Uint8Array>({
write: (chunk) => {
buf += sanitizeText(decoder.decode(chunk, { stream: true }));
flushLines();
return new WritableStream<Uint8Array | string>({
write: async (chunk) => {
if (typeof chunk === "string") {
await this.push(chunk);
} else {
if (!decoder) {
const dec = new TextDecoder("utf-8", { ignoreBOM: true });
decoder = dec;
finalize = async () => {
await this.push(dec.decode());
};
}
await this.push(decoder.decode(chunk, { stream: true }));
}
},
close: finalize,
abort: finalize,
});
}
createStringWritable(): WritableStream<string> {
let buf = "";
async dump(notice?: string): Promise<OutputResult> {
const noticeLine = notice ? `[${notice}]\n` : "";
const flushLines = () => {
let start = 0;
while (true) {
const nl = buf.indexOf("\n", start);
if (nl === -1) break;
this.pushChunk(buf.slice(start, nl + 1));
start = nl + 1;
}
buf = buf.slice(start);
};
const finalize = () => {
flushLines();
buf = buf.trimEnd();
if (buf) {
this.pushChunk(`${buf}\n`);
}
};
return new WritableStream<string>({
write: (chunk) => {
buf += sanitizeText(chunk);
flushLines();
},
close: finalize,
abort: finalize,
});
}
async close(): Promise<void> {
await this.fileSink?.end();
}
dump(annotation?: string): OutputResult {
let output = this.buffer;
if (annotation) {
output += `\n${annotation}\n`;
if (this.#file) {
await this.#file.sink.end();
return { output: `${noticeLine}...${this.#buffer}`, truncated: true, fullOutputPath: this.#file.path };
} else {
return { output: `${noticeLine}${this.#buffer}`, truncated: false };
}
if (!this.filePath) {
return { output, truncated: false };
}
this.fileSink!.write(this.buffer);
this.fileSink!.flush();
return {
output,
truncated: true,
fullOutputPath: this.filePath,
};
}
}
+10 -44
View File
@@ -6,14 +6,14 @@ import { Type } from "@sinclair/typebox";
import { truncateToVisualLines } from "../../modes/interactive/components/visual-truncate";
import type { Theme } from "../../modes/interactive/theme/theme";
import bashDescription from "../../prompts/tools/bash.md" with { type: "text" };
import { type BashExecutorOptions, executeBash, executeBashWithOperations } from "../bash-executor";
import { type BashExecutorOptions, executeBash } from "../bash-executor";
import type { RenderResultOptions } from "../custom-tools/types";
import { renderPromptTemplate } from "../prompt-templates";
import { checkBashInterception, checkSimpleLsInterception } from "./bash-interceptor";
import type { ToolSession } from "./index";
import { resolveToCwd } from "./path-utils";
import { ToolUIKit } from "./render-utils";
import { DEFAULT_MAX_BYTES, formatSize, type TruncationResult, truncateTail } from "./truncate";
import { formatTailTruncationNotice, type TruncationResult, truncateTail } from "./truncate";
export const BASH_DEFAULT_PREVIEW_LINES = 10;
@@ -31,32 +31,12 @@ export interface BashToolDetails {
fullOutput?: string;
}
/**
* Pluggable operations for bash execution.
* Override to delegate command execution to remote systems.
*/
export interface BashOperations {
exec: (
command: string,
cwd: string,
options: {
onData: (data: Buffer) => void;
signal?: AbortSignal;
timeout?: number;
},
) => Promise<{ exitCode: number | null }>;
}
export interface BashToolOptions {
/** Custom operations for command execution. Default: local shell */
operations?: BashOperations;
}
export interface BashToolOptions {}
/**
* Bash tool implementation.
*
* Executes bash commands with optional timeout and working directory.
* Supports custom operations for remote execution.
*/
export class BashTool implements AgentTool<typeof bashSchema, BashToolDetails> {
public readonly name = "bash";
@@ -65,11 +45,9 @@ export class BashTool implements AgentTool<typeof bashSchema, BashToolDetails> {
public readonly parameters = bashSchema;
private readonly session: ToolSession;
private readonly options?: BashToolOptions;
constructor(session: ToolSession, options?: BashToolOptions) {
constructor(session: ToolSession) {
this.session = session;
this.options = options;
this.description = renderPromptTemplate(bashDescription);
}
@@ -125,12 +103,8 @@ export class BashTool implements AgentTool<typeof bashSchema, BashToolDetails> {
},
};
// Use custom operations if provided, otherwise use default local executor
const result = this.options?.operations
? await executeBashWithOperations(command, commandCwd, this.options.operations, executorOptions)
: await executeBash(command, executorOptions);
// Handle errors
const result = await executeBash(command, executorOptions);
if (result.cancelled) {
throw new Error(result.output || "Command aborted");
}
@@ -147,18 +121,10 @@ export class BashTool implements AgentTool<typeof bashSchema, BashToolDetails> {
fullOutputPath: result.fullOutputPath,
fullOutput: currentOutput,
};
const startLine = truncation.totalLines - truncation.outputLines + 1;
const endLine = truncation.totalLines;
if (truncation.lastLinePartial) {
const lastLineSize = formatSize(Buffer.byteLength(result.output.split("\n").pop() || "", "utf-8"));
outputText += `\n\n[Showing last ${formatSize(truncation.outputBytes)} of line ${endLine} (line is ${lastLineSize}). Full output: ${result.fullOutputPath}]`;
} else if (truncation.truncatedBy === "lines") {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines}. Full output: ${result.fullOutputPath}]`;
} else {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines} (${formatSize(DEFAULT_MAX_BYTES)} limit). Full output: ${result.fullOutputPath}]`;
}
outputText += formatTailTruncationNotice(truncation, {
fullOutputPath: result.fullOutputPath,
originalContent: result.output,
});
}
if (result.exitCode !== 0 && result.exitCode !== undefined) {
@@ -263,7 +229,7 @@ export const bashToolRenderer = {
warnings.push(`Truncated: showing ${truncation.outputLines} of ${truncation.totalLines} lines`);
} else {
warnings.push(
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes ?? DEFAULT_MAX_BYTES)} limit)`,
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes)} limit)`,
);
}
}
@@ -1,5 +1,5 @@
export { AskTool, type AskToolDetails } from "./ask";
export { type BashOperations, BashTool, type BashToolDetails, type BashToolOptions } from "./bash";
export { BashTool, type BashToolDetails, type BashToolOptions } from "./bash";
export { CalculatorTool, type CalculatorToolDetails } from "./calculator";
export { CompleteTool } from "./complete";
// Exa MCP tools (22 tools)
+6 -15
View File
@@ -14,7 +14,7 @@ import type { PreludeHelper, PythonStatusEvent } from "../python-kernel";
import type { ToolSession } from "./index";
import { resolveToCwd } from "./path-utils";
import { getTreeBranch, getTreeContinuePrefix, shortenPath, ToolUIKit, truncate } from "./render-utils";
import { DEFAULT_MAX_BYTES, formatSize, type TruncationResult, truncateTail } from "./truncate";
import { DEFAULT_MAX_BYTES, formatTailTruncationNotice, type TruncationResult, truncateTail } from "./truncate";
export const PYTHON_DEFAULT_PREVIEW_LINES = 10;
@@ -234,7 +234,6 @@ export class PythonTool implements AgentTool<typeof pythonSchema> {
let details: PythonToolDetails | undefined;
if (truncation.truncated) {
const fullOutputSuffix = result.fullOutputPath ? ` Full output: ${result.fullOutputPath}` : "";
details = {
truncation,
fullOutputPath: result.fullOutputPath,
@@ -242,18 +241,10 @@ export class PythonTool implements AgentTool<typeof pythonSchema> {
images,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
};
const startLine = truncation.totalLines - truncation.outputLines + 1;
const endLine = truncation.totalLines;
if (truncation.lastLinePartial) {
const lastLineSize = formatSize(Buffer.byteLength(result.output.split("\n").pop() || "", "utf-8"));
outputText += `\n\n[Showing last ${formatSize(truncation.outputBytes)} of line ${endLine} (line is ${lastLineSize})${fullOutputSuffix}]`;
} else if (truncation.truncatedBy === "lines") {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines}${fullOutputSuffix}]`;
} else {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines} (${formatSize(DEFAULT_MAX_BYTES)} limit)${fullOutputSuffix}]`;
}
outputText += formatTailTruncationNotice(truncation, {
fullOutputPath: result.fullOutputPath,
originalContent: result.output,
});
}
if (!details && (jsonOutputs.length > 0 || images.length > 0 || statusEvents.length > 0)) {
@@ -685,7 +676,7 @@ export const pythonToolRenderer = {
warnings.push(`Truncated: showing ${truncation.outputLines} of ${truncation.totalLines} lines`);
} else {
warnings.push(
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes ?? DEFAULT_MAX_BYTES)} limit)`,
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes)} limit)`,
);
}
}
+6 -14
View File
@@ -14,7 +14,7 @@ import { ensureHostInfo, getHostInfoForHost } from "../ssh/connection-manager";
import { executeSSH } from "../ssh/ssh-executor";
import type { ToolSession } from "./index";
import { ToolUIKit } from "./render-utils";
import { DEFAULT_MAX_BYTES, formatSize, type TruncationResult, truncateTail } from "./truncate";
import { formatTailTruncationNotice, type TruncationResult, truncateTail } from "./truncate";
const sshSchema = Type.Object({
host: Type.String({ description: "Host name from ssh.json or .ssh.json" }),
@@ -193,18 +193,10 @@ export class SshTool implements AgentTool<typeof sshSchema, SSHToolDetails> {
truncation,
fullOutputPath: result.fullOutputPath,
};
const startLine = truncation.totalLines - truncation.outputLines + 1;
const endLine = truncation.totalLines;
if (truncation.lastLinePartial) {
const lastLineSize = formatSize(Buffer.byteLength(result.output.split("\n").pop() || "", "utf-8"));
outputText += `\n\n[Showing last ${formatSize(truncation.outputBytes)} of line ${endLine} (line is ${lastLineSize}). Full output: ${result.fullOutputPath}]`;
} else if (truncation.truncatedBy === "lines") {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines}. Full output: ${result.fullOutputPath}]`;
} else {
outputText += `\n\n[Showing lines ${startLine}-${endLine} of ${truncation.totalLines} (${formatSize(DEFAULT_MAX_BYTES)} limit). Full output: ${result.fullOutputPath}]`;
}
outputText += formatTailTruncationNotice(truncation, {
fullOutputPath: result.fullOutputPath,
originalContent: result.output,
});
}
if (result.exitCode !== 0 && result.exitCode !== undefined) {
@@ -311,7 +303,7 @@ export const sshToolRenderer = {
warnings.push(`Truncated: showing ${truncation.outputLines} of ${truncation.totalLines} lines`);
} else {
warnings.push(
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes ?? DEFAULT_MAX_BYTES)} limit)`,
`Truncated: ${truncation.outputLines} lines shown (${ui.formatBytes(truncation.maxBytes)} limit)`,
);
}
}
@@ -289,3 +289,95 @@ export function truncateLine(
}
return { text: `${line.slice(0, maxChars)}... [truncated]`, wasTruncated: true };
}
// =============================================================================
// Truncation notice formatting
// =============================================================================
export interface TailTruncationNoticeOptions {
/** Path to full output file (e.g., from bash/python executor) */
fullOutputPath?: string;
/** Original content for computing last line size when lastLinePartial */
originalContent?: string;
/** Additional suffix to append inside the brackets */
suffix?: string;
}
/**
* Format a truncation notice for tail-truncated output (bash, python, ssh).
* Returns empty string if not truncated.
*
* Examples:
* - "[Showing last 50KB of line 1000 (line is 2.1MB). Full output: /tmp/out.txt]"
* - "[Showing lines 500-1000 of 1000. Full output: /tmp/out.txt]"
* - "[Showing lines 500-1000 of 1000 (50KB limit). Full output: /tmp/out.txt]"
*/
export function formatTailTruncationNotice(
truncation: TruncationResult,
options: TailTruncationNoticeOptions = {},
): string {
if (!truncation.truncated) {
return "";
}
const { fullOutputPath, originalContent, suffix = "" } = options;
const startLine = truncation.totalLines - truncation.outputLines + 1;
const endLine = truncation.totalLines;
const fullOutputPart = fullOutputPath ? `. Full output: ${fullOutputPath}` : "";
let notice: string;
if (truncation.lastLinePartial) {
let lastLineSizePart = "";
if (originalContent) {
const lastLine = originalContent.split("\n").pop() || "";
lastLineSizePart = ` (line is ${formatSize(Buffer.byteLength(lastLine, "utf-8"))})`;
}
notice = `[Showing last ${formatSize(truncation.outputBytes)} of line ${endLine}${lastLineSizePart}${fullOutputPart}${suffix}]`;
} else if (truncation.truncatedBy === "lines") {
notice = `[Showing lines ${startLine}-${endLine} of ${truncation.totalLines}${fullOutputPart}${suffix}]`;
} else {
notice = `[Showing lines ${startLine}-${endLine} of ${truncation.totalLines} (${formatSize(truncation.maxBytes)} limit)${fullOutputPart}${suffix}]`;
}
return `\n\n${notice}`;
}
export interface HeadTruncationNoticeOptions {
/** 1-indexed start line number (default: 1) */
startLine?: number;
/** Total lines in the original file (for "of N" display) */
totalFileLines?: number;
}
/**
* Format a truncation notice for head-truncated output (read tool).
* Returns empty string if not truncated.
*
* Examples:
* - "[Showing lines 1-2000 of 5000. Use offset=2001 to continue]"
* - "[Showing lines 100-2099 of 5000 (50KB limit). Use offset=2100 to continue]"
*/
export function formatHeadTruncationNotice(
truncation: TruncationResult,
options: HeadTruncationNoticeOptions = {},
): string {
if (!truncation.truncated) {
return "";
}
const startLineDisplay = options.startLine ?? 1;
const totalFileLines = options.totalFileLines ?? truncation.totalLines;
const endLineDisplay = startLineDisplay + truncation.outputLines - 1;
const nextOffset = endLineDisplay + 1;
let notice: string;
if (truncation.truncatedBy === "lines") {
notice = `[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines}. Use offset=${nextOffset} to continue]`;
} else {
notice = `[Showing lines ${startLineDisplay}-${endLineDisplay} of ${totalFileLines} (${formatSize(truncation.maxBytes)} limit). Use offset=${nextOffset} to continue]`;
}
return `\n\n${notice}`;
}
@@ -24,7 +24,9 @@ import { convertWithMarkitdown, fetchBinary } from "./web-scrapers/utils";
// Types and Constants
// =============================================================================
const DEFAULT_TIMEOUT = 20;
const MIN_TIMEOUT = 1_000;
const DEFAULT_TIMEOUT = 20_000;
const MAX_TIMEOUT = 45_000;
// Convertible document types (markitdown supported)
const CONVERTIBLE_MIMES = new Set([
@@ -109,7 +111,7 @@ async function exec(
): Promise<{ stdout: string; stderr: string; ok: boolean }> {
const proc = ptree.cspawn([cmd, ...args], {
stdin: options?.input ? "pipe" : null,
timeout: options?.timeout,
timeout: options?.timeout ? options.timeout * 1000 : undefined,
});
if (options?.input) {
@@ -244,7 +246,7 @@ async function tryMdSuffix(url: string, timeout: number, signal?: AbortSignal):
if (signal?.aborted) {
return null;
}
const result = await loadPage(candidate, { timeout: Math.min(timeout, 5), signal });
const result = await loadPage(candidate, { timeout: Math.min(timeout, MAX_TIMEOUT), signal });
if (result.ok && result.content.trim().length > 100 && !looksLikeHtml(result.content)) {
return result.content;
}
@@ -910,7 +912,7 @@ export class WebFetchTool implements AgentTool<typeof webFetchSchema, WebFetchTo
}
// Clamp timeout
const effectiveTimeout = Math.min(Math.max(timeout, 1), 120);
const effectiveTimeout = Math.min(Math.max(timeout, MIN_TIMEOUT), MAX_TIMEOUT);
const result = await renderUrl(url, effectiveTimeout, raw, signal);
@@ -1,36 +1,13 @@
import { rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import * as path from "node:path";
import { $ } from "bun";
import { ptree } from "@oh-my-pi/pi-utils";
import { nanoid } from "nanoid";
import { ensureTool } from "../../../utils/tools-manager";
import { createRequestSignal } from "./types";
const MAX_BYTES = 50 * 1024 * 1024; // 50MB for binary files
interface ExecResult {
stdout: string;
stderr: string;
ok: boolean;
exitCode: number;
}
async function exec(
cmd: string,
args: string[],
options?: { timeout?: number; input?: string | Buffer },
): Promise<ExecResult> {
void options;
const result = await $`${cmd} ${args}`.quiet().nothrow();
const decoder = new TextDecoder();
return {
stdout: result.stdout ? decoder.decode(result.stdout) : "",
stderr: result.stderr ? decoder.decode(result.stderr) : "",
ok: result.exitCode === 0,
exitCode: result.exitCode ?? -1,
};
}
export interface ConvertResult {
content: string;
ok: boolean;
@@ -72,16 +49,16 @@ export async function convertWithMarkitdown(
try {
await Bun.write(tmpFile, content);
const result = await exec(markitdown, [tmpFile], { timeout });
if (!result.ok) {
const stderr = result.stderr.trim();
const result = await ptree.cspawn([markitdown, tmpFile], { timeout });
const [stdout, stderr, exitCode] = await Promise.all([result.stdout.text(), result.stderr.text(), result.exited]);
if (exitCode !== 0) {
return {
content: result.stdout,
content: stdout,
ok: false,
error: stderr.length > 0 ? stderr : `markitdown failed (exit ${result.exitCode})`,
error: stderr.length > 0 ? stderr : `markitdown failed (exit ${exitCode})`,
};
}
return { content: result.stdout, ok: true };
return { content: stdout, ok: true };
} finally {
try {
await rm(tmpFile, { force: true });
@@ -64,7 +64,6 @@ function buildSystemBlocks(
return buildAnthropicSystemBlocks(systemPrompt, {
includeClaudeCodeInstruction: includeClaudeCode,
includeCacheControl: auth.isOAuth,
extraInstructions,
});
}
-1
View File
@@ -194,7 +194,6 @@ export {
export { type FileSlashCommand, loadSlashCommands as discoverSlashCommands } from "./core/slash-commands";
// Tools (detail types and utilities)
export {
type BashOperations,
type BashToolDetails,
DEFAULT_MAX_BYTES,
DEFAULT_MAX_LINES,
+3 -3
View File
@@ -142,7 +142,7 @@ export class ChildProcess {
#stderrStream?: ReadableStream<Uint8Array>;
#exitReason?: Exception;
#exitReasonPending?: Exception;
#exited: Promise<void>;
#exited: Promise<number>;
#resolveExited: (ex?: PromiseLike<Exception> | Exception) => void;
constructor(proc: PipedSubprocess) {
@@ -221,7 +221,7 @@ export class ChildProcess {
const { promise, resolve } = Promise.withResolvers<Exception | undefined>();
this.#exited = promise.then((ex?: Exception) => {
if (!ex) return; // success, no exception
if (!ex) return proc.exitCode ?? -1337; // success, no exception
if (proc.killed && this.#exitReasonPending) {
ex = this.#exitReasonPending; // propagate reason if killed
}
@@ -245,7 +245,7 @@ export class ChildProcess {
get pid(): number | undefined {
return this.#proc.pid;
}
get exited(): Promise<void> {
get exited(): Promise<number> {
return this.#exited;
}
get exitCode(): number | null {