diff --git a/AGENTS.md b/AGENTS.md index 923d4f74a..a02d655de 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -8,15 +8,15 @@ This repo contains multiple packages, but **`packages/coding-agent/`** is the pr ### Package Structure -| Package | Description | -| ----------------------- | ----------------------------------------------------- | -| `packages/ai` | Multi-provider LLM client with streaming support | -| `packages/agent` | Agent runtime with tool calling and state management | -| `packages/coding-agent` | Main CLI application (primary focus) | -| `packages/tui` | Terminal UI library with differential rendering | -| `packages/natives` | WASM bindings for native text/image/grep operations | -| `packages/stats` | Local observability dashboard (`omp stats`) | -| `packages/utils` | Shared utilities (logger, streams, temp files) | +| Package | Description | +| ----------------------- | ------------------------------------------------------ | +| `packages/ai` | Multi-provider LLM client with streaming support | +| `packages/agent` | Agent runtime with tool calling and state management | +| `packages/coding-agent` | Main CLI application (primary focus) | +| `packages/tui` | Terminal UI library with differential rendering | +| `packages/natives` | WASM bindings for native text/image/grep operations | +| `packages/stats` | Local observability dashboard (`omp stats`) | +| `packages/utils` | Shared utilities (logger, streams, temp files) | | `crates/pi-natives` | Rust WASM crate for performance-critical text/grep ops | ## Code Quality @@ -29,6 +29,7 @@ This repo contains multiple packages, but **`packages/coding-agent/`** is the pr - **NEVER build prompts in code** — no inline strings, no template literals, no string concatenation. Prompts live in static `.md` files; use Handlebars for any dynamic content. - **Import static text files via Bun** — use `import content from "./prompt.md" with { type: "text" }` instead of `readFileSync` - **Use `Promise.withResolvers()`** instead of `new Promise((resolve, reject) => ...)` — cleaner, avoids callback nesting, and the resolver functions are properly typed: + ```typescript // BAD: Verbose, callback nesting const promise = new Promise((resolve, reject) => { ... }); @@ -311,28 +312,7 @@ const entries = lines.map((line) => JSON.parse(line)); const entries = Bun.JSONL.parse(text); ``` -**For streaming JSONL** (SSE, JSON-RPC, subprocess output), use `Bun.JSONL.parseChunk()`: - -```typescript -// BAD: Manual buffering and line splitting -let buffer = ""; -for await (const chunk of stream) { - buffer += decoder.decode(chunk); - const lines = buffer.split("\n"); - buffer = lines.pop() ?? ""; - for (const line of lines) { - if (line.trim()) yield JSON.parse(line); - } -} - -// GOOD: Bun handles buffering and parsing -let buffer: Uint8Array | undefined; -for await (const chunk of stream) { - const { values, remainder } = Bun.JSONL.parseChunk(chunk, buffer); - buffer = remainder; - for (const value of values) yield value; -} -``` +**For streaming JSONL** (SSE, JSON-RPC, subprocess output), use `Bun.JSONL.parseChunk() | Bun.JSONL.parse()` without decoding to string: ### Terminal Width and Wrapping @@ -428,20 +408,20 @@ Logs go to `~/.omp/logs/omp.YYYY-MM-DD.log` with automatic rotation. ## Commands -| Command | Description | -| -------------- | ------------------------------------------------ | -| `bun check` | Check all (TypeScript + Rust) | -| `bun check:ts` | Biome check + tsgo type checking | -| `bun check:rs` | Cargo fmt --check + clippy | -| `bun lint` | Lint all | -| `bun lint:ts` | Biome lint | -| `bun lint:rs` | Cargo clippy | -| `bun fmt` | Format all | -| `bun fmt:ts` | Biome format | -| `bun fmt:rs` | Cargo fmt | -| `bun fix` | Fix all (unsafe fixes + format) | -| `bun fix:ts` | Biome --unsafe + format-prompts | -| `bun fix:rs` | Clippy --fix + cargo fmt | +| Command | Description | +| -------------- | -------------------------------- | +| `bun check` | Check all (TypeScript + Rust) | +| `bun check:ts` | Biome check + tsgo type checking | +| `bun check:rs` | Cargo fmt --check + clippy | +| `bun lint` | Lint all | +| `bun lint:ts` | Biome lint | +| `bun lint:rs` | Cargo clippy | +| `bun fmt` | Format all | +| `bun fmt:ts` | Biome format | +| `bun fmt:rs` | Cargo fmt | +| `bun fix` | Fix all (unsafe fixes + format) | +| `bun fix:ts` | Biome --unsafe + format-prompts | +| `bun fix:rs` | Clippy --fix + cargo fmt | - NEVER run: `bun run dev`, `bun test` unless user instructs - Only run specific tests if user instructs: `bun test test/specific.test.ts` diff --git a/bun.lock b/bun.lock index 271ff6665..4c1a7d548 100644 --- a/bun.lock +++ b/bun.lock @@ -91,6 +91,7 @@ "devDependencies": { "@types/bun": "^1.3.8", "@types/diff": "^7.0.2", + "@types/jsdom": "27.0.0", "@types/ms": "^2.1.0", "ms": "^2.1.3", }, @@ -490,6 +491,8 @@ "@types/diff": ["@types/diff@7.0.2", "", {}, "sha512-JSWRMozjFKsGlEjiiKajUjIJVKuKdE3oVy2DNtK+fUo8q82nhFZ2CPQwicAIkXrofahDXrWJ7mjelvZphMS98Q=="], + "@types/jsdom": ["@types/jsdom@27.0.0", "", { "dependencies": { "@types/node": "*", "@types/tough-cookie": "*", "parse5": "^7.0.0" } }, "sha512-NZyFl/PViwKzdEkQg96gtnB8wm+1ljhdDay9ahn4hgb+SfVtPCbm3TlmDUFXTA+MGN3CijicnMhG18SI5H3rFw=="], + "@types/mime-types": ["@types/mime-types@3.0.1", "", {}, "sha512-xRMsfuQbnRq1Ef+C+RKaENOxXX87Ygl38W1vDfPHRku02TgQr+Qd8iivLtAMcR0KF5/29xlnFihkTlbqFrGOVQ=="], "@types/ms": ["@types/ms@2.1.0", "", {}, "sha512-GsCCIZDE/p3i96vtEqx+7dBUGXrc7zeSK3wwPHIaRThS+9OhWIXRqzs4d6k1SVU8g91DrNRWxWUGhp5KXQb2VA=="], @@ -500,6 +503,8 @@ "@types/react-dom": ["@types/react-dom@19.2.3", "", { "peerDependencies": { "@types/react": "^19.2.0" } }, "sha512-jp2L/eY6fn+KgVVQAOqYItbF0VY/YApe5Mz2F0aykSO8gx31bYCZyvSeYxCHKvzHG5eZjc+zyaS5BrBWya2+kQ=="], + "@types/tough-cookie": ["@types/tough-cookie@4.0.5", "", {}, "sha512-/Ad8+nIOV7Rl++6f1BdKxFSMgmoqEoYbHRpPcx3JEfv8VRsQe9Z4mCXeJBzxs7mbHY/XOZZuXlRNfhpVPbs6ZA=="], + "@types/triple-beam": ["@types/triple-beam@1.3.5", "", {}, "sha512-6WaYesThRMCl19iryMYP7/x2OVgCtbIVflDGFpWnb9irXI3UjYE4AzmYuiUKY1AJstGijoY+MgUszMgRxIYTYw=="], "@types/use-sync-external-store": ["@types/use-sync-external-store@0.0.6", "", {}, "sha512-zFDAD+tlpf2r4asuHEj0XH6pY6i0g5NeAHPn+15wk3BV6JA69eERFXC1gyGThDkVa1zCyKr5jox1+2LbV/AMLg=="], @@ -988,7 +993,7 @@ "parse-json": ["parse-json@5.2.0", "", { "dependencies": { "@babel/code-frame": "^7.0.0", "error-ex": "^1.3.1", "json-parse-even-better-errors": "^2.3.0", "lines-and-columns": "^1.1.6" } }, "sha512-ayCKvm/phCGxOkYRSCM82iDwct8/EonSEgCSxWxD7ve6jHggsFl4fZVQBPRNgQoKiuV/odhFrGzQXZwbifC8Rg=="], - "parse5": ["parse5@8.0.0", "", { "dependencies": { "entities": "^6.0.0" } }, "sha512-9m4m5GSgXjL4AjumKzq1Fgfp3Z8rsvjRNbnkVwfu2ImRqE5D0LnY2QfDen18FSY9C573YU5XxSapdHZTZ2WolA=="], + "parse5": ["parse5@7.3.0", "", { "dependencies": { "entities": "^6.0.0" } }, "sha512-IInvU7fabl34qmi9gY8XOVxhYyMyuH2xUNpb2q8/Y+7552KlejkRvqvD19nMoUW/uQGGbqNpA6Tufu5FL5BZgw=="], "parseurl": ["parseurl@1.3.3", "", {}, "sha512-CiyeOxFT/JZyN5m0z9PfXw4SCBJ6Sygz1Dpl0wqjlhDEGGBP1GnsUVEL0p63hoG1fcj3fHynXi9NYO4nWOL+qQ=="], @@ -1296,6 +1301,8 @@ "get-uri/data-uri-to-buffer": ["data-uri-to-buffer@6.0.2", "", {}, "sha512-7hvf7/GW8e86rW0ptuwS3OcBGDjIi6SZva7hCyWC0yYry2cOPmLIjXAUHI6DK2HsnwJd9ifmt57i8eV2n4YNpw=="], + "jsdom/parse5": ["parse5@8.0.0", "", { "dependencies": { "entities": "^6.0.0" } }, "sha512-9m4m5GSgXjL4AjumKzq1Fgfp3Z8rsvjRNbnkVwfu2ImRqE5D0LnY2QfDen18FSY9C573YU5XxSapdHZTZ2WolA=="], + "log-update/strip-ansi": ["strip-ansi@7.1.2", "", { "dependencies": { "ansi-regex": "^6.0.1" } }, "sha512-gmBGslpoQJtgnMAvOVqGZpEz9dyoKTCzy2nfz/n8aIFhN/jCE/rCmcxabB6jOOHV+0WNnylOxaxBQPSvcWklhA=="], "proxy-agent/lru-cache": ["lru-cache@7.18.3", "", {}, "sha512-jumlc0BIUrS3qJGgIkWZsyfAM7NCWiBcCDhnd+3NNM5KbBmLTgHVfWBcg6W+rLUsIpzpERPsvwUP7CckAQSOoA=="], diff --git a/packages/agent/src/proxy.ts b/packages/agent/src/proxy.ts index 6f2e57bcf..08c8ffa94 100644 --- a/packages/agent/src/proxy.ts +++ b/packages/agent/src/proxy.ts @@ -13,7 +13,7 @@ import { type ToolCall, } from "@oh-my-pi/pi-ai"; import { parseStreamingJson } from "@oh-my-pi/pi-ai/utils/json-parse"; -import { readSseEvents } from "@oh-my-pi/pi-utils"; +import { readSseJson } from "@oh-my-pi/pi-utils"; // Create stream class matching ProxyMessageEventStream class ProxyMessageEventStream extends EventStream { @@ -147,24 +147,13 @@ export function streamProxy(model: Model, context: Context, options: ProxyStream throw new Error(errorMessage); } - for await (const event of readSseEvents(response.body!)) { - if (options.signal?.aborted) { - throw new Error("Request aborted by user"); - } - - const data = event.data?.trim(); - if (!data || data === "[DONE]") continue; - const proxyEvent = JSON.parse(data) as ProxyAssistantMessageEvent; - const parsedEvent = processProxyEvent(proxyEvent, partial); + for await (const event of readSseJson(response.body!, options.signal)) { + const parsedEvent = processProxyEvent(event, partial); if (parsedEvent) { stream.push(parsedEvent); } } - if (options.signal?.aborted) { - throw new Error("Request aborted by user"); - } - stream.end(); } catch (error) { const errorMessage = error instanceof Error ? error.message : String(error); diff --git a/packages/ai/README.md b/packages/ai/README.md index 310415a64..75f30c38e 100644 --- a/packages/ai/README.md +++ b/packages/ai/README.md @@ -10,38 +10,38 @@ Unified LLM API with automatic model discovery, provider configuration, token an - [Installation](#installation) - [Quick Start](#quick-start) - [Tools](#tools) - - [Defining Tools](#defining-tools) - - [Handling Tool Calls](#handling-tool-calls) - - [Streaming Tool Calls with Partial JSON](#streaming-tool-calls-with-partial-json) - - [Validating Tool Arguments](#validating-tool-arguments) - - [Complete Event Reference](#complete-event-reference) + - [Defining Tools](#defining-tools) + - [Handling Tool Calls](#handling-tool-calls) + - [Streaming Tool Calls with Partial JSON](#streaming-tool-calls-with-partial-json) + - [Validating Tool Arguments](#validating-tool-arguments) + - [Complete Event Reference](#complete-event-reference) - [Image Input](#image-input) - [Thinking/Reasoning](#thinkingreasoning) - - [Unified Interface](#unified-interface-streamsimplecompletesimple) - - [Provider-Specific Options](#provider-specific-options-streamcomplete) - - [Streaming Thinking Content](#streaming-thinking-content) + - [Unified Interface](#unified-interface-streamsimplecompletesimple) + - [Provider-Specific Options](#provider-specific-options-streamcomplete) + - [Streaming Thinking Content](#streaming-thinking-content) - [Stop Reasons](#stop-reasons) - [Error Handling](#error-handling) - - [Aborting Requests](#aborting-requests) - - [Continuing After Abort](#continuing-after-abort) + - [Aborting Requests](#aborting-requests) + - [Continuing After Abort](#continuing-after-abort) - [APIs, Models, and Providers](#apis-models-and-providers) - - [Providers and Models](#providers-and-models) - - [Querying Providers and Models](#querying-providers-and-models) - - [Custom Models](#custom-models) - - [OpenAI Compatibility Settings](#openai-compatibility-settings) - - [Type Safety](#type-safety) + - [Providers and Models](#providers-and-models) + - [Querying Providers and Models](#querying-providers-and-models) + - [Custom Models](#custom-models) + - [OpenAI Compatibility Settings](#openai-compatibility-settings) + - [Type Safety](#type-safety) - [Cross-Provider Handoffs](#cross-provider-handoffs) - [Context Serialization](#context-serialization) - [Browser Usage](#browser-usage) - - [Environment Variables](#environment-variables-nodejs-only) - - [Checking Environment Variables](#checking-environment-variables) + - [Environment Variables](#environment-variables-nodejs-only) + - [Checking Environment Variables](#checking-environment-variables) - [OAuth Providers](#oauth-providers) - - [Vertex AI (ADC)](#vertex-ai-adc) - - [CLI Login](#cli-login) - - [Programmatic OAuth](#programmatic-oauth) - - [Login Flow Example](#login-flow-example) - - [Using OAuth Tokens](#using-oauth-tokens) - - [Provider Notes](#provider-notes) + - [Vertex AI (ADC)](#vertex-ai-adc) + - [CLI Login](#cli-login) + - [Programmatic OAuth](#programmatic-oauth) + - [Login Flow Example](#login-flow-example) + - [Using OAuth Tokens](#using-oauth-tokens) + - [Provider Notes](#provider-notes) - [License](#license) ## Supported Providers @@ -267,7 +267,7 @@ context.messages.push({ toolName: "generate_chart", content: [ { type: "text", text: "Generated chart showing temperature trends" }, - { type: "image", data: imageBuffer.toString("base64"), mimeType: "image/png" }, + { type: "image", data: imageBuffer.toBase64(), mimeType: "image/png" }, ], isError: false, timestamp: Date.now(), @@ -390,7 +390,7 @@ if (model.input.includes("image")) { } const imageBuffer = fs.readFileSync("image.png"); -const base64Image = imageBuffer.toString("base64"); +const base64Image = imageBuffer.toBase64(); const response = await complete(model, { messages: [ @@ -443,7 +443,7 @@ const response = await completeSimple( }, { reasoning: "medium", // 'minimal' | 'low' | 'medium' | 'high' | 'xhigh' (xhigh maps to high on non-OpenAI providers) - }, + } ); // Access thinking and text blocks @@ -562,7 +562,7 @@ const s = stream( }, { signal, - }, + } ); for await (const event of s) { @@ -877,7 +877,7 @@ const response = await complete( }, { apiKey: "your-api-key", - }, + } ); ``` @@ -1065,7 +1065,7 @@ const response = await complete( { messages: [{ role: "user", content: "Hello!" }], }, - { apiKey: result.apiKey }, + { apiKey: result.apiKey } ); ``` diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 03624f822..664e7e782 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -1,5 +1,5 @@ import * as os from "node:os"; -import { $env, abortableSleep } from "@oh-my-pi/pi-utils"; +import { $env, abortableSleep, readSseJson } from "@oh-my-pi/pi-utils"; import type { ResponseFunctionToolCall, ResponseInput, @@ -38,7 +38,7 @@ import { URL_PATHS, } from "./openai-codex/constants"; import { type CodexRequestOptions, type RequestBody, transformRequestBody } from "./openai-codex/request-transformer"; -import { parseCodexError, parseCodexSseStream } from "./openai-codex/response-handler"; +import { parseCodexError } from "./openai-codex/response-handler"; import { transformMessages } from "./transform-messages"; export interface OpenAICodexResponsesOptions extends StreamOptions { @@ -234,7 +234,7 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" const blocks = output.content; const blockIndex = () => blocks.length - 1; - for await (const rawEvent of parseCodexSseStream(response)) { + for await (const rawEvent of readSseJson>(response.body!, options?.signal)) { const eventType = typeof rawEvent.type === "string" ? rawEvent.type : ""; if (!eventType) continue; diff --git a/packages/ai/src/providers/openai-codex/response-handler.ts b/packages/ai/src/providers/openai-codex/response-handler.ts index 8c405f9ec..9b3b1cbee 100644 --- a/packages/ai/src/providers/openai-codex/response-handler.ts +++ b/packages/ai/src/providers/openai-codex/response-handler.ts @@ -1,5 +1,3 @@ -import { readSseData } from "@oh-my-pi/pi-utils"; - export type CodexRateLimit = { used_percent?: number; window_minutes?: number; @@ -71,16 +69,6 @@ export async function parseCodexError(response: Response): Promise> { - if (!response.body) { - return; - } - - for await (const data of readSseData>(response.body)) { - yield data; - } -} - function toNumber(v: string | null): number | undefined { if (v == null) return undefined; const n = Number(v); diff --git a/packages/ai/test/image-limits.test.ts b/packages/ai/test/image-limits.test.ts index 5c0e8ce70..7a185f313 100644 --- a/packages/ai/test/image-limits.test.ts +++ b/packages/ai/test/image-limits.test.ts @@ -85,7 +85,7 @@ async function generateImage(width: number, height: number, filename: string): P const filepath = path.join(TEMP_DIR, filename); execSync(`magick -size ${width}x${height} xc:red "${filepath}"`, { stdio: "ignore" }); const buffer = await fs.promises.readFile(filepath); - return buffer.toString("base64"); + return buffer.toBase64(); } /** @@ -111,7 +111,7 @@ async function generateImageWithSize(targetBytes: number, filename: string): Pro } const buffer = await fs.promises.readFile(filepath); - return buffer.toString("base64"); + return buffer.toBase64(); } /** diff --git a/packages/ai/test/image-tool-result.test.ts b/packages/ai/test/image-tool-result.test.ts index 1f58f152a..9a82f1b5f 100644 --- a/packages/ai/test/image-tool-result.test.ts +++ b/packages/ai/test/image-tool-result.test.ts @@ -34,7 +34,7 @@ async function handleToolWithImageResult(model: Model, o // Read the test image const imagePath = path.join(import.meta.dir, "data", "red-circle.png"); const imageBuffer = await fs.readFile(imagePath); - const base64Image = imageBuffer.toString("base64"); + const base64Image = imageBuffer.toBase64(); // Define a tool that returns only an image (no text) const getImageSchema = Type.Object({}); @@ -122,7 +122,7 @@ async function handleToolWithTextAndImageResult(model: Model { const payload = Buffer.from( JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: "acc_test" } }), "utf8", - ).toString("base64"); + ).toBase64(); const token = `aaa.${payload}.bbb`; const sse = `${[ @@ -135,7 +135,7 @@ describe("openai-codex streaming", () => { const payload = Buffer.from( JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: "acc_test" } }), "utf8", - ).toString("base64"); + ).toBase64(); const token = `aaa.${payload}.bbb`; const sse = `${[ @@ -236,7 +236,7 @@ describe("openai-codex streaming", () => { const payload = Buffer.from( JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: "acc_test" } }), "utf8", - ).toString("base64"); + ).toBase64(); const token = `aaa.${payload}.bbb`; const sse = `${[ @@ -270,14 +270,6 @@ describe("openai-codex streaming", () => { })}`, ].join("\n\n")}\n\n`; - const encoder = new TextEncoder(); - const stream = new ReadableStream({ - start(controller) { - controller.enqueue(encoder.encode(sse)); - controller.close(); - }, - }); - const fetchMock = vi.fn(async (input: string | URL, init?: RequestInit) => { const url = typeof input === "string" ? input : input.toString(); if (url === "https://api.github.com/repos/openai/codex/releases/latest") { @@ -292,7 +284,7 @@ describe("openai-codex streaming", () => { expect(headers?.has("conversation_id")).toBe(false); expect(headers?.has("session_id")).toBe(false); - return new Response(stream, { + return new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" }, }); diff --git a/packages/ai/test/stream.test.ts b/packages/ai/test/stream.test.ts index 239a7e9ea..cb9087636 100644 --- a/packages/ai/test/stream.test.ts +++ b/packages/ai/test/stream.test.ts @@ -220,7 +220,7 @@ async function handleImage(model: Model, options?: Optio // Read the test image const imagePath = path.join(import.meta.dir, "data", "red-circle.png"); const imageBuffer = await fs.readFile(imagePath); - const base64Image = imageBuffer.toString("base64"); + const base64Image = imageBuffer.toBase64(); const imageContent: ImageContent = { type: "image", diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index 3581f9101..9b18ec2ae 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -105,7 +105,8 @@ "@types/diff": "^7.0.2", "@types/ms": "^2.1.0", "@types/bun": "^1.3.8", - "ms": "^2.1.3" + "ms": "^2.1.3", + "@types/jsdom": "27.0.0" }, "keywords": [ "coding-agent", diff --git a/packages/coding-agent/src/cli/file-processor.ts b/packages/coding-agent/src/cli/file-processor.ts index 239f907f4..33f8eab5e 100644 --- a/packages/coding-agent/src/cli/file-processor.ts +++ b/packages/coding-agent/src/cli/file-processor.ts @@ -49,7 +49,7 @@ export async function processFileArguments(fileArgs: string[], options?: Process if (mimeType) { // Handle image file - const base64Content = buffer.toString("base64"); + const base64Content = buffer.toBase64(); let attachment: ImageContent; let dimensionNote: string | undefined; diff --git a/packages/coding-agent/src/export/html/index.ts b/packages/coding-agent/src/export/html/index.ts index d9b43224f..f8e28db6b 100644 --- a/packages/coding-agent/src/export/html/index.ts +++ b/packages/coding-agent/src/export/html/index.ts @@ -102,7 +102,7 @@ interface SessionData { /** Generate HTML from bundled template with runtime substitutions. */ async function generateHtml(sessionData: SessionData, themeName?: string): Promise { const themeVars = await generateThemeVars(themeName); - const sessionDataBase64 = Buffer.from(JSON.stringify(sessionData)).toString("base64"); + const sessionDataBase64 = Buffer.from(JSON.stringify(sessionData)).toBase64(); return TEMPLATE.replace("", ``).replace( "{{SESSION_DATA}}", diff --git a/packages/coding-agent/src/mcp/json-rpc.ts b/packages/coding-agent/src/mcp/json-rpc.ts index 779213e12..838901483 100644 --- a/packages/coding-agent/src/mcp/json-rpc.ts +++ b/packages/coding-agent/src/mcp/json-rpc.ts @@ -13,8 +13,8 @@ export function parseSSE(text: string): unknown { if (line.startsWith("data: ")) { const data = line.slice(6).trim(); if (data === "[DONE]") continue; - const result = Bun.JSONL.parseChunk(`${data}\n`); - if (result.values.length > 0) return result.values[0]; + const result = JSON.parse(data) as unknown; + if (result) return result; } } // Fallback: try parsing entire response as JSON diff --git a/packages/coding-agent/src/mcp/transports/http.ts b/packages/coding-agent/src/mcp/transports/http.ts index aef702ce7..bafc01799 100644 --- a/packages/coding-agent/src/mcp/transports/http.ts +++ b/packages/coding-agent/src/mcp/transports/http.ts @@ -4,7 +4,7 @@ * Implements JSON-RPC 2.0 over HTTP POST with optional SSE streaming. * Based on MCP spec 2025-03-26. */ -import { readSseEvents } from "@oh-my-pi/pi-utils"; +import { readSseJson } from "@oh-my-pi/pi-utils"; import type { JsonRpcMessage, JsonRpcResponse, @@ -86,26 +86,11 @@ export class HttpTransport implements MCPTransport { return; } - let buffer = ""; // Read SSE stream - for await (const event of readSseEvents(response.body)) { + for await (const message of readSseJson(response.body, this.sseConnection.signal)) { if (!this._connected) break; - const data = event.data?.trim(); - if (!data || data === "[DONE]") continue; - buffer += data; - if (!data.endsWith("\n")) { - buffer += "\n"; - } - const result = Bun.JSONL.parseChunk(buffer); - buffer = buffer.slice(result.read); - if (result.error) { - buffer = ""; - continue; - } - for (const message of result.values as JsonRpcMessage[]) { - if ("method" in message && !("id" in message)) { - this.onNotification?.(message.method, message.params); - } + if ("method" in message && !("id" in message)) { + this.onNotification?.(message.method, message.params); } } } catch (error) { @@ -182,40 +167,18 @@ export class HttpTransport implements MCPTransport { const timeout = this.config.timeout ?? 30000; const parse = async (): Promise => { - let buffer = ""; - for await (const event of readSseEvents(response.body!)) { - const data = event.data?.trim(); - if (!data || data === "[DONE]") continue; - buffer += data; - if (!data.endsWith("\n")) { - buffer += "\n"; - } - const result = Bun.JSONL.parseChunk(buffer); - buffer = buffer.slice(result.read); - if (result.error) { - buffer = ""; - continue; + for await (const message of readSseJson(response.body!)) { + if ("id" in message && message.id === expectedId && ("result" in message || "error" in message)) { + if (message.error) { + throw new Error(`MCP error ${message.error.code}: ${message.error.message}`); + } + return message.result as T; } - for (const message of result.values as JsonRpcMessage[]) { - if ( - "id" in message && - (message as JsonRpcResponse).id === expectedId && - ("result" in message || "error" in message) - ) { - const response = message as JsonRpcResponse; - if (response.error) { - throw new Error(`MCP error ${response.error.code}: ${response.error.message}`); - } - return response.result as T; - } - - if ("method" in message && !("id" in message)) { - this.onNotification?.(message.method, message.params); - } + if ("method" in message && !("id" in message)) { + this.onNotification?.(message.method, message.params); } } - throw new Error(`No response received for request ID ${expectedId}`); }; diff --git a/packages/coding-agent/src/modes/controllers/input-controller.ts b/packages/coding-agent/src/modes/controllers/input-controller.ts index 757f504a6..7ff3884d6 100644 --- a/packages/coding-agent/src/modes/controllers/input-controller.ts +++ b/packages/coding-agent/src/modes/controllers/input-controller.ts @@ -599,7 +599,7 @@ export class InputController { try { const image = await readImageFromClipboard(); if (image) { - const base64Data = Buffer.from(image.data).toString("base64"); + const base64Data = image.data.toBase64(); let imageData = { data: base64Data, mimeType: image.mimeType }; if (settings.get("images.autoResize")) { try { diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index d0d140f2a..57407d410 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -5,7 +5,7 @@ */ import type { AgentEvent, AgentMessage, ThinkingLevel } from "@oh-my-pi/pi-agent-core"; import type { ImageContent } from "@oh-my-pi/pi-ai"; -import { createSanitizerStream, createSplitterStream, createTextDecoderStream, ptree } from "@oh-my-pi/pi-utils"; +import { createTextLineSplitter, ptree } from "@oh-my-pi/pi-utils"; import type { BashResult } from "../../exec/bash-executor"; import type { SessionStats } from "../../session/agent-session"; import type { CompactionResult } from "../../session/compaction"; @@ -119,10 +119,7 @@ export class RpcClient { }); // Process lines in background - const lines = this.process.stdout - .pipeThrough(createTextDecoderStream()) - .pipeThrough(createSanitizerStream()) - .pipeThrough(createSplitterStream("\n")); + const lines = this.process.stdout.pipeThrough(createTextLineSplitter(true)); this.lineReader = lines; void (async () => { try { diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index 0af31fe74..aae69e683 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -10,7 +10,7 @@ * - Events: AgentSessionEvent objects streamed as they occur * - Extension UI: Extension UI requests are emitted, client responds with extension_ui_response */ -import { readLines, Snowflake } from "@oh-my-pi/pi-utils"; +import { createTextLineSplitter, Snowflake } from "@oh-my-pi/pi-utils"; import type { ExtensionUIContext, ExtensionUIDialogOptions } from "../../extensibility/extensions"; import { type Theme, theme } from "../../modes/theme/theme"; import type { AgentSession } from "../../session/agent-session"; @@ -633,7 +633,7 @@ export async function runRpcMode(session: AgentSession): Promise { } // Listen for JSON input using Bun's stdin - for await (const line of readLines(Bun.stdin.stream())) { + for await (const line of Bun.stdin.stream().pipeThrough(createTextLineSplitter())) { if (!line.trim()) continue; const result = Bun.JSONL.parseChunk(`${line}\n`); diff --git a/packages/coding-agent/src/session/auth-storage.ts b/packages/coding-agent/src/session/auth-storage.ts index 891524b11..88e3dbdc8 100644 --- a/packages/coding-agent/src/session/auth-storage.ts +++ b/packages/coding-agent/src/session/auth-storage.ts @@ -2,7 +2,6 @@ * Credential storage for API keys and OAuth tokens. * Handles loading, saving, and refreshing credentials from agent.db. */ -import { Buffer } from "node:buffer"; import * as path from "node:path"; import { antigravityUsageProvider, @@ -386,9 +385,10 @@ export class AuthStorage { const parts = token.split("."); if (parts.length !== 3) return undefined; const payloadRaw = parts[1]; + const decoder = new TextDecoder("utf-8"); try { const payload = JSON.parse( - Buffer.from(payloadRaw.replace(/-/g, "+").replace(/_/g, "/"), "base64").toString("utf8"), + decoder.decode(Uint8Array.fromBase64(payloadRaw, { alphabet: "base64url" })), ) as Record; if (!payload || typeof payload !== "object") return undefined; const identifiers: string[] = []; diff --git a/packages/coding-agent/src/session/messages.ts b/packages/coding-agent/src/session/messages.ts index 9e3698932..e08489f3a 100644 --- a/packages/coding-agent/src/session/messages.ts +++ b/packages/coding-agent/src/session/messages.ts @@ -275,23 +275,23 @@ export function convertToLlm(messages: AgentMessage[]): Message[] { timestamp: m.timestamp, }; case "fileMention": { - const fileContents = m.files - .map(file => { - const inner = file.content ? `\n${file.content}\n` : "\n"; - return `${inner}`; - }) - .join("\n\n"); - const content: (TextContent | ImageContent)[] = [ - { type: "text" as const, text: `\n${fileContents}\n` }, - ]; - for (const file of m.files) { - if (file.image) { - content.push(file.image); - } - } + const fileContents = m.files + .map(file => { + const inner = file.content ? `\n${file.content}\n` : "\n"; + return `${inner}`; + }) + .join("\n\n"); + const content: (TextContent | ImageContent)[] = [ + { type: "text" as const, text: `\n${fileContents}\n` }, + ]; + for (const file of m.files) { + if (file.image) { + content.push(file.image); + } + } return { role: "user", - content, + content, timestamp: m.timestamp, }; } diff --git a/packages/coding-agent/src/session/session-manager.ts b/packages/coding-agent/src/session/session-manager.ts index b349c688c..78321c0a6 100644 --- a/packages/coding-agent/src/session/session-manager.ts +++ b/packages/coding-agent/src/session/session-manager.ts @@ -279,26 +279,30 @@ function migrateToCurrentVersion(entries: FileEntry[]): boolean { return true; } -function parseJsonlEntries(content: string): T[] { - if (!content.trim()) return []; - const entries: T[] = []; - let buffer = content; +function parseJsonlEntries(buffer: string): T[] { + let entries: T[] | undefined; + while (buffer.length > 0) { - const result = Bun.JSONL.parseChunk(buffer); - if (result.values.length > 0) { - entries.push(...(result.values as T[])); + const { values, error, read, done } = Bun.JSONL.parseChunk(buffer); + if (values.length > 0) { + const ext = values as T[]; + if (!entries) { + entries = ext; + } else { + entries.push(...ext); + } } - if (result.error) { - const nextNewline = buffer.indexOf("\n", result.read); + if (error) { + const nextNewline = buffer.indexOf("\n", read); if (nextNewline === -1) break; - buffer = buffer.slice(nextNewline + 1); + buffer = buffer.substring(nextNewline + 1); continue; } - if (result.read === 0) break; - buffer = buffer.slice(result.read); - if (result.done) break; + if (read === 0) break; + buffer = buffer.substring(read); + if (done) break; } - return entries; + return entries ?? []; } /** Exported for testing */ diff --git a/packages/coding-agent/src/session/streaming-output.ts b/packages/coding-agent/src/session/streaming-output.ts index 6ca9330bc..62d15f103 100644 --- a/packages/coding-agent/src/session/streaming-output.ts +++ b/packages/coding-agent/src/session/streaming-output.ts @@ -140,7 +140,7 @@ export class OutputSink { await this.push(dec.decode()); }; - return new WritableStream({ + return new WritableStream({ write: async chunk => { if (typeof chunk === "string") { await this.push(chunk); diff --git a/packages/coding-agent/src/ssh/ssh-executor.ts b/packages/coding-agent/src/ssh/ssh-executor.ts index 4cdfc9f42..17c48f210 100644 --- a/packages/coding-agent/src/ssh/ssh-executor.ts +++ b/packages/coding-agent/src/ssh/ssh-executor.ts @@ -79,6 +79,7 @@ export async function executeSSH( using child = ptree.spawn(["ssh", ...(await buildRemoteCommand(host, resolvedCommand))], { signal: options?.signal, timeout: options?.timeout, + exposeStderr: true, }); const sink = new OutputSink({ @@ -87,9 +88,11 @@ export async function executeSSH( artifactId: options?.artifactId, }); - await Promise.allSettled([child.stdout.pipeTo(sink.createInput()), child.stderr.pipeTo(sink.createInput())]).catch( - () => {}, - ); + const streams = [child.stdout.pipeTo(sink.createInput())]; + if (child.stderr) { + streams.push(child.stderr.pipeTo(sink.createInput())); + } + await Promise.allSettled(streams).catch(() => {}); try { return { diff --git a/packages/coding-agent/src/tools/browser.ts b/packages/coding-agent/src/tools/browser.ts index 2eb9c8323..275d9bc8a 100644 --- a/packages/coding-agent/src/tools/browser.ts +++ b/packages/coding-agent/src/tools/browser.ts @@ -1031,10 +1031,10 @@ export class BrowserTool implements AgentTool page.content())) as string; const url = page.url(); const virtualConsole = new VirtualConsole(); - virtualConsole.on("jsdomError", error => { - if (error?.message?.includes("Could not parse CSS stylesheet")) return; + virtualConsole.on("jsdomError", err => { + if (err?.message?.includes("Could not parse CSS stylesheet")) return; logger.debug("JSDOM error during readable extraction", { - error: error instanceof Error ? error.message : String(error), + error: err instanceof Error ? err.message : String(err), }); }); const dom = new JSDOM(html, { url, virtualConsole }); @@ -1092,7 +1092,7 @@ export class BrowserTool implements AgentTool { @@ -585,68 +585,37 @@ interface AntigravitySseResult { usage?: GeminiUsageMetadata; } +const _prefix = Buffer.from("data: ", "utf-8"); + async function parseAntigravitySseForImage(response: Response, signal?: AbortSignal): Promise { if (!response.body) { throw new Error("No response body"); } - const reader = response.body.getReader(); - const decoder = new TextDecoder(); - let buffer = ""; const textParts: string[] = []; const images: InlineImageData[] = []; let usage: GeminiUsageMetadata | undefined; - try { - while (true) { - if (signal?.aborted) { - throw new Error("Request was aborted"); - } - - const { done, value } = await reader.read(); - if (done) break; - - buffer += decoder.decode(value, { stream: true }); - const lines = buffer.split("\n"); - buffer = lines.pop() ?? ""; - - for (const line of lines) { - if (!line.startsWith("data:")) continue; - const jsonStr = line.slice(5).trim(); - if (!jsonStr) continue; - - const parsed = Bun.JSONL.parseChunk(`${jsonStr}\n`); - if (parsed.error || parsed.values.length === 0) continue; - - for (const value of parsed.values) { - const chunk = value as AntigravityResponseChunk; - const responseData = chunk.response; - if (!responseData?.candidates) continue; - - if (responseData.usageMetadata) { - usage = responseData.usageMetadata; - } - - for (const candidate of responseData.candidates) { - const parts = candidate.content?.parts; - if (!parts) continue; - for (const part of parts) { - if (part.text) { - textParts.push(part.text); - } - if (part.inlineData?.data && part.inlineData?.mimeType) { - images.push({ - data: part.inlineData.data, - mimeType: part.inlineData.mimeType, - }); - } - } - } + for await (const chunk of readSseJson(response.body, signal)) { + const responseData = chunk.response; + if (!responseData) continue; + if (!responseData.candidates) continue; + for (const candidate of responseData.candidates) { + const parts = candidate.content?.parts; + if (!parts) continue; + for (const part of parts) { + if (part.text) { + textParts.push(part.text); + } + const inlineData = part.inlineData; + if (inlineData?.data && inlineData.mimeType) { + images.push({ data: inlineData.data, mimeType: inlineData.mimeType }); } } } - } finally { - reader.releaseLock(); + if (responseData.usageMetadata) { + usage = responseData.usageMetadata; + } } return { images, text: textParts, usage }; diff --git a/packages/coding-agent/src/tools/read.ts b/packages/coding-agent/src/tools/read.ts index f363b1e8d..17a6a8e6f 100644 --- a/packages/coding-agent/src/tools/read.ts +++ b/packages/coding-agent/src/tools/read.ts @@ -631,7 +631,7 @@ export class ReadTool implements AgentTool { const maxStr = formatSize(MAX_IMAGE_SIZE); throw new ToolError(`Image file too large: ${sizeStr} exceeds ${maxStr} limit.`); } else { - const base64 = Buffer.from(buffer).toString("base64"); + const base64 = new Uint8Array(buffer).toBase64(); if (this.autoResizeImages) { // Resize image if needed - catch errors from Photon diff --git a/packages/coding-agent/src/utils/file-mentions.ts b/packages/coding-agent/src/utils/file-mentions.ts index 9474aa97a..4f3559ae3 100644 --- a/packages/coding-agent/src/utils/file-mentions.ts +++ b/packages/coding-agent/src/utils/file-mentions.ts @@ -188,7 +188,7 @@ export async function generateFileMentionMessages( continue; } - const base64Content = buffer.toString("base64"); + const base64Content = buffer.toBase64(); let image = { type: "image" as const, mimeType, data: base64Content }; let dimensionNote: string | undefined; diff --git a/packages/coding-agent/src/utils/image-convert.ts b/packages/coding-agent/src/utils/image-convert.ts index e8eb83615..2c32d2fad 100644 --- a/packages/coding-agent/src/utils/image-convert.ts +++ b/packages/coding-agent/src/utils/image-convert.ts @@ -17,7 +17,7 @@ export async function convertToPng( const image = await PhotonImage.parse(new Uint8Array(Buffer.from(base64Data, "base64"))); const pngBuffer = await image.encode(ImageFormat.PNG, 100); return { - data: Buffer.from(pngBuffer).toString("base64"), + data: Buffer.from(pngBuffer).toBase64(), mimeType: "image/png", }; } catch { diff --git a/packages/coding-agent/src/utils/image-resize.ts b/packages/coding-agent/src/utils/image-resize.ts index 19272164e..437ca6225 100644 --- a/packages/coding-agent/src/utils/image-resize.ts +++ b/packages/coding-agent/src/utils/image-resize.ts @@ -119,7 +119,7 @@ export async function resizeImage(img: ImageContent, options?: ImageResizeOption if (best.buffer.length <= opts.maxBytes) { return { - data: Buffer.from(best.buffer).toString("base64"), + data: best.buffer.toBase64(), mimeType: best.mimeType, originalWidth, originalHeight, @@ -135,7 +135,7 @@ export async function resizeImage(img: ImageContent, options?: ImageResizeOption if (best.buffer.length <= opts.maxBytes) { return { - data: Buffer.from(best.buffer).toString("base64"), + data: best.buffer.toBase64(), mimeType: best.mimeType, originalWidth, originalHeight, @@ -160,7 +160,7 @@ export async function resizeImage(img: ImageContent, options?: ImageResizeOption if (best.buffer.length <= opts.maxBytes) { return { - data: Buffer.from(best.buffer).toString("base64"), + data: best.buffer.toBase64(), mimeType: best.mimeType, originalWidth, originalHeight, @@ -174,7 +174,7 @@ export async function resizeImage(img: ImageContent, options?: ImageResizeOption // Last resort: return smallest version we produced return { - data: Buffer.from(best.buffer).toString("base64"), + data: best.buffer.toBase64(), mimeType: best.mimeType, originalWidth, originalHeight, diff --git a/packages/coding-agent/src/web/search/providers/codex.ts b/packages/coding-agent/src/web/search/providers/codex.ts index ada748293..f1bd6de8d 100644 --- a/packages/coding-agent/src/web/search/providers/codex.ts +++ b/packages/coding-agent/src/web/search/providers/codex.ts @@ -6,7 +6,7 @@ * Returns synthesized answers with web search sources. */ import * as os from "node:os"; -import { readSseData } from "@oh-my-pi/pi-utils"; +import { readSseJson } from "@oh-my-pi/pi-utils"; import packageJson from "../../../../package.json" with { type: "json" }; import { getAgentDbPath, getConfigDirPaths } from "../../../config"; import { AgentStorage } from "../../../session/agent-storage"; @@ -21,6 +21,7 @@ const DEFAULT_INSTRUCTIONS = "You are a helpful assistant with web search capabilities. Search the web to answer the user's question accurately and cite your sources."; export interface CodexSearchParams { + signal?: AbortSignal; query: string; system_prompt?: string; num_results?: number; @@ -178,7 +179,7 @@ function buildCodexHeaders(accessToken: string, accountId: string): Record>(response.body)) { + for await (const rawEvent of readSseJson>(response.body, options.signal)) { const eventType = typeof rawEvent.type === "string" ? rawEvent.type : ""; if (!eventType) continue; diff --git a/packages/stats/src/parser.ts b/packages/stats/src/parser.ts index 3a0ea2e5f..e09faa559 100644 --- a/packages/stats/src/parser.ts +++ b/packages/stats/src/parser.ts @@ -64,24 +64,21 @@ export async function parseSessionFile( } const text = await file.text(); - const entries = Bun.JSONL.parse(text) as SessionEntry[]; const lines = text.split("\n"); const folder = extractFolderFromPath(sessionPath); const stats: MessageStats[] = []; let currentOffset = 0; - let entryIndex = 0; for (const line of lines) { const lineLength = line.length + 1; // +1 for newline if (line.trim()) { - const entry = entries[entryIndex]; + const entry = JSON.parse(line) as SessionEntry; if (currentOffset >= fromOffset && entry && isAssistantMessage(entry)) { const msgStats = extractStats(sessionPath, folder, entry); if (msgStats) { stats.push(msgStats); } } - entryIndex += 1; } currentOffset += lineLength; } @@ -142,7 +139,7 @@ export async function getSessionEntry(sessionPath: string, entryId: string): Pro const file = Bun.file(sessionPath); if (!(await file.exists())) return null; - const text = await file.text(); + const text = await file.bytes(); const entries = Bun.JSONL.parse(text) as SessionEntry[]; for (const entry of entries) { diff --git a/packages/tui/src/mermaid.ts b/packages/tui/src/mermaid.ts index ad28e855c..db8ec0de3 100644 --- a/packages/tui/src/mermaid.ts +++ b/packages/tui/src/mermaid.ts @@ -63,8 +63,8 @@ export async function renderMermaidToPng( return null; } - const buffer = Buffer.from(await outputFile.bytes()); - const base64 = buffer.toString("base64"); + const buffer = await outputFile.bytes(); + const base64 = buffer.toBase64(); const dims = parsePngDimensions(buffer); if (!dims) { @@ -83,14 +83,15 @@ export async function renderMermaidToPng( } } -function parsePngDimensions(buffer: Buffer): { width: number; height: number } | null { +function parsePngDimensions(buffer: Uint8Array): { width: number; height: number } | null { if (buffer.length < 24) return null; if (buffer[0] !== 0x89 || buffer[1] !== 0x50 || buffer[2] !== 0x4e || buffer[3] !== 0x47) { return null; } + const view = new DataView(buffer.buffer, buffer.byteOffset, buffer.byteLength); return { - width: buffer.readUInt32BE(16), - height: buffer.readUInt32BE(20), + width: view.getUint32(16, false), + height: view.getUint32(20, false), }; } diff --git a/packages/tui/src/terminal-capabilities.ts b/packages/tui/src/terminal-capabilities.ts index da5f764c4..5f8c9c53e 100644 --- a/packages/tui/src/terminal-capabilities.ts +++ b/packages/tui/src/terminal-capabilities.ts @@ -192,7 +192,7 @@ export function encodeITerm2( if (options.width !== undefined) params.push(`width=${options.width}`); if (options.height !== undefined) params.push(`height=${options.height}`); if (options.name) { - const nameBase64 = Buffer.from(options.name).toString("base64"); + const nameBase64 = Buffer.from(options.name).toBase64(); params.push(`name=${nameBase64}`); } if (options.preserveAspectRatio === false) { diff --git a/packages/tui/test/image-test.ts b/packages/tui/test/image-test.ts index b3d77778c..d35014d24 100644 --- a/packages/tui/test/image-test.ts +++ b/packages/tui/test/image-test.ts @@ -20,7 +20,7 @@ try { process.exit(1); } -const base64Data = Buffer.from(imageBuffer).toString("base64"); +const base64Data = imageBuffer.toBase64(); const dims = getImageDimensions(base64Data, "image/png"); console.log("Image dimensions:", dims); diff --git a/packages/utils/src/ptree.ts b/packages/utils/src/ptree.ts index ee9664a51..3f8a99f21 100644 --- a/packages/utils/src/ptree.ts +++ b/packages/utils/src/ptree.ts @@ -1,125 +1,24 @@ /** * Process tree management utilities for Bun subprocesses. * - * Exposes the same public interface as the original implementation, but with - * much less code: * - Track managed child processes for cleanup on shutdown (postmortem). * - Drain stdout/stderr to avoid subprocess pipe deadlocks. * - Cross-platform tree kill for process groups (Windows taskkill, Unix -pid). * - Convenience helpers: captureText / execText, AbortSignal, timeouts. */ - import type { Spawn, Subprocess } from "bun"; import { terminate } from "./procmgr"; +type InMask = "pipe" | "ignore" | Buffer | Uint8Array | null; + /** A Bun subprocess with stdout/stderr always piped (stdin may vary). */ type PipedSubprocess = Subprocess; -/** Minimal push-based ReadableStream that buffers unboundedly (like the old queue). */ -function pushStream() { - let controller!: ReadableStreamDefaultController; - let closed = false; - - const stream = new ReadableStream({ - start(c) { - controller = c; - }, - cancel() { - closed = true; // consumer no longer cares; keep draining but drop - }, - }); - - return { - stream, - push(value: T) { - if (closed) return; - try { - controller.enqueue(value); - } catch { - closed = true; - } - }, - close() { - if (closed) return; - closed = true; - try { - controller.close(); - } catch {} - }, - }; -} - -const DONE = { done: true, value: undefined } as const; - -function abortRead(signal: AbortSignal) { - if (signal.aborted) return Promise.resolve(DONE); - const { promise, resolve } = Promise.withResolvers(); - signal.addEventListener("abort", () => resolve(DONE), { once: true }); - return promise; -} - -/** Drain a ReadableStream into a pushStream, optionally tapping each chunk. */ -async function pump( - src: ReadableStream, - dst: ReturnType>, - opts?: { signal?: AbortSignal; onChunk?: (chunk: Uint8Array) => void; onFinally?: () => void }, -) { - const reader = src.getReader(); - const stop = opts?.signal ? abortRead(opts.signal) : null; - - try { - while (true) { - const r = stop ? await Promise.race([reader.read(), stop]) : await reader.read(); - if (r.done) break; - if (!r.value) continue; - opts?.onChunk?.(r.value); - dst.push(r.value); - } - } catch { - // ignore; this module is "best effort" for streaming/cleanup - } finally { - try { - await reader.cancel(); - } catch {} - try { - reader.releaseLock(); - } catch {} - dst.close(); - opts?.onFinally?.(); - } -} +// ── Exceptions ─────────────────────────────────────────────────────────────── /** - * Kill a child process and its descendents. - * - Windows: taskkill /T, add /F on SIGKILL - * - Unix: negative PID signals the process group - */ -async function killChild(child: ChildProcess) { - await terminate({ target: child.proc }); -} - -/** - * Options for waiting for process exit and capturing output. - */ -export interface WaitOptions { - allowNonZero?: boolean; - allowAbort?: boolean; - stderr?: "full" | "buffer"; -} - -/** - * Result from wait and captureText. - */ -export interface ExecResult { - stdout: string; - stderr: string; - exitCode: number | null; - ok: boolean; - exitError?: Exception; -} - -/** - * Base for all exceptions representing child process nonzero exit, killed, or cancellation. + * Base for all exceptions representing child process nonzero exit, killed, or + * cancellation. */ export abstract class Exception extends Error { constructor( @@ -133,119 +32,114 @@ export abstract class Exception extends Error { abstract get aborted(): boolean; } -/** - * Exception for nonzero exit codes (not cancellation). - */ +/** Exception for nonzero exit codes (not cancellation). */ export class NonZeroExitError extends Exception { static readonly MAX_TRACE = 32 * 1024; - constructor( - public readonly exitCode: number, - public readonly stderr: string, - ) { + constructor(exitCode: number, stderr: string) { super(`Process exited with code ${exitCode}:\n${stderr}`, exitCode, stderr); } - get aborted(): boolean { + get aborted() { return false; } } -/** - * Exception for explicit process abortion (via signal). - */ +/** Exception for explicit process abortion (via signal). */ export class AbortError extends Exception { constructor( public readonly reason: unknown, stderr: string, ) { - const reasonString = reason instanceof Error ? reason.message : String(reason ?? "aborted"); - super(`Operation cancelled: ${reasonString}`, -1, stderr); + const msg = reason instanceof Error ? reason.message : String(reason ?? "aborted"); + super(`Operation cancelled: ${msg}`, -1, stderr); } - get aborted(): boolean { + get aborted() { return true; } } -/** - * Exception for process timeout. - */ +/** Exception for process timeout. */ export class TimeoutError extends AbortError { constructor(timeout: number, stderr: string) { super(new Error(`Timed out after ${Math.round(timeout / 1000)}s`), stderr); } } -type InMask = "pipe" | "ignore" | Buffer | Uint8Array | null; +// ── Wait / Exec types ──────────────────────────────────────────────────────── + +/** Options for waiting for process exit and capturing output. */ +export interface WaitOptions { + allowNonZero?: boolean; + allowAbort?: boolean; + stderr?: "full" | "buffer"; +} + +/** Result from wait and exec. */ +export interface ExecResult { + stdout: string; + stderr: string; + exitCode: number | null; + ok: boolean; + exitError?: Exception; +} + +// ── ChildProcess ───────────────────────────────────────────────────────────── /** * ChildProcess wraps a managed subprocess, capturing stderr tail, providing * cross-platform kill/detach logic plus AbortSignal integration. + * + * Stdout is exposed directly from the underlying Bun subprocess; consumers + * must read it (via text(), wait(), etc.) to prevent pipe deadlock. + * Stderr is eagerly drained into an internal buffer. */ export class ChildProcess { #nothrow = false; - - #stderrBuffer = ""; + #stderrTail = ""; + #stderrChunks: Uint8Array[] = []; #exitReason?: Exception; #exitReasonPending?: Exception; - - #stop = new AbortController(); - - #stdoutOut = pushStream(); - #stderrOut = pushStream(); - #stderrDone: Promise; #exited: Promise; + #stderrStream?: ReadableStream; - constructor(public readonly proc: PipedSubprocess) { - const { promise: stderrDone, resolve: resolveStderrDone } = Promise.withResolvers(); - this.#stderrDone = stderrDone; - - // Drain stdout always -> expose our buffered stream to the user. - void pump(proc.stdout, this.#stdoutOut, { signal: this.#stop.signal }).catch(() => this.#stdoutOut.close()); - - // Drain stderr always -> expose stream + keep a bounded tail buffer. - const decoder = new TextDecoder(); + constructor(public readonly proc: PipedSubprocess, options?: { exposeStderr?: boolean }) { + // Eagerly drain stderr into a truncated tail string + raw chunks. + const dec = new TextDecoder(); const trim = () => { - if (this.#stderrBuffer.length > NonZeroExitError.MAX_TRACE) { - this.#stderrBuffer = this.#stderrBuffer.slice(-NonZeroExitError.MAX_TRACE); - } + if (this.#stderrTail.length > NonZeroExitError.MAX_TRACE) + this.#stderrTail = this.#stderrTail.slice(-NonZeroExitError.MAX_TRACE); }; - void pump(proc.stderr, this.#stderrOut, { - signal: this.#stop.signal, - onChunk: chunk => { - this.#stderrBuffer += decoder.decode(chunk, { stream: true }); - trim(); - }, - onFinally: () => { - this.#stderrBuffer += decoder.decode(); - trim(); - resolveStderrDone(); - }, - }).catch(() => { + let stderrStream = proc.stderr; + if (options?.exposeStderr) { + const [teeStream, drainStream] = stderrStream.tee(); + this.#stderrStream = teeStream; + stderrStream = drainStream; + } + this.#stderrDone = (async () => { try { - this.#stderrBuffer += decoder.decode(); - trim(); + for await (const chunk of stderrStream) { + this.#stderrChunks.push(chunk); + this.#stderrTail += dec.decode(chunk, { stream: true }); + trim(); + } } catch {} - this.#stderrOut.close(); - resolveStderrDone(); - }); + this.#stderrTail += dec.decode(); + trim(); + })(); + // Normalize Bun's exited promise into our exitReason / exitedCleanly model. const { promise, resolve, reject } = Promise.withResolvers(); this.#exited = promise; - // Normalize Bun's exited promise into our "exitReason / exitedCleanly" model. proc.exited .catch(() => null) .then(async exitCode => { - // Stop pumping streams - process has exited, no more data coming - this.#stop.abort(); - if (this.#exitReasonPending) { this.#exitReason = this.#exitReasonPending; reject(this.#exitReasonPending); return; } - if (exitCode === 0) { resolve(0); return; @@ -254,54 +148,61 @@ export class ChildProcess { await this.#stderrDone; if (exitCode !== null) { - this.#exitReason = new NonZeroExitError(exitCode, this.#stderrBuffer); + this.#exitReason = new NonZeroExitError(exitCode, this.#stderrTail); resolve(exitCode); return; } const ex = this.proc.killed - ? new AbortError(new Error("process killed"), this.#stderrBuffer) - : new NonZeroExitError(-1, this.#stderrBuffer); - + ? new AbortError(new Error("process killed"), this.#stderrTail) + : new NonZeroExitError(-1, this.#stderrTail); this.#exitReason = ex; reject(ex); }); } - get pid(): number | undefined { + // ── Properties ─────────────────────────────────────────────────────── + + get pid() { return this.proc.pid; } - get exited(): Promise { + get exited() { 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 { + get exitCode() { return this.proc.exitCode; } - get exitReason(): Exception | undefined { + get exitReason() { return this.#exitReason; } - get killed(): boolean { + get killed() { return this.proc.killed; } get stdin(): Bun.SpawnOptions.WritableToIO { return this.proc.stdin; } - get stdout(): ReadableStream { - return this.#stdoutOut.stream; - } - get stderr(): ReadableStream { - return this.#stderrOut.stream; + + /** Raw stdout stream. Must be consumed to prevent pipe deadlock. */ + get stdout() { + return this.proc.stdout; } - peekStderr(): string { - return this.#stderrBuffer; + /** Optional stderr stream (only when requested in spawn options). */ + get stderr() { + return this.#stderrStream; + } + + get exitedCleanly(): Promise { + if (this.#nothrow) return this.#exited; + return this.#exited.then(code => { + if (code !== 0) throw new NonZeroExitError(code, this.#stderrTail); + return code; + }); + } + + /** Returns the truncated stderr tail (last 32KB). */ + peekStderr() { + return this.#stderrTail; } nothrow(): this { @@ -311,48 +212,53 @@ export class ChildProcess { kill(reason?: Exception) { if (reason && !this.#exitReasonPending) this.#exitReasonPending = reason; - this.#stop.abort(); - if (this.proc.killed) return; - void killChild(this); + if (!this.proc.killed) void terminate({ target: this.proc }); + } + + // ── Output helpers ─────────────────────────────────────────────────── + + async text(): Promise { + const p = new Response(this.stdout).text(); + if (this.#nothrow) return p; + const [text] = await Promise.all([p, this.exitedCleanly]); + return text; } - // Output helpers async blob(): Promise { - const blobPromise = new Response(this.stdout).blob(); - if (this.#nothrow) return await blobPromise; - const [blob] = await Promise.all([blobPromise, this.exitedCleanly]); + const p = new Response(this.stdout).blob(); + if (this.#nothrow) return p; + const [blob] = await Promise.all([p, this.exitedCleanly]); return blob; } - async text(): Promise { - return (await this.blob()).text(); - } + async json(): Promise { - return await new Response(await this.blob()).json(); + return new Response(this.stdout).json(); } + async arrayBuffer(): Promise { - return (await this.blob()).arrayBuffer(); + return new Response(this.stdout).arrayBuffer(); } + async bytes(): Promise { - return new Uint8Array(await this.arrayBuffer()); + return new Response(this.stdout).bytes(); } - async wait(options?: WaitOptions): Promise { - const { allowNonZero = false, allowAbort = false, stderr: stderrMode = "buffer" } = options ?? {}; + // ── Wait ───────────────────────────────────────────────────────────── - const stdoutPromise = new Response(this.stdout).text(); - const stderrPromise = + async wait(opts?: WaitOptions): Promise { + const { allowNonZero = false, allowAbort = false, stderr: stderrMode = "buffer" } = opts ?? {}; + + const stdoutP = this.text(); + const stderrP = stderrMode === "full" - ? new Response(this.stderr).text() - : (async () => { - await Promise.allSettled([stdoutPromise, this.exited, this.#stderrDone]); - return this.peekStderr(); - })(); + ? this.#stderrDone.then(() => new TextDecoder().decode(Buffer.concat(this.#stderrChunks))) + : this.#stderrDone.then(() => this.#stderrTail); - const [stdout, stderr] = await Promise.all([stdoutPromise, stderrPromise]); + const [stdout, stderr] = await Promise.all([stdoutP, stderrP]); let exitError: Exception | undefined; try { - await this.exited; + await this.#exited; } catch (err) { if (err instanceof Exception) exitError = err; else throw err; @@ -362,69 +268,55 @@ export class ChildProcess { const ok = exitCode === 0; if (exitError) { - if ((exitError.aborted && !allowAbort) || (!exitError.aborted && !allowNonZero)) { - throw exitError; - } + if ((exitError.aborted && !allowAbort) || (!exitError.aborted && !allowNonZero)) throw exitError; } return { stdout, stderr, exitCode, ok, exitError }; } + // ── Signal / timeout ───────────────────────────────────────────────── + attachSignal(signal: AbortSignal): void { const onAbort = () => this.kill(new AbortError(signal.reason, "")); if (signal.aborted) return void onAbort(); - signal.addEventListener("abort", onAbort, { once: true }); - this.#exited - .catch(() => {}) - .finally(() => { - signal.removeEventListener("abort", onAbort); - }); + this.#exited.catch(() => {}).finally(() => signal.removeEventListener("abort", onAbort)); } - attachTimeout(timeout: number): void { - if (timeout <= 0 || this.proc.killed) return; - void (async () => { - const timedOut = await Promise.race([ - Bun.sleep(timeout).then(() => true), - this.proc.exited.then( - () => false, - () => false, - ), - ]); - if (timedOut) this.kill(new TimeoutError(timeout, this.#stderrBuffer)); - })(); + attachTimeout(ms: number): void { + if (ms <= 0 || this.proc.killed) return; + Promise.race([ + Bun.sleep(ms).then(() => true), + this.proc.exited.then( + () => false, + () => false, + ), + ]).then(timedOut => { + if (timedOut) this.kill(new TimeoutError(ms, this.#stderrTail)); + }); } [Symbol.dispose](): void { - // Don't kill if process already exited - avoids race where dispose runs - // before the proc.exited.then() callback, causing spurious AbortError if (this.proc.exitCode !== null) return; - this.kill(new AbortError("process disposed", this.#stderrBuffer)); + this.kill(new AbortError("process disposed", this.#stderrTail)); } } -/** - * Options for cspawn (child spawn). Always pipes stdout/stderr, allows signal. - */ +// ── Spawn / exec ───────────────────────────────────────────────────────────── + +/** Options for child spawn. Always pipes stdout/stderr. */ type ChildSpawnOptions = Omit< Spawn.SpawnOptions, "stdout" | "stderr" | "detached" > & { - /** AbortSignal to cancel the process */ signal?: AbortSignal; - /** Whether to detach the process */ detached?: boolean; + exposeStderr?: boolean; }; -/** - * Spawn a child process. - * @param cmd - The command to spawn. - * @param options - The options for the spawn. - * @returns A ChildProcess instance. - */ -export function spawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess { - const { timeout = -1, signal, ...rest } = options ?? {}; +/** Spawn a child process with piped stdout/stderr. */ +export function spawn(cmd: string[], opts?: ChildSpawnOptions): ChildProcess { + const { timeout = -1, signal, exposeStderr, ...rest } = opts ?? {}; const child = Bun.spawn(cmd, { stdin: "ignore", stdout: "pipe", @@ -432,55 +324,55 @@ export function spawn(cmd: string[], options?: Child windowsHide: true, ...rest, }); - const cproc = new ChildProcess(child); - if (signal) cproc.attachSignal(signal); - if (timeout > 0) cproc.attachTimeout(timeout); - return cproc; + const cp = new ChildProcess(child, { exposeStderr }); + if (signal) cp.attachSignal(signal); + if (timeout > 0) cp.attachTimeout(timeout); + return cp; } -/** - * Options for execText. - */ +/** Options for exec. */ export interface ExecOptions extends Omit, WaitOptions { input?: string | Buffer | Uint8Array; } -export async function exec(cmd: string[], options?: ExecOptions): Promise { - const { input, stderr, allowAbort, allowNonZero, ...spawnOptions } = options ?? {}; +/** Spawn, wait, and return captured output. */ +export async function exec(cmd: string[], opts?: ExecOptions): Promise { + const { input, stderr, allowAbort, allowNonZero, ...spawnOpts } = opts ?? {}; const stdin = typeof input === "string" ? Buffer.from(input) : input; - const resolvedOptions: ChildSpawnOptions = stdin === undefined ? spawnOptions : { ...spawnOptions, stdin }; - using child = spawn(cmd, resolvedOptions); - return await child.wait({ stderr, allowAbort, allowNonZero }); + const resolved: ChildSpawnOptions = stdin === undefined ? spawnOpts : { ...spawnOpts, stdin }; + using child = spawn(cmd, resolved); + return child.wait({ stderr, allowAbort, allowNonZero }); } +// ── Signal combinators ─────────────────────────────────────────────────────── + type SignalValue = AbortSignal | number | null | undefined; +/** Combine AbortSignals and timeout values into a single signal. */ export function combineSignals(...signals: SignalValue[]): AbortSignal | undefined { let timeout: number | undefined; + let n = 0; for (let i = 0; i < signals.length; i++) { const s = signals[i]; if (s instanceof AbortSignal) { if (s.aborted) return s; - signals[n++] = s; + if (i !== n) signals[n] = s; + n++; } else if (typeof s === "number" && s > 0) { - timeout = Math.min(timeout ?? s, s); + timeout = timeout === undefined ? s : Math.min(timeout, s); } } - - // Create timeout signal. if (timeout !== undefined) { - signals[n++] = AbortSignal.timeout(timeout); + signals[n] = AbortSignal.timeout(timeout); + n++; } - - // Single signal, ezpz. - const rawSignals = signals.slice(0, n) as AbortSignal[]; - switch (rawSignals.length) { + switch (n) { case 0: return undefined; case 1: - return rawSignals[0]; + return signals[0] as AbortSignal; default: - return AbortSignal.any(rawSignals); + return AbortSignal.any(signals.slice(0, n) as AbortSignal[]); } } diff --git a/packages/utils/src/stream.ts b/packages/utils/src/stream.ts index a7f2a634d..aa7479951 100644 --- a/packages/utils/src/stream.ts +++ b/packages/utils/src/stream.ts @@ -1,3 +1,5 @@ +import { ArrayBufferSink } from "bun"; + /** * Sanitize binary output for display/storage. * Removes characters that crash string-width or cause display issues: @@ -43,27 +45,57 @@ export function sanitizeText(text: string): string { /** * Create a transform stream that splits lines. */ -export function createSplitterStream(delimiter: string): TransformStream { - let buf = ""; - return new TransformStream({ - transform(chunk, controller) { - buf = buf ? `${buf}${chunk}` : chunk; +export function createSplitterStream(options: { + newLine?: boolean; + mapFn: (chunk: Uint8Array) => T; +}): TransformStream { + const { newLine = false, mapFn } = options; + const LF = 0x0a; + const sink = new Bun.ArrayBufferSink(); + sink.start({ asUint8Array: true, stream: true, highWaterMark: 4096 }); + let pending = false; // whether the sink has unflushed data - while (true) { - const nl = buf.indexOf(delimiter); - if (nl === -1) break; - controller.enqueue(buf.slice(0, nl)); - buf = buf.slice(nl + delimiter.length); + return new TransformStream({ + transform(chunk, ctrl) { + let pos = 0; + + while (pos < chunk.length) { + const nl = chunk.indexOf(LF, pos); + if (nl === -1) { + sink.write(chunk.subarray(pos)); + pending = true; + break; + } + + const slice = chunk.subarray(pos, newLine ? nl + 1 : nl); + + if (pending) { + if (slice.length > 0) sink.write(slice); + ctrl.enqueue(mapFn(sink.flush() as Uint8Array)); + pending = false; + } else { + ctrl.enqueue(mapFn(slice)); + } + pos = nl + 1; } }, - flush(controller) { - if (buf) { - controller.enqueue(buf); + flush(ctrl) { + if (pending) { + const tail = sink.end() as Uint8Array; + if (tail.length > 0) ctrl.enqueue(mapFn(tail)); } }, }); } +export function createTextLineSplitter(sanitize = false): TransformStream { + const dec = new TextDecoder("utf-8", { ignoreBOM: true, fatal: true }); + if (sanitize) { + return createSplitterStream({ mapFn: chunk => sanitizeText(dec.decode(chunk)) }); + } + return createSplitterStream({ mapFn: dec.decode.bind(dec) }); +} + /** * Create a transform stream that sanitizes text. */ @@ -82,152 +114,130 @@ export function createTextDecoderStream(): TransformStream { return new TextDecoderStream() as TransformStream; } -/** - * Read stream line-by-line - * - * @param delimiter Line delimiter (default: "\n") - */ -export function readLines(stream: ReadableStream, delimiter = "\n"): AsyncIterable { - return stream.pipeThrough(createTextDecoderStream()).pipeThrough(createSplitterStream(delimiter)); -} - // ============================================================================= // SSE (Server-Sent Events) // ============================================================================= -/** - * Parsed SSE event. - */ -export interface SseEvent { - /** Event type (from `event:` field, default: "message") */ - event: string; - /** Event data (from `data:` field(s), joined with newlines) */ - data: string; - /** Event ID (from `id:` field) */ - id?: string; - /** Retry interval in ms (from `retry:` field) */ - retry?: number; -} +const LF = 0x0a; +const CR = 0x0d; +const SPACE = 0x20; -/** - * Parse a single SSE event block (lines between blank lines). - * Returns null if the block contains no data. - */ -export function parseSseEvent(block: string): SseEvent | null { - const lines = block.split("\n"); - let event = "message"; - const dataLines: string[] = []; - let id: string | undefined; - let retry: number | undefined; +// "data:" = [0x64, 0x61, 0x74, 0x61, 0x3a] +const DATA_0 = 0x64; // d +const DATA_1 = 0x61; // a +const DATA_2 = 0x74; // t +const DATA_3 = 0x61; // a +const DATA_4 = 0x3a; // : - for (const line of lines) { - // Comments start with ':' - if (line.startsWith(":")) continue; +// "[DONE]" = [0x5b, 0x44, 0x4f, 0x4e, 0x45, 0x5d] +const DONE = Uint8Array.from([0x5b, 0x44, 0x4f, 0x4e, 0x45, 0x5d]); - const colonIdx = line.indexOf(":"); - if (colonIdx === -1) continue; - - const field = line.slice(0, colonIdx); - // Value starts after colon, with optional leading space trimmed - let value = line.slice(colonIdx + 1); - if (value.startsWith(" ")) value = value.slice(1); - - switch (field) { - case "event": - event = value; - break; - case "data": - dataLines.push(value); - break; - case "id": - id = value; - break; - case "retry": { - const n = parseInt(value, 10); - if (!Number.isNaN(n)) retry = n; - break; - } - } +function isDone(buf: Uint8Array, start: number, end: number): boolean { + if (end - start !== 6) return false; + for (let i = 0; i < 6; i++) { + if (buf[start + i] !== DONE[i]) return false; } - - if (dataLines.length === 0) return null; - - return { - event, - data: dataLines.join("\n"), - id, - retry, - }; + return true; } /** - * Read SSE events from a stream. - * - * Handles the SSE wire format: - * - Events separated by blank lines - * - Fields: event, data, id, retry - * - Comments (lines starting with :) are ignored - * - Multiple data: lines are joined with newlines + * Stream parsed JSON objects from SSE `data:` lines. * * @example * ```ts - * for await (const event of readSseEvents(response.body)) { - * if (event.data === "[DONE]") break; - * const payload = JSON.parse(event.data); - * console.log(event.event, payload); + * for await (const obj of readSseJson(response.body!)) { + * console.log(obj); * } * ``` */ -export async function* readSseEvents(stream: ReadableStream): AsyncGenerator { - const blockLines: string[] = []; - - for await (const rawLine of readLines(stream)) { - const line = rawLine.replace(/\r$/, ""); - if (line === "") { - if (blockLines.length > 0) { - const event = parseSseEvent(blockLines.join("\n")); - if (event) yield event; - blockLines.length = 0; - } - continue; - } - - blockLines.push(line); - } - - if (blockLines.length > 0) { - const event = parseSseEvent(blockLines.join("\n")); - if (event) yield event; - } -} - -/** - * Read SSE data payloads from a stream, parsing JSON automatically. - * - * Convenience wrapper over readSseEvents that: - * - Skips [DONE] markers - * - Parses JSON data - * - Optionally filters by event type - * - * @example - * ```ts - * for await (const data of readSseData(response.body)) { - * console.log(data.choices[0].delta); - * } - * ``` - */ -export async function* readSseData( +export async function* readSseJson( stream: ReadableStream, - eventType?: string, -): AsyncGenerator { - for await (const event of readSseEvents(stream)) { - if (eventType && event.event !== eventType) continue; - if (event.data === "[DONE]") continue; + abortSignal?: AbortSignal, +): AsyncGenerator { + const sink = new ArrayBufferSink(); + sink.start({ asUint8Array: true, stream: true, highWaterMark: 4096 }); + let pending = false; - try { - yield JSON.parse(event.data) as T; - } catch { - // Skip malformed JSON + const cleanup = () => { + stream.cancel(abortSignal?.reason ?? new Error("Request aborted")); + }; + abortSignal?.addEventListener("abort", cleanup, { once: true }); + + try { + for await (const chunk of stream) { + let pos = 0; + while (pos < chunk.length) { + const nl = chunk.indexOf(LF, pos); + if (nl === -1) { + sink.write(chunk.subarray(pos)); + pending = true; + break; + } + + let line: Uint8Array; + if (pending) { + if (nl > pos) sink.write(chunk.subarray(pos, nl)); + line = sink.flush() as Uint8Array; + pending = false; + } else { + line = chunk.subarray(pos, nl); + } + pos = nl + 1; + + // Strip trailing CR, skip blank/short lines. + const len = line.length > 0 && line[line.length - 1] === CR ? line.length - 1 : line.length; + if (len < 6) continue; // "data:" + at least 1 byte + + // Check "data:" prefix. + if ( + line[0] !== DATA_0 || + line[1] !== DATA_1 || + line[2] !== DATA_2 || + line[3] !== DATA_3 || + line[4] !== DATA_4 + ) + continue; + + // Payload start — skip optional space after colon. + const pStart = line[5] === SPACE ? 6 : 5; + if (pStart >= len) continue; + if (isDone(line, pStart, len)) return; + + // Build payload + \n for JSONL.parse. + const pLen = len - pStart; + const buf = new Uint8Array(pLen + 1); + buf.set(line.subarray(pStart, len)); + buf[pLen] = LF; + + const [parsed] = Bun.JSONL.parse(buf); + if (parsed !== undefined) yield parsed as T; + } + } + } finally { + abortSignal?.removeEventListener("abort", cleanup); + } + + // Trailing line without final newline. + if (pending) { + const tail = sink.end() as Uint8Array; + const len = tail.length > 0 && tail[tail.length - 1] === CR ? tail.length - 1 : tail.length; + if ( + len >= 6 && + tail[0] === DATA_0 && + tail[1] === DATA_1 && + tail[2] === DATA_2 && + tail[3] === DATA_3 && + tail[4] === DATA_4 + ) { + const pStart = tail[5] === SPACE ? 6 : 5; + if (pStart < len && !isDone(tail, pStart, len)) { + const pLen = len - pStart; + const buf = new Uint8Array(pLen + 1); + buf.set(tail.subarray(pStart, len)); + buf[pLen] = LF; + const [parsed] = Bun.JSONL.parse(buf); + if (parsed !== undefined) yield parsed as T; + } } } } diff --git a/packages/utils/test/stream.test.ts b/packages/utils/test/stream.test.ts new file mode 100644 index 000000000..7a5827ed3 --- /dev/null +++ b/packages/utils/test/stream.test.ts @@ -0,0 +1,171 @@ +import { describe, expect, it } from "bun:test"; +import { + createSanitizerStream, + createSplitterStream, + createTextDecoderStream, + createTextLineSplitter, + readSseJson, + sanitizeBinaryOutput, + sanitizeText, +} from "../src/stream"; + +const encoder = new TextEncoder(); + +async function runTransform(transform: TransformStream, chunks: Uint8Array[]): Promise { + const readable = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const reader = readable.pipeThrough(transform).getReader(); + const output: T[] = []; + while (true) { + const { value, done } = await reader.read(); + if (done) break; + output.push(value); + } + return output; +} + +async function runStringTransform(transform: TransformStream, chunks: string[]): Promise { + const readable = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const reader = readable.pipeThrough(transform).getReader(); + const output: string[] = []; + while (true) { + const { value, done } = await reader.read(); + if (done) break; + output.push(value); + } + return output; +} + +async function collectAsync(iter: AsyncIterable): Promise { + const output: T[] = []; + for await (const item of iter) output.push(item); + return output; +} + +describe("sanitizeBinaryOutput", () => { + it("removes control characters but keeps tabs/newlines", () => { + const input = "a\u0000b\tline\ncarriage\r\u0001"; + expect(sanitizeBinaryOutput(input)).toBe("ab\tline\ncarriage\r"); + }); +}); + +describe("sanitizeText", () => { + it("strips ANSI and normalizes CR", () => { + const input = "\u001b[31mred\u001b[0m\r\n"; + expect(sanitizeText(input)).toBe("red\n"); + }); +}); + +describe("createSplitterStream", () => { + it("splits lines across chunks without newlines", async () => { + const transform = createSplitterStream({ + mapFn: chunk => new TextDecoder().decode(chunk), + }); + + const output = await runTransform(transform, [ + encoder.encode("alpha\nbe"), + encoder.encode("ta\ngam"), + encoder.encode("ma"), + ]); + + expect(output).toEqual(["alpha", "beta", "gamma"]); + }); + + it("includes newlines when requested", async () => { + const transform = createSplitterStream({ + newLine: true, + mapFn: chunk => new TextDecoder().decode(chunk), + }); + + const output = await runTransform(transform, [encoder.encode("one\ntwo\n")]); + + expect(output).toEqual(["one\n", "two\n"]); + }); +}); + +describe("createTextLineSplitter", () => { + it("decodes utf-8 and sanitizes when requested", async () => { + const transform = createTextLineSplitter(true); + const output = await runTransform(transform, [ + encoder.encode("\u001b[32mgreen\u001b[0m\r\nblue\n"), + ]); + + expect(output).toEqual(["green", "blue"]); + }); +}); + +describe("createSanitizerStream", () => { + it("sanitizes text chunks", async () => { + const transform = createSanitizerStream(); + const output = await runStringTransform(transform, ["\u001b[34mhi\u001b[0m\r\n"]); + + expect(output).toEqual(["hi\n"]); + }); +}); + +describe("createTextDecoderStream", () => { + it("decodes utf-8 byte streams", async () => { + const transform = createTextDecoderStream(); + const output = await runTransform(transform, [encoder.encode("hello"), encoder.encode(" world")]); + + expect(output.join("")).toBe("hello world"); + }); +}); + +describe("readSseJson", () => { + it("parses data lines and stops at [DONE]", async () => { + const chunks = [ + encoder.encode("data: {\"a\":1}\n"), + encoder.encode("event: ping\n"), + encoder.encode("data: {\"b\":2}\r\n"), + encoder.encode("data: [DONE]\n"), + encoder.encode("data: {\"c\":3}\n"), + ]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const output = await collectAsync(readSseJson(stream)); + expect(output).toEqual([{ a: 1 }, { b: 2 }]); + }); + + it("parses trailing line without newline", async () => { + const chunks = [encoder.encode("data: {\"c\":3}")]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const output = await collectAsync(readSseJson(stream)); + expect(output).toEqual([{ c: 3 }]); + }); + + it("handles data lines split across chunks", async () => { + const chunks = [encoder.encode("data: {\"a\""), encoder.encode(":1}\n")]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const output = await collectAsync(readSseJson(stream)); + expect(output).toEqual([{ a: 1 }]); + }); +}); diff --git a/scripts/dump-edit-history.ts b/scripts/dump-edit-history.ts index 441ecefc9..c4fbd4eec 100644 --- a/scripts/dump-edit-history.ts +++ b/scripts/dump-edit-history.ts @@ -72,7 +72,7 @@ function classifyError(resultText: string): string { } async function extractEditAttempts(sessionPath: string): Promise { - const content = await Bun.file(sessionPath).text(); + const content = await Bun.file(sessionPath).bytes(); const messages = Bun.JSONL.parse(content) as Message[]; const editAttempts: EditAttempt[] = []; diff --git a/tsconfig.base.json b/tsconfig.base.json index e0b90a58b..c4d28c082 100644 --- a/tsconfig.base.json +++ b/tsconfig.base.json @@ -2,7 +2,7 @@ "compilerOptions": { "target": "ES2024", "module": "ESNext", - "lib": ["ES2024"], + "lib": ["ES2024", "DOM.AsyncIterable"], "strict": true, "esModuleInterop": true, "skipLibCheck": true,