diff --git a/Cargo.lock b/Cargo.lock index 8e51d1ea9..d0d9a384e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1989,7 +1989,7 @@ dependencies = [ [[package]] name = "pi-natives" -version = "11.2.1" +version = "11.2.2" dependencies = [ "arboard", "brush-builtins", diff --git a/bun.lock b/bun.lock index fec22beb4..c9399c312 100644 --- a/bun.lock +++ b/bun.lock @@ -67,8 +67,8 @@ }, "dependencies": { "@mozilla/readability": "0.6.0", - "@oclif/core": "^4.5.6", - "@oclif/plugin-autocomplete": "^3.2.23", + "@oclif/core": "^4.8.0", + "@oclif/plugin-autocomplete": "^3.2.40", "@oh-my-pi/omp-stats": "workspace:*", "@oh-my-pi/pi-agent-core": "workspace:*", "@oh-my-pi/pi-ai": "workspace:*", @@ -86,7 +86,7 @@ "jsdom": "28.0.0", "marked": "^17.0.1", "node-html-parser": "^7.0.2", - "puppeteer": "^24.36.1", + "puppeteer": "^24.37.1", "smol-toml": "^1.6.0", "zod": "^4.3.6", }, @@ -915,7 +915,7 @@ "onetime": ["onetime@7.0.0", "", { "dependencies": { "mimic-function": "^5.0.0" } }, "sha512-VXJjc87FScF88uafS3JllDgvAm+c/Slfz06lorj2uAY34rlUu0Nt+v8wreiImcrgAjjIHp1rXpTDlLOGw29WwQ=="], - "openai": ["openai@6.17.0", "", { "peerDependencies": { "ws": "^8.18.0", "zod": "^3.25 || ^4.0" }, "optionalPeers": ["ws", "zod"], "bin": { "openai": "bin/cli" } }, "sha512-NHRpPEUPzAvFOAFs9+9pC6+HCw/iWsYsKCMPXH5Kw7BpMxqd8g/A07/1o7Gx2TWtCnzevVRyKMRFqyiHyAlqcA=="], + "openai": ["openai@6.18.0", "", { "peerDependencies": { "ws": "^8.18.0", "zod": "^3.25 || ^4.0" }, "optionalPeers": ["ws", "zod"], "bin": { "openai": "bin/cli" } }, "sha512-odLRYyz9rlzz6g8gKn61RM2oP5UUm428sE2zOxZqS9MzVfD5/XW8UoEjpnRkzTuScXP7ZbP/m7fC+bl8jCOZZw=="], "pac-proxy-agent": ["pac-proxy-agent@7.2.0", "", { "dependencies": { "@tootallnate/quickjs-emscripten": "^0.23.0", "agent-base": "^7.1.2", "debug": "^4.3.4", "get-uri": "^6.0.1", "http-proxy-agent": "^7.0.0", "https-proxy-agent": "^7.0.6", "pac-resolver": "^7.0.1", "socks-proxy-agent": "^8.0.5" } }, "sha512-TEB8ESquiLMc0lV8vcd5Ql/JAKAoyzHFXaStwjkzpOpC5Yv+pIzLfHvjTSdf3vpa2bMiUQrg9i6276yn8666aA=="], @@ -1003,7 +1003,7 @@ "scheduler": ["scheduler@0.27.0", "", {}, "sha512-eNv+WrVbKu1f3vbYJT/xtiF5syA5HPIMtf9IgY/nKg0sWqzAUEvqY/xm7OcZc/qafLx/iO9FgOmeSAp4v5ti/Q=="], - "semver": ["semver@7.7.3", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-SdsKMrI9TdgjdweUSR9MweHA4EJ8YxHn8DFaDisvhVlUOe4BF1tLD7GAj0lIqWVl+dPb/rExr0Btby5loQm20Q=="], + "semver": ["semver@7.7.4", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-vFKC2IEtQnVhpT78h1Yp8wzwrf8CM+MzKMHGJZfBtzhZNycRFnXsHk6E5TxIkkMsgNS7mdX3AGB7x2QM2di4lA=="], "shebang-command": ["shebang-command@2.0.0", "", { "dependencies": { "shebang-regex": "^3.0.0" } }, "sha512-kHxr2zZpYtdmrN1qDjrrX/Z1rR1kG8Dx+gkpK1G4eXmvXswmcE1hTWBWYUzlraYw1/yZp6YuDY77YtvbN0dmDA=="], diff --git a/package.json b/package.json index 2ce4d0794..9a8da7621 100644 --- a/package.json +++ b/package.json @@ -52,7 +52,6 @@ "engines": { "bun": ">=1.3.7" }, - "packageManager": "bun@1.3.7", "version": "0.0.3", "dependencies": { "@sinclair/typebox": "^0.34.48" diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 699827f40..2abfe6bcd 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -1,6 +1,22 @@ # Changelog ## [Unreleased] +### Added + +- Added Claude Opus 4.6 model support across multiple providers (Anthropic, Amazon Bedrock, GitHub Copilot, OpenRouter, OpenCode, Vercel AI Gateway) +- Added GPT-5.3 Codex model support for OpenAI +- Added `readSseJson` utility import for improved SSE stream handling in Google Gemini CLI provider + +### Changed + +- Updated Google Gemini CLI provider to use `readSseJson` utility for cleaner SSE stream parsing +- Updated pricing for Llama 3.1 405B model on Vercel AI Gateway (cache read rate adjusted) +- Updated Llama 3.1 405B context window and max tokens on Vercel AI Gateway (256000 for both) + +### Removed + +- Removed Kimi K2, Kimi K2 Turbo Preview, and Kimi K2.5 models +- Removed Deep Cogito Cogito V2 Preview models from OpenRouter ## [11.0.0] - 2026-02-05 @@ -765,4 +781,4 @@ _Dedicated to Peter's shoulder ([@steipete](https://twitter.com/steipete))_ ## [0.9.4] - 2025-11-26 -Initial release with multi-provider LLM support. +Initial release with multi-provider LLM support. \ No newline at end of file diff --git a/packages/ai/scripts/generate-models.ts b/packages/ai/scripts/generate-models.ts index 4f662bfe0..eb238c210 100644 --- a/packages/ai/scripts/generate-models.ts +++ b/packages/ai/scripts/generate-models.ts @@ -192,7 +192,7 @@ interface KimiModelInfo { async function fetchKimiCodeModels(): Promise[]> { // Kimi Code /models endpoint requires authentication // Use KIMI_API_KEY env var if available, otherwise return fallback models - + const apiKey = $env.KIMI_API_KEY; if (apiKey) { try { console.log("Fetching models from Kimi Code API..."); @@ -924,6 +924,18 @@ async function generateModels() { contextWindow: CODEX_CONTEXT, maxTokens: CODEX_MAX_TOKENS, }, + { + id: "gpt-5.3-codex", + name: "GPT-5.3 Codex", + api: "openai-codex-responses", + provider: "openai-codex", + baseUrl: CODEX_BASE_URL, + reasoning: true, + input: ["text", "image"], + cost: { input: 1.75, output: 14, cacheRead: 0.175, cacheWrite: 0 }, + contextWindow: 400000, + maxTokens: CODEX_MAX_TOKENS, + }, ]; allModels.push(...codexModels); diff --git a/packages/ai/src/models.generated.ts b/packages/ai/src/models.generated.ts index c7ed7eaa4..142275715 100644 --- a/packages/ai/src/models.generated.ts +++ b/packages/ai/src/models.generated.ts @@ -107,6 +107,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 4096, } satisfies Model<"bedrock-converse-stream">, + "anthropic.claude-opus-4-6-v1:0": { + id: "anthropic.claude-opus-4-6-v1:0", + name: "Claude Opus 4.6", + api: "bedrock-converse-stream", + provider: "amazon-bedrock", + baseUrl: "https://bedrock-runtime.us-east-1.amazonaws.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 200000, + maxTokens: 128000, + } satisfies Model<"bedrock-converse-stream">, "cohere.command-r-plus-v1:0": { id: "cohere.command-r-plus-v1:0", name: "Command R+", @@ -192,6 +209,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"bedrock-converse-stream">, + "eu.anthropic.claude-opus-4-6-v1:0": { + id: "eu.anthropic.claude-opus-4-6-v1:0", + name: "Claude Opus 4.6 (EU)", + api: "bedrock-converse-stream", + provider: "amazon-bedrock", + baseUrl: "https://bedrock-runtime.us-east-1.amazonaws.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 1.5, + cacheWrite: 18.75, + }, + contextWindow: 200000, + maxTokens: 128000, + } satisfies Model<"bedrock-converse-stream">, "eu.anthropic.claude-sonnet-4-20250514-v1:0": { id: "eu.anthropic.claude-sonnet-4-20250514-v1:0", name: "Claude Sonnet 4 (EU)", @@ -277,6 +311,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"bedrock-converse-stream">, + "global.anthropic.claude-opus-4-6-v1:0": { + id: "global.anthropic.claude-opus-4-6-v1:0", + name: "Claude Opus 4.6 (Global)", + api: "bedrock-converse-stream", + provider: "amazon-bedrock", + baseUrl: "https://bedrock-runtime.us-east-1.amazonaws.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 200000, + maxTokens: 128000, + } satisfies Model<"bedrock-converse-stream">, "global.anthropic.claude-sonnet-4-20250514-v1:0": { id: "global.anthropic.claude-sonnet-4-20250514-v1:0", name: "Claude Sonnet 4", @@ -855,6 +906,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"bedrock-converse-stream">, + "us.anthropic.claude-opus-4-6-v1:0": { + id: "us.anthropic.claude-opus-4-6-v1:0", + name: "Claude Opus 4.6 (US)", + api: "bedrock-converse-stream", + provider: "amazon-bedrock", + baseUrl: "https://bedrock-runtime.us-east-1.amazonaws.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 1.5, + cacheWrite: 18.75, + }, + contextWindow: 200000, + maxTokens: 128000, + } satisfies Model<"bedrock-converse-stream">, "us.anthropic.claude-sonnet-4-20250514-v1:0": { id: "us.anthropic.claude-sonnet-4-20250514-v1:0", name: "Claude Sonnet 4 (US)", @@ -1316,6 +1384,40 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"anthropic-messages">, + "claude-opus-4-6": { + id: "claude-opus-4-6", + name: "Claude Opus 4.6", + api: "anthropic-messages", + provider: "anthropic", + baseUrl: "https://api.anthropic.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 1000000, + maxTokens: 128000, + } satisfies Model<"anthropic-messages">, + "claude-opus-4-6-20260205": { + id: "claude-opus-4-6-20260205", + name: "Claude Opus 4.6", + api: "anthropic-messages", + provider: "anthropic", + baseUrl: "https://api.anthropic.com", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 200000, + maxTokens: 128000, + } satisfies Model<"anthropic-messages">, "claude-sonnet-4-0": { id: "claude-sonnet-4-0", name: "Claude Sonnet 4 (latest)", @@ -1700,6 +1802,25 @@ export const MODELS = { contextWindow: 128000, maxTokens: 16000, } satisfies Model<"openai-completions">, + "claude-opus-4.6": { + id: "claude-opus-4.6", + name: "Claude Opus 4.6", + api: "openai-completions", + provider: "github-copilot", + baseUrl: "https://api.individual.githubcopilot.com", + headers: {"User-Agent":"GitHubCopilotChat/0.35.0","Editor-Version":"vscode/1.107.0","Editor-Plugin-Version":"copilot-chat/0.35.0","Copilot-Integration-Id":"vscode-chat"}, + compat: {"supportsStore":false,"supportsDeveloperRole":false,"supportsReasoningEffort":false}, + reasoning: true, + input: ["text", "image"], + cost: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + }, + contextWindow: 128000, + maxTokens: 16000, + } satisfies Model<"openai-completions">, "claude-sonnet-4": { id: "claude-sonnet-4", name: "Claude Sonnet 4", @@ -3030,63 +3151,6 @@ export const MODELS = { contextWindow: 262144, maxTokens: 32000, } satisfies Model<"openai-completions">, - "kimi-k2": { - id: "kimi-k2", - name: "Kimi K2", - api: "openai-completions", - provider: "kimi-code", - baseUrl: "https://api.kimi.com/coding/v1", - headers: {"User-Agent":"KimiCLI/1.0","X-Msh-Platform":"kimi_cli"}, - compat: {"thinkingFormat":"zai","reasoningContentField":"reasoning_content","supportsDeveloperRole":false}, - reasoning: true, - input: ["text"], - cost: { - input: 0, - output: 0, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 262144, - maxTokens: 32000, - } satisfies Model<"openai-completions">, - "kimi-k2-turbo-preview": { - id: "kimi-k2-turbo-preview", - name: "Kimi K2 Turbo Preview", - api: "openai-completions", - provider: "kimi-code", - baseUrl: "https://api.kimi.com/coding/v1", - headers: {"User-Agent":"KimiCLI/1.0","X-Msh-Platform":"kimi_cli"}, - compat: {"thinkingFormat":"zai","reasoningContentField":"reasoning_content","supportsDeveloperRole":false}, - reasoning: true, - input: ["text"], - cost: { - input: 0, - output: 0, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 262144, - maxTokens: 32000, - } satisfies Model<"openai-completions">, - "kimi-k2.5": { - id: "kimi-k2.5", - name: "Kimi K2.5", - api: "openai-completions", - provider: "kimi-code", - baseUrl: "https://api.kimi.com/coding/v1", - headers: {"User-Agent":"KimiCLI/1.0","X-Msh-Platform":"kimi_cli"}, - compat: {"thinkingFormat":"zai","reasoningContentField":"reasoning_content","supportsDeveloperRole":false}, - reasoning: true, - input: ["text", "image"], - cost: { - input: 0, - output: 0, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 262144, - maxTokens: 32000, - } satisfies Model<"openai-completions">, }, "minimax": { "MiniMax-M2": { @@ -4030,6 +4094,23 @@ export const MODELS = { contextWindow: 400000, maxTokens: 128000, } satisfies Model<"openai-responses">, + "gpt-5.3-codex": { + id: "gpt-5.3-codex", + name: "GPT-5.3 Codex", + api: "openai-responses", + provider: "openai", + baseUrl: "https://api.openai.com/v1", + reasoning: true, + input: ["text", "image"], + cost: { + input: 1.75, + output: 14, + cacheRead: 0.175, + cacheWrite: 0, + }, + contextWindow: 400000, + maxTokens: 128000, + } satisfies Model<"openai-responses">, "o1": { id: "o1", name: "o1", @@ -4253,6 +4334,23 @@ export const MODELS = { contextWindow: 272000, maxTokens: 128000, } satisfies Model<"openai-codex-responses">, + "gpt-5.3-codex": { + id: "gpt-5.3-codex", + name: "GPT-5.3 Codex", + api: "openai-codex-responses", + provider: "openai-codex", + baseUrl: "https://chatgpt.com/backend-api", + reasoning: true, + input: ["text", "image"], + cost: { + input: 1.75, + output: 14, + cacheRead: 0.175, + cacheWrite: 0, + }, + contextWindow: 400000, + maxTokens: 128000, + } satisfies Model<"openai-codex-responses">, }, "opencode": { "big-pickle": { @@ -4340,6 +4438,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"anthropic-messages">, + "claude-opus-4-6": { + id: "claude-opus-4-6", + name: "Claude Opus 4.6", + api: "anthropic-messages", + provider: "opencode", + baseUrl: "https://opencode.ai/zen", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 1000000, + maxTokens: 128000, + } satisfies Model<"anthropic-messages">, "claude-sonnet-4": { id: "claude-sonnet-4", name: "Claude Sonnet 4", @@ -5060,6 +5175,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"openai-completions">, + "anthropic/claude-opus-4.6": { + id: "anthropic/claude-opus-4.6", + name: "Anthropic: Claude Opus 4.6", + api: "openai-completions", + provider: "openrouter", + baseUrl: "https://openrouter.ai/api/v1", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 1000000, + maxTokens: 128000, + } satisfies Model<"openai-completions">, "anthropic/claude-sonnet-4": { id: "anthropic/claude-sonnet-4", name: "Anthropic: Claude Sonnet 4", @@ -5265,57 +5397,6 @@ export const MODELS = { contextWindow: 128000, maxTokens: 4000, } satisfies Model<"openai-completions">, - "deepcogito/cogito-v2-preview-llama-109b-moe": { - id: "deepcogito/cogito-v2-preview-llama-109b-moe", - name: "Cogito V2 Preview Llama 109B", - api: "openai-completions", - provider: "openrouter", - baseUrl: "https://openrouter.ai/api/v1", - reasoning: true, - input: ["text", "image"], - cost: { - input: 0.18, - output: 0.59, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 32767, - maxTokens: 4096, - } satisfies Model<"openai-completions">, - "deepcogito/cogito-v2-preview-llama-405b": { - id: "deepcogito/cogito-v2-preview-llama-405b", - name: "Deep Cogito: Cogito V2 Preview Llama 405B", - api: "openai-completions", - provider: "openrouter", - baseUrl: "https://openrouter.ai/api/v1", - reasoning: true, - input: ["text"], - cost: { - input: 3.5, - output: 3.5, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 32768, - maxTokens: 4096, - } satisfies Model<"openai-completions">, - "deepcogito/cogito-v2-preview-llama-70b": { - id: "deepcogito/cogito-v2-preview-llama-70b", - name: "Deep Cogito: Cogito V2 Preview Llama 70B", - api: "openai-completions", - provider: "openrouter", - baseUrl: "https://openrouter.ai/api/v1", - reasoning: true, - input: ["text"], - cost: { - input: 0.88, - output: 0.88, - cacheRead: 0, - cacheWrite: 0, - }, - contextWindow: 32768, - maxTokens: 4096, - } satisfies Model<"openai-completions">, "deepseek/deepseek-chat": { id: "deepseek/deepseek-chat", name: "DeepSeek: DeepSeek V3", @@ -5412,7 +5493,7 @@ export const MODELS = { cost: { input: 0.21, output: 0.7899999999999999, - cacheRead: 0.16799999999999998, + cacheRead: 0.1300000002, cacheWrite: 0, }, contextWindow: 163840, @@ -9042,6 +9123,23 @@ export const MODELS = { contextWindow: 200000, maxTokens: 64000, } satisfies Model<"anthropic-messages">, + "anthropic/claude-opus-4.6": { + id: "anthropic/claude-opus-4.6", + name: "Claude Opus 4.6", + api: "anthropic-messages", + provider: "vercel-ai-gateway", + baseUrl: "https://ai-gateway.vercel.sh", + reasoning: true, + input: ["text", "image"], + cost: { + input: 5, + output: 25, + cacheRead: 0.5, + cacheWrite: 6.25, + }, + contextWindow: 1000000, + maxTokens: 128000, + } satisfies Model<"anthropic-messages">, "anthropic/claude-sonnet-4": { id: "anthropic/claude-sonnet-4", name: "Claude Sonnet 4", @@ -9799,13 +9897,13 @@ export const MODELS = { reasoning: true, input: ["text", "image"], cost: { - input: 0.44999999999999996, + input: 0.5, output: 2.8, cacheRead: 0, cacheWrite: 0, }, - contextWindow: 262144, - maxTokens: 252144, + contextWindow: 256000, + maxTokens: 256000, } satisfies Model<"anthropic-messages">, "nvidia/nemotron-nano-12b-v2-vl": { id: "nvidia/nemotron-nano-12b-v2-vl", diff --git a/packages/ai/src/providers/google-gemini-cli.ts b/packages/ai/src/providers/google-gemini-cli.ts index 539c594aa..e39658b6c 100644 --- a/packages/ai/src/providers/google-gemini-cli.ts +++ b/packages/ai/src/providers/google-gemini-cli.ts @@ -5,7 +5,7 @@ */ import { createHash } from "node:crypto"; import type { Content, ThinkingConfig } from "@google/genai"; -import { abortableSleep } from "@oh-my-pi/pi-utils"; +import { abortableSleep, readSseJson } from "@oh-my-pi/pi-utils"; import { calculateCost } from "../models"; import type { Api, @@ -523,211 +523,168 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( const blocks = output.content; const blockIndex = () => blocks.length - 1; - // Read SSE stream - const reader = activeResponse.body.getReader(); - const decoder = new TextDecoder(); - let buffer = ""; - let jsonlBuffer = ""; + for await (const chunk of readSseJson( + activeResponse.body!, + options?.signal, + )) { + const responseData = chunk.response; + if (!responseData) continue; - // Set up abort handler to cancel reader when signal fires - const abortHandler = () => { - void reader.cancel().catch(() => {}); - }; - options?.signal?.addEventListener("abort", abortHandler); - - try { - while (true) { - // Check abort signal before each read - if (options?.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; - jsonlBuffer += `${jsonStr}\n`; - const parsed = Bun.JSONL.parseChunk(jsonlBuffer); - jsonlBuffer = jsonlBuffer.slice(parsed.read); - if (parsed.error) { - jsonlBuffer = ""; - continue; - } - - const chunk = parsed.values[0] as CloudCodeAssistResponseChunk | undefined; - if (!chunk) continue; - - // Unwrap the response - const responseData = chunk.response; - if (!responseData) continue; - - const candidate = responseData.candidates?.[0]; - if (candidate?.content?.parts) { - for (const part of candidate.content.parts) { - if (part.text !== undefined) { - hasContent = true; - const isThinking = isThinkingPart(part); - if ( - !currentBlock || - (isThinking && currentBlock.type !== "thinking") || - (!isThinking && currentBlock.type !== "text") - ) { - if (currentBlock) { - if (currentBlock.type === "text") { - stream.push({ - type: "text_end", - contentIndex: blocks.length - 1, - content: currentBlock.text, - partial: output, - }); - } else { - stream.push({ - type: "thinking_end", - contentIndex: blockIndex(), - content: currentBlock.thinking, - partial: output, - }); - } - } - if (isThinking) { - currentBlock = { type: "thinking", thinking: "", thinkingSignature: undefined }; - output.content.push(currentBlock); - ensureStarted(); - stream.push({ - type: "thinking_start", - contentIndex: blockIndex(), - partial: output, - }); - } else { - currentBlock = { type: "text", text: "" }; - output.content.push(currentBlock); - ensureStarted(); - stream.push({ type: "text_start", contentIndex: blockIndex(), partial: output }); - } - } - if (currentBlock.type === "thinking") { - currentBlock.thinking += part.text; - currentBlock.thinkingSignature = retainThoughtSignature( - currentBlock.thinkingSignature, - part.thoughtSignature, - ); + const candidate = responseData.candidates?.[0]; + if (candidate?.content?.parts) { + for (const part of candidate.content.parts) { + if (part.text !== undefined) { + hasContent = true; + const isThinking = isThinkingPart(part); + if ( + !currentBlock || + (isThinking && currentBlock.type !== "thinking") || + (!isThinking && currentBlock.type !== "text") + ) { + if (currentBlock) { + if (currentBlock.type === "text") { stream.push({ - type: "thinking_delta", - contentIndex: blockIndex(), - delta: part.text, + type: "text_end", + contentIndex: blocks.length - 1, + content: currentBlock.text, partial: output, }); } else { - currentBlock.text += part.text; - currentBlock.textSignature = retainThoughtSignature( - currentBlock.textSignature, - part.thoughtSignature, - ); stream.push({ - type: "text_delta", + type: "thinking_end", contentIndex: blockIndex(), - delta: part.text, + content: currentBlock.thinking, partial: output, }); } } - - if (part.functionCall) { - hasContent = true; - if (currentBlock) { - if (currentBlock.type === "text") { - stream.push({ - type: "text_end", - contentIndex: blockIndex(), - content: currentBlock.text, - partial: output, - }); - } else { - stream.push({ - type: "thinking_end", - contentIndex: blockIndex(), - content: currentBlock.thinking, - partial: output, - }); - } - currentBlock = null; - } - - const providedId = part.functionCall.id; - const needsNewId = - !providedId || output.content.some(b => b.type === "toolCall" && b.id === providedId); - const toolCallId = needsNewId - ? `${part.functionCall.name}_${Date.now()}_${++toolCallCounter}` - : providedId; - - const toolCall: ToolCall = { - type: "toolCall", - id: toolCallId, - name: part.functionCall.name || "", - arguments: part.functionCall.args as Record, - ...(part.thoughtSignature && { thoughtSignature: part.thoughtSignature }), - }; - - output.content.push(toolCall); + if (isThinking) { + currentBlock = { type: "thinking", thinking: "", thinkingSignature: undefined }; + output.content.push(currentBlock); ensureStarted(); - stream.push({ type: "toolcall_start", contentIndex: blockIndex(), partial: output }); stream.push({ - type: "toolcall_delta", + type: "thinking_start", contentIndex: blockIndex(), - delta: JSON.stringify(toolCall.arguments), partial: output, }); + } else { + currentBlock = { type: "text", text: "" }; + output.content.push(currentBlock); + ensureStarted(); + stream.push({ type: "text_start", contentIndex: blockIndex(), partial: output }); + } + } + if (currentBlock.type === "thinking") { + currentBlock.thinking += part.text; + currentBlock.thinkingSignature = retainThoughtSignature( + currentBlock.thinkingSignature, + part.thoughtSignature, + ); + stream.push({ + type: "thinking_delta", + contentIndex: blockIndex(), + delta: part.text, + partial: output, + }); + } else { + currentBlock.text += part.text; + currentBlock.textSignature = retainThoughtSignature( + currentBlock.textSignature, + part.thoughtSignature, + ); + stream.push({ + type: "text_delta", + contentIndex: blockIndex(), + delta: part.text, + partial: output, + }); + } + } + + if (part.functionCall) { + hasContent = true; + if (currentBlock) { + if (currentBlock.type === "text") { stream.push({ - type: "toolcall_end", + type: "text_end", contentIndex: blockIndex(), - toolCall, + content: currentBlock.text, + partial: output, + }); + } else { + stream.push({ + type: "thinking_end", + contentIndex: blockIndex(), + content: currentBlock.thinking, partial: output, }); } + currentBlock = null; } - } - if (candidate?.finishReason) { - output.stopReason = mapStopReasonString(candidate.finishReason); - if (output.content.some(b => b.type === "toolCall")) { - output.stopReason = "toolUse"; - } - } + const providedId = part.functionCall.id; + const needsNewId = + !providedId || output.content.some(b => b.type === "toolCall" && b.id === providedId); + const toolCallId = needsNewId + ? `${part.functionCall.name}_${Date.now()}_${++toolCallCounter}` + : providedId; - if (responseData.usageMetadata) { - // promptTokenCount includes cachedContentTokenCount, so subtract to get fresh input - const promptTokens = responseData.usageMetadata.promptTokenCount || 0; - const cacheReadTokens = responseData.usageMetadata.cachedContentTokenCount || 0; - output.usage = { - input: promptTokens - cacheReadTokens, - output: - (responseData.usageMetadata.candidatesTokenCount || 0) + - (responseData.usageMetadata.thoughtsTokenCount || 0), - cacheRead: cacheReadTokens, - cacheWrite: 0, - totalTokens: responseData.usageMetadata.totalTokenCount || 0, - cost: { - input: 0, - output: 0, - cacheRead: 0, - cacheWrite: 0, - total: 0, - }, + const toolCall: ToolCall = { + type: "toolCall", + id: toolCallId, + name: part.functionCall.name || "", + arguments: part.functionCall.args as Record, + ...(part.thoughtSignature && { thoughtSignature: part.thoughtSignature }), }; - calculateCost(model, output.usage); + + output.content.push(toolCall); + ensureStarted(); + stream.push({ type: "toolcall_start", contentIndex: blockIndex(), partial: output }); + stream.push({ + type: "toolcall_delta", + contentIndex: blockIndex(), + delta: JSON.stringify(toolCall.arguments), + partial: output, + }); + stream.push({ + type: "toolcall_end", + contentIndex: blockIndex(), + toolCall, + partial: output, + }); } } } - } finally { - options?.signal?.removeEventListener("abort", abortHandler); + + if (candidate?.finishReason) { + output.stopReason = mapStopReasonString(candidate.finishReason); + if (output.content.some(b => b.type === "toolCall")) { + output.stopReason = "toolUse"; + } + } + + if (responseData.usageMetadata) { + // promptTokenCount includes cachedContentTokenCount, so subtract to get fresh input + const promptTokens = responseData.usageMetadata.promptTokenCount || 0; + const cacheReadTokens = responseData.usageMetadata.cachedContentTokenCount || 0; + output.usage = { + input: promptTokens - cacheReadTokens, + output: + (responseData.usageMetadata.candidatesTokenCount || 0) + + (responseData.usageMetadata.thoughtsTokenCount || 0), + cacheRead: cacheReadTokens, + cacheWrite: 0, + totalTokens: responseData.usageMetadata.totalTokenCount || 0, + cost: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + total: 0, + }, + }; + calculateCost(model, output.usage); + } } if (currentBlock) { diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index 099e42cf9..54c7feb6c 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -87,8 +87,8 @@ "test": "bun test" }, "dependencies": { - "@oclif/core": "^4.5.6", - "@oclif/plugin-autocomplete": "^3.2.23", + "@oclif/core": "^4.8.0", + "@oclif/plugin-autocomplete": "^3.2.40", "@mozilla/readability": "0.6.0", "@oh-my-pi/omp-stats": "workspace:*", "@oh-my-pi/pi-agent-core": "workspace:*", @@ -107,7 +107,7 @@ "jsdom": "28.0.0", "marked": "^17.0.1", "node-html-parser": "^7.0.2", - "puppeteer": "^24.36.1", + "puppeteer": "^24.37.1", "smol-toml": "^1.6.0", "zod": "^4.3.6" }, diff --git a/packages/coding-agent/src/mcp/transports/stdio.ts b/packages/coding-agent/src/mcp/transports/stdio.ts index 4bf2321cc..43256a235 100644 --- a/packages/coding-agent/src/mcp/transports/stdio.ts +++ b/packages/coding-agent/src/mcp/transports/stdio.ts @@ -4,6 +4,8 @@ * Implements JSON-RPC 2.0 over subprocess stdin/stdout. * Messages are newline-delimited JSON. */ + +import { readLines } from "@oh-my-pi/pi-utils"; import { type Subprocess, spawn } from "bun"; import type { JsonRpcResponse, MCPStdioServerConfig, MCPTransport } from "../../mcp/types"; @@ -25,7 +27,6 @@ export class StdioTransport implements MCPTransport { reject: (error: Error) => void; } >(); - private buffer = ""; private _connected = false; private readLoop: Promise | null = null; @@ -72,23 +73,21 @@ export class StdioTransport implements MCPTransport { private async startReadLoop(): Promise { if (!this.process?.stdout) return; - const reader = this.process.stdout.getReader(); const decoder = new TextDecoder(); - try { - while (this._connected) { - const { done, value } = await reader.read(); - if (done) break; - - this.buffer += decoder.decode(value, { stream: true }); - this.processBuffer(); + for await (const line of readLines(this.process.stdout)) { + if (!this._connected) break; + try { + this.handleMessage(JSON.parse(decoder.decode(line)) as JsonRpcResponse); + } catch { + // Skip malformed lines + } } } catch (error) { if (this._connected) { this.onError?.(error instanceof Error ? error : new Error(String(error))); } } finally { - reader.releaseLock(); this.handleClose(); } } @@ -117,29 +116,6 @@ export class StdioTransport implements MCPTransport { } } - private processBuffer(): void { - while (this.buffer.length > 0) { - const result = Bun.JSONL.parseChunk(this.buffer); - for (const message of result.values) { - this.handleMessage(message as JsonRpcResponse); - } - - if (result.error) { - const nextNewline = this.buffer.indexOf("\n", result.read); - if (nextNewline === -1) { - this.buffer = ""; - break; - } - this.buffer = this.buffer.slice(nextNewline + 1); - continue; - } - - if (result.read === 0) break; - this.buffer = this.buffer.slice(result.read); - if (result.done) break; - } - } - private handleMessage(message: JsonRpcResponse): void { // Check if it's a response (has id) if ("id" in message && message.id !== null) { diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index 57407d410..d7948257a 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 { createTextLineSplitter, ptree } from "@oh-my-pi/pi-utils"; +import { ptree, readJsonl } 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"; @@ -83,11 +83,11 @@ function isAgentEvent(value: unknown): value is AgentEvent { export class RpcClient { private process: ptree.ChildProcess | null = null; - private lineReader: ReadableStream | null = null; private eventListeners: RpcEventListener[] = []; private pendingRequests: Map void; reject: (error: Error) => void }> = new Map(); private requestId = 0; + private abortController = new AbortController(); constructor(private options: RpcClientOptions = {}) {} @@ -119,19 +119,12 @@ export class RpcClient { }); // Process lines in background - const lines = this.process.stdout.pipeThrough(createTextLineSplitter(true)); - this.lineReader = lines; + const lines = readJsonl(this.process.stdout, this.abortController.signal); void (async () => { - try { - for await (const line of lines) { - this.handleLine(line); - } - } catch { - // Stream closed - } finally { - lines.cancel(); + for await (const line of lines) { + this.handleLine(line); } - })(); + })().catch(() => {}); // Wait a moment for process to initialize await Bun.sleep(100); @@ -154,11 +147,9 @@ export class RpcClient { stop() { if (!this.process) return; - this.lineReader?.cancel(); this.process.kill(); - + this.abortController.abort(); this.process = null; - this.lineReader = null; this.pendingRequests.clear(); } @@ -461,29 +452,23 @@ export class RpcClient { // Internal // ========================================================================= - private handleLine(line: string): void { - const result = Bun.JSONL.parseChunk(line); - if (result.error) return; - - for (const data of result.values) { - // Check if it's a response to a pending request - if (isRpcResponse(data)) { - const id = data.id; - if (id && this.pendingRequests.has(id)) { - const pending = this.pendingRequests.get(id)!; - this.pendingRequests.delete(id); - pending.resolve(data); - return; - } - continue; + private handleLine(data: unknown): void { + // Check if it's a response to a pending request + if (isRpcResponse(data)) { + const id = data.id; + if (id && this.pendingRequests.has(id)) { + const pending = this.pendingRequests.get(id)!; + this.pendingRequests.delete(id); + pending.resolve(data); + return; } + } - if (!isAgentEvent(data)) continue; + if (!isAgentEvent(data)) return; - // Otherwise it's an event - for (const listener of this.eventListeners) { - listener(data); - } + // Otherwise it's an event + for (const listener of this.eventListeners) { + listener(data); } } diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index aae69e683..f5064f6c7 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 { createTextLineSplitter, Snowflake } from "@oh-my-pi/pi-utils"; +import { readJsonl, 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,37 +633,27 @@ export async function runRpcMode(session: AgentSession): Promise { } // Listen for JSON input using Bun's stdin - for await (const line of Bun.stdin.stream().pipeThrough(createTextLineSplitter())) { - if (!line.trim()) continue; - - const result = Bun.JSONL.parseChunk(`${line}\n`); - if (result.error) { - output(error(undefined, "parse", `Failed to parse command: ${result.error.message}`)); - continue; - } - - for (const parsed of result.values) { - try { - // Handle extension UI responses - if ((parsed as RpcExtensionUIResponse).type === "extension_ui_response") { - const response = parsed as RpcExtensionUIResponse; - const pending = pendingExtensionRequests.get(response.id); - if (pending) { - pending.resolve(response); - } - continue; + for await (const parsed of readJsonl(Bun.stdin.stream())) { + try { + // Handle extension UI responses + if ((parsed as RpcExtensionUIResponse).type === "extension_ui_response") { + const response = parsed as RpcExtensionUIResponse; + const pending = pendingExtensionRequests.get(response.id); + if (pending) { + pending.resolve(response); } - - // Handle regular commands - const command = parsed as RpcCommand; - const response = await handleCommand(command); - output(response); - - // Check for deferred shutdown request (idle between commands) - await checkShutdownRequested(); - } catch (e: any) { - output(error(undefined, "parse", `Failed to parse command: ${e.message}`)); + continue; } + + // Handle regular commands + const command = parsed as RpcCommand; + const response = await handleCommand(command); + output(response); + + // Check for deferred shutdown request (idle between commands) + await checkShutdownRequested(); + } catch (e: any) { + output(error(undefined, "parse", `Failed to parse command: ${e.message}`)); } } diff --git a/packages/coding-agent/src/session/session-manager.ts b/packages/coding-agent/src/session/session-manager.ts index 78321c0a6..5f0228fc8 100644 --- a/packages/coding-agent/src/session/session-manager.ts +++ b/packages/coding-agent/src/session/session-manager.ts @@ -1,7 +1,7 @@ import * as path from "node:path"; import type { AgentMessage } from "@oh-my-pi/pi-agent-core"; import type { ImageContent, Message, TextContent, Usage } from "@oh-my-pi/pi-ai"; -import { isEnoent, logger, Snowflake } from "@oh-my-pi/pi-utils"; +import { isEnoent, logger, parseJsonlLenient, Snowflake } from "@oh-my-pi/pi-utils"; import { getAgentDir as getDefaultAgentDir } from "../config"; import { resizeImage } from "../utils/image-resize"; import { @@ -279,32 +279,6 @@ function migrateToCurrentVersion(entries: FileEntry[]): boolean { return true; } -function parseJsonlEntries(buffer: string): T[] { - let entries: T[] | undefined; - - while (buffer.length > 0) { - 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 (error) { - const nextNewline = buffer.indexOf("\n", read); - if (nextNewline === -1) break; - buffer = buffer.substring(nextNewline + 1); - continue; - } - if (read === 0) break; - buffer = buffer.substring(read); - if (done) break; - } - return entries ?? []; -} - /** Exported for testing */ export function migrateSessionEntries(entries: FileEntry[]): void { migrateToCurrentVersion(entries); @@ -312,7 +286,7 @@ export function migrateSessionEntries(entries: FileEntry[]): void { /** Exported for compaction.test.ts */ export function parseSessionEntries(content: string): FileEntry[] { - return parseJsonlEntries(content); + return parseJsonlLenient(content); } export function getLatestCompactionEntry(entries: SessionEntry[]): CompactionEntry | null { @@ -485,7 +459,7 @@ export async function loadEntriesFromFile( if (isEnoent(err)) return []; throw err; } - const entries = parseJsonlEntries(content); + const entries = parseJsonlLenient(content); // Validate session header if (entries.length === 0) return entries; @@ -587,7 +561,7 @@ async function getSortedSessions(sessionDir: string, storage: SessionStorage): P files.map(async (path: string) => { try { const content = await storage.readTextPrefix(path, 4096); - const entries = parseJsonlEntries>(content); + const entries = parseJsonlLenient>(content); if (entries.length === 0) return; const header = entries[0] as Record; if (header.type !== "session" || typeof header.id !== "string") return; @@ -915,7 +889,7 @@ async function collectSessionsFromFiles(files: string[], storage: SessionStorage files.map(async file => { try { const content = await storage.readText(file); - const entries = parseJsonlEntries>(content); + const entries = parseJsonlLenient>(content); if (entries.length === 0) return; // Check first entry for valid session header diff --git a/packages/coding-agent/src/session/session-storage.ts b/packages/coding-agent/src/session/session-storage.ts index 31f8bf390..ce16334b3 100644 --- a/packages/coding-agent/src/session/session-storage.ts +++ b/packages/coding-agent/src/session/session-storage.ts @@ -169,7 +169,7 @@ export class FileSessionStorage implements SessionStorage { async readTextPrefix(path: string, maxBytes: number): Promise { const handle = await fs.promises.open(path, "r"); try { - const buffer = Buffer.alloc(maxBytes); + const buffer = Buffer.allocUnsafe(maxBytes); const { bytesRead } = await handle.read(buffer, 0, maxBytes, 0); return buffer.subarray(0, bytesRead).toString("utf-8"); } finally { diff --git a/packages/coding-agent/src/tools/read.ts b/packages/coding-agent/src/tools/read.ts index 17a6a8e6f..f67ef8aad 100644 --- a/packages/coding-agent/src/tools/read.ts +++ b/packages/coding-agent/src/tools/read.ts @@ -43,7 +43,7 @@ function isRemoteMountPath(absolutePath: string): boolean { return absolutePath.startsWith(REMOTE_MOUNT_PREFIX); } -const READ_CHUNK_SIZE = 64 * 1024; +const READ_CHUNK_SIZE = 8 * 1024; async function streamLinesFromFile( filePath: string, diff --git a/packages/coding-agent/src/utils/mime.ts b/packages/coding-agent/src/utils/mime.ts index 3653f632b..06dae82ef 100644 --- a/packages/coding-agent/src/utils/mime.ts +++ b/packages/coding-agent/src/utils/mime.ts @@ -8,7 +8,7 @@ const FILE_TYPE_SNIFF_BYTES = 4100; export async function detectSupportedImageMimeTypeFromFile(filePath: string): Promise { const fileHandle = await fs.open(filePath, "r"); try { - const buffer = Buffer.alloc(FILE_TYPE_SNIFF_BYTES); + const buffer = Buffer.allocUnsafe(FILE_TYPE_SNIFF_BYTES); const { bytesRead } = await fileHandle.read(buffer, 0, FILE_TYPE_SNIFF_BYTES, 0); if (bytesRead === 0) { return null; diff --git a/packages/coding-agent/src/web/scrapers/utils.ts b/packages/coding-agent/src/web/scrapers/utils.ts index 2c1c50a09..d0d052f4f 100644 --- a/packages/coding-agent/src/web/scrapers/utils.ts +++ b/packages/coding-agent/src/web/scrapers/utils.ts @@ -68,9 +68,11 @@ export async function convertWithMarkitdown( } } +const kEmptyBuffer = Buffer.alloc(0); + export async function fetchBinary(url: string, timeout: number, userSignal?: AbortSignal): Promise { if (userSignal?.aborted) { - return { buffer: Buffer.alloc(0), contentType: "", ok: false, error: "aborted" }; + return { buffer: kEmptyBuffer, contentType: "", ok: false, error: "aborted" }; } const signal = ptree.combineSignals(userSignal, timeout * 1000); @@ -89,7 +91,7 @@ export async function fetchBinary(url: string, timeout: number, userSignal?: Abo if (!response.ok) { return { - buffer: Buffer.alloc(0), + buffer: kEmptyBuffer, contentType, contentDisposition, ok: false, @@ -103,7 +105,7 @@ export async function fetchBinary(url: string, timeout: number, userSignal?: Abo const size = Number.parseInt(contentLength, 10); if (Number.isFinite(size) && size > MAX_BYTES) { return { - buffer: Buffer.alloc(0), + buffer: kEmptyBuffer, contentType, contentDisposition, ok: false, @@ -116,7 +118,7 @@ export async function fetchBinary(url: string, timeout: number, userSignal?: Abo const buffer = Buffer.from(await response.arrayBuffer()); if (buffer.length > MAX_BYTES) { return { - buffer: Buffer.alloc(0), + buffer: kEmptyBuffer, contentType, contentDisposition, ok: false, @@ -128,10 +130,10 @@ export async function fetchBinary(url: string, timeout: number, userSignal?: Abo return { buffer, contentType, contentDisposition, ok: true, status: response.status }; } catch (err) { if (signal?.aborted) { - return { buffer: Buffer.alloc(0), contentType: "", ok: false, error: "aborted" }; + return { buffer: kEmptyBuffer, contentType: "", ok: false, error: "aborted" }; } return { - buffer: Buffer.alloc(0), + buffer: kEmptyBuffer, contentType: "", ok: false, error: `request failed: ${String(err)}`, diff --git a/packages/utils/src/stream.ts b/packages/utils/src/stream.ts index 72875f7eb..10ae0fb0e 100644 --- a/packages/utils/src/stream.ts +++ b/packages/utils/src/stream.ts @@ -1,5 +1,3 @@ -import { ArrayBufferSink } from "bun"; - /** * Sanitize binary output for display/storage. * Removes characters that crash string-width or cause display issues: @@ -42,58 +40,70 @@ export function sanitizeText(text: string): string { return sanitizeBinaryOutput(Bun.stripANSI(text)).replace(/\r/g, ""); } -/** - * Create a transform stream that splits lines. - */ -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 +const LF = 0x0a; - 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; +export async function* readLines(stream: ReadableStream, signal?: AbortSignal): AsyncGenerator { + const buffer = new ConcatSink(); + const source = signal ? stream.pipeThrough(new TransformStream(), { signal }) : stream; + try { + for await (const chunk of source) { + for (const line of buffer.appendAndFlushLines(chunk)) { + yield line; } - }, - flush(ctrl) { - if (pending) { - const tail = sink.end() as Uint8Array; - if (tail.length > 0) ctrl.enqueue(mapFn(tail)); + } + if (!buffer.isEmpty) { + const tail = buffer.flush(); + if (tail) { + buffer.clear(); + yield tail; } - }, - }); + } + } catch (err) { + // Abort errors are expected — just stop the generator. + if (signal?.aborted) return; + throw err; + } } -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)) }); +export async function* readJsonl(stream: ReadableStream, signal?: AbortSignal): AsyncGenerator { + const buffer = new ConcatSink(); + const source = signal ? stream.pipeThrough(new TransformStream(), { signal }) : stream; + try { + const yieldBuffer: T[] = []; + for await (const chunk of source) { + buffer.appendAndConsume(chunk, 0, chunk.length, (payload, beg, end) => { + const { values, error, read, done } = Bun.JSONL.parseChunk(payload, beg, end); + if (values.length > 0) { + yieldBuffer.push(...(values as T[])); + } + if (error) throw error; + if (done) return 0; + return end - read; + }); + if (yieldBuffer.length > 0) { + yield* yieldBuffer; + yieldBuffer.length = 0; + } + } + if (!buffer.isEmpty) { + const tail = buffer.flush(); + if (tail) { + buffer.clear(); + const { values, error, done } = Bun.JSONL.parseChunk(tail, 0, tail.length); + if (values.length > 0) { + yield* values as T[]; + } + if (error) throw error; + if (!done) { + throw new Error("JSONL stream ended unexpectedly"); + } + } + } + } catch (err) { + // Abort errors are expected — just stop the generator. + if (signal?.aborted) return; + throw err; } - return createSplitterStream({ mapFn: dec.decode.bind(dec) }); } /** @@ -118,28 +128,169 @@ export function createTextDecoderStream(): TransformStream { // SSE (Server-Sent Events) // ============================================================================= -const LF = 0x0a; -const CR = 0x0d; -const SPACE = 0x20; - -// "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; // : - -// "[DONE]" = [0x5b, 0x44, 0x4f, 0x4e, 0x45, 0x5d] -const DONE = Uint8Array.from([0x5b, 0x44, 0x4f, 0x4e, 0x45, 0x5d]); - -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; +class Bitmap { + private bits: Uint32Array; + constructor(n: number) { + this.bits = new Uint32Array((n + 31) >>> 5); + } + + set(i: number, value: boolean) { + const index = i >>> 5; + const mask = 1 << (i & 31); + if (value) { + this.bits[index] |= mask; + } else { + this.bits[index] &= ~mask; + } + } + get(i: number) { + const index = i >>> 5; + const mask = 1 << (i & 31); + const word = this.bits[index]; + return word !== undefined && (word & mask) !== 0; } - return true; } +const WHITESPACE = new Bitmap(256); +for (let i = 0; i <= 0x7f; i++) { + const c = String.fromCharCode(i); + switch (c) { + case " ": + case "\t": + case "\n": + case "\r": + WHITESPACE.set(i, true); + break; + default: + WHITESPACE.set(i, !c.trim()); + break; + } +} + +const createPattern = (prefix: string) => { + const pre = Buffer.from(prefix, "utf-8"); + return { + strip(buf: Uint8Array): number | null { + const n = pre.length; + if (buf.length < n) return null; + if (pre.equals(buf.subarray(0, n))) { + return n; + } + return null; + }, + }; +}; + +const PAT_DATA = createPattern("data:"); + +const PAT_DONE = createPattern("[DONE]"); + +class ConcatSink { + #space?: Buffer; + #length = 0; + + #ensureCapacity(size: number): Buffer { + const space = this.#space; + if (space && space.length >= size) return space; + const nextSize = space ? Math.max(size, space.length * 2) : size; + const next = Buffer.allocUnsafe(nextSize); + if (space && this.#length > 0) { + space.copy(next, 0, 0, this.#length); + } + this.#space = next; + return next; + } + + append(chunk: Uint8Array) { + const n = chunk.length; + if (!n) return; + const offset = this.#length; + const space = this.#ensureCapacity(offset + n); + space.set(chunk, offset); + this.#length += n; + } + + reset(chunk: Uint8Array) { + const n = chunk.length; + if (!n) { + this.#length = 0; + return; + } + const space = this.#ensureCapacity(n); + space.set(chunk, 0); + this.#length = n; + } + + get isEmpty(): boolean { + return this.#length === 0; + } + + flush(): Uint8Array | undefined { + if (!this.#length) return undefined; + return this.#space!.subarray(0, this.#length); + } + + clear() { + this.#length = 0; + } + + *appendAndFlushLines(chunk: Uint8Array) { + let pos = 0; + while (pos < chunk.length) { + const nl = chunk.indexOf(LF, pos); + if (nl === -1) { + this.append(chunk.subarray(pos)); + return; + } + const suffix = chunk.subarray(pos, nl); + pos = nl + 1; + if (this.isEmpty) { + yield suffix; + } else { + this.append(suffix); + const payload = this.flush(); + if (payload) { + yield payload; + this.clear(); + } + } + } + } + + appendAndConsume( + chunk: Uint8Array, + beg: number, + end: number, + // (slice) => [remaining length] + consumer: (payload: Uint8Array, beg: number, end: number) => number, + ) { + if (this.isEmpty) { + const rem = consumer(chunk, beg, end); + if (!rem) return; + this.reset(chunk.subarray(end - rem, end)); + return; + } + + const offset = this.#length; + const n = end - beg; + const total = offset + n; + const space = this.#ensureCapacity(total); + space.set(chunk.subarray(beg, end), offset); + this.#length = total; + const rem = consumer(space.subarray(0, total), 0, total); + if (!rem) { + this.#length = 0; + return; + } + if (rem < total) { + space.copyWithin(0, total - rem, total); + } + this.#length = rem; + } +} + +const kDoneError = new Error("SSE stream done"); + /** * Stream parsed JSON objects from SSE `data:` lines. * @@ -150,96 +301,113 @@ function isDone(buf: Uint8Array, start: number, end: number): boolean { * } * ``` */ -export async function* readSseJson( - stream: ReadableStream, - abortSignal?: AbortSignal, -): AsyncGenerator { - const sink = new ArrayBufferSink(); - sink.start({ asUint8Array: true, stream: true, highWaterMark: 4096 }); - let pending = false; +export async function* readSseJson(stream: ReadableStream, signal?: AbortSignal): AsyncGenerator { + const lineBuffer = new ConcatSink(); + const jsonBuffer = new ConcatSink(); // pipeThrough with { signal } makes the stream abort-aware: the pipe // cancels the source and errors the output when the signal fires, // so for-await-of exits cleanly without manual reader/listener management. - const source = abortSignal ? stream.pipeThrough(new TransformStream(), { signal: abortSignal }) : stream; - + stream = signal ? stream.pipeThrough(new TransformStream(), { signal }) : stream; try { - for await (const chunk of source) { - 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 yieldBuffer: T[] = []; + const processLine = (line: Uint8Array) => { + // Strip trailing spaces including \r. + let end = line.length; + while (end && WHITESPACE.get(line[end - 1])) { + --end; + } + if (!end) return; // blank line + + const trimmed = end === line.length ? line : line.subarray(0, end); + + // Check "data:" prefix and optional space afterwards. + let beg = PAT_DATA.strip(trimmed); + if (beg === null) return; + while (beg < end && WHITESPACE.get(trimmed[beg])) { + ++beg; + } + if (beg >= end) return; + + jsonBuffer.appendAndConsume(trimmed, beg, end, (payload, beg, end) => { + const { values, error, read, done } = Bun.JSONL.parseChunk(payload, beg, end); + if (values.length > 0) { + yieldBuffer.push(...(values as T[])); } - - 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); + if (error) { + if (PAT_DONE.strip(payload.subarray(beg, end))) { + throw kDoneError; + } + throw error; + } + if (done) return 0; + return end - read; + }); + }; + for await (const chunk of stream) { + for (const line of lineBuffer.appendAndFlushLines(chunk)) { + processLine(line); + if (yieldBuffer.length > 0) { + yield* yieldBuffer; + yieldBuffer.length = 0; + } + } + } + if (!lineBuffer.isEmpty) { + const tail = lineBuffer.flush(); + if (tail) { + lineBuffer.clear(); + processLine(tail); + if (yieldBuffer.length > 0) { + yield* yieldBuffer; + yieldBuffer.length = 0; } - 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; } } } catch (err) { + if (err === kDoneError) return; // Abort errors are expected — just stop the generator. - if (abortSignal?.aborted) return; + if (signal?.aborted) return; throw err; } - - // 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; - } - } + if (!jsonBuffer.isEmpty) { + throw new Error("SSE stream ended unexpectedly"); } } + +/** + * Parse a complete JSONL string, skipping malformed lines instead of throwing. + * + * Uses `Bun.JSONL.parseChunk` internally. On parse errors, the malformed + * region is skipped up to the next newline and parsing continues. + * + * @example + * ```ts + * const entries = parseJsonlLenient(fileContents); + * ``` + */ +export function parseJsonlLenient(buffer: string): T[] { + let entries: T[] | undefined; + + while (buffer.length > 0) { + 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 (error) { + const nextNewline = buffer.indexOf("\n", read); + if (nextNewline === -1) break; + buffer = buffer.substring(nextNewline + 1); + continue; + } + if (read === 0) break; + buffer = buffer.substring(read); + if (done) break; + } + return entries ?? []; +} diff --git a/packages/utils/test/stream.test.ts b/packages/utils/test/stream.test.ts index 44473d3d0..3832e6c94 100644 --- a/packages/utils/test/stream.test.ts +++ b/packages/utils/test/stream.test.ts @@ -1,9 +1,10 @@ import { describe, expect, it } from "bun:test"; import { createSanitizerStream, - createSplitterStream, createTextDecoderStream, - createTextLineSplitter, + parseJsonlLenient, + readJsonl, + readLines, readSseJson, sanitizeBinaryOutput, sanitizeText, @@ -67,39 +68,51 @@ describe("sanitizeText", () => { }); }); -describe("createSplitterStream", () => { +describe("readLines", () => { it("splits lines across chunks without newlines", async () => { - const transform = createSplitterStream({ - mapFn: chunk => new TextDecoder().decode(chunk), + const readable = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("alpha\nbe")); + controller.enqueue(encoder.encode("ta\ngam")); + controller.enqueue(encoder.encode("ma")); + controller.close(); + }, }); - const output = await runTransform(transform, [ - encoder.encode("alpha\nbe"), - encoder.encode("ta\ngam"), - encoder.encode("ma"), - ]); + const output: string[] = []; + const dec = new TextDecoder(); + for await (const line of readLines(readable)) { + output.push(dec.decode(line)); + } 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")]); +describe("readJsonl", () => { + it("parses JSONL across chunk boundaries", async () => { + const readable = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode('{"a":1}\n{"b":')); + controller.enqueue(encoder.encode('2}\n{"c":3}\n')); + controller.close(); + }, + }); - expect(output).toEqual(["green", "blue"]); + const output = await collectAsync(readJsonl(readable)); + expect(output).toEqual([{ a: 1 }, { b: 2 }, { c: 3 }]); + }); + + it("parses trailing line without newline", async () => { + const readable = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode('{"z":9}')); + controller.close(); + }, + }); + + const output = await collectAsync(readJsonl(readable)); + expect(output).toEqual([{ z: 9 }]); }); }); @@ -121,6 +134,27 @@ describe("createTextDecoderStream", () => { }); }); +describe("parseJsonlLenient", () => { + it("parses valid JSONL", () => { + const result = parseJsonlLenient<{ a: number }>('{"a":1}\n{"a":2}\n{"a":3}\n'); + expect(result).toEqual([{ a: 1 }, { a: 2 }, { a: 3 }]); + }); + + it("skips malformed lines and continues", () => { + const result = parseJsonlLenient<{ a: number }>('{"a":1}\n{bad json}\n{"a":3}\n'); + expect(result).toEqual([{ a: 1 }, { a: 3 }]); + }); + + it("returns empty array for empty input", () => { + expect(parseJsonlLenient("")).toEqual([]); + }); + + it("handles input without trailing newline", () => { + const result = parseJsonlLenient<{ x: number }>('{"x":42}'); + expect(result).toEqual([{ x: 42 }]); + }); +}); + describe("readSseJson", () => { it("parses data lines and stops at [DONE]", async () => { const chunks = [