refactor(stream): improved SSE parsing & added GPT-5.3, Claude 4.6 Opus

- Extracted stream parsing utilities into reusable functions in pi-utils package (readLines, readJsonl, parseJsonlLenient, readSseJson).
- Replaced manual buffer management with ConcatSink class for efficient stream handling across multiple modules.
- Simplified SSE stream parsing in google-gemini-cli provider by using readSseJson utility instead of manual reader setup.
- Refactored MCP stdio transport to use readLines utility for cleaner line-based stream processing.
- Updated RPC client and mode to use readJsonl utility, removing duplicate JSONL parsing logic.
- Optimized buffer allocations by using allocUnsafe where data is immediately populated and reusing empty buffer instances.
This commit is contained in:
can1357
2026-02-05 21:00:13 +01:00
parent 0a8bd08565
commit b12a27830c
18 changed files with 823 additions and 612 deletions
Generated
+1 -1
View File
@@ -1989,7 +1989,7 @@ dependencies = [
[[package]]
name = "pi-natives"
version = "11.2.1"
version = "11.2.2"
dependencies = [
"arboard",
"brush-builtins",
+5 -5
View File
@@ -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=="],
-1
View File
@@ -52,7 +52,6 @@
"engines": {
"bun": ">=1.3.7"
},
"packageManager": "bun@1.3.7",
"version": "0.0.3",
"dependencies": {
"@sinclair/typebox": "^0.34.48"
+17 -1
View File
@@ -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.
+13 -1
View File
@@ -192,7 +192,7 @@ interface KimiModelInfo {
async function fetchKimiCodeModels(): Promise<Model<"openai-completions">[]> {
// 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);
+210 -112
View File
@@ -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",
+135 -178
View File
@@ -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<CloudCodeAssistResponseChunk>(
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<string, unknown>,
...(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<string, unknown>,
...(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) {
+3 -3
View File
@@ -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"
},
@@ -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<void> | null = null;
@@ -72,23 +73,21 @@ export class StdioTransport implements MCPTransport {
private async startReadLoop(): Promise<void> {
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) {
@@ -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<string> | null = null;
private eventListeners: RpcEventListener[] = [];
private pendingRequests: Map<string, { resolve: (response: RpcResponse) => 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);
}
}
+20 -30
View File
@@ -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<never> {
}
// 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}`));
}
}
@@ -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<T>(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<FileEntry>(content);
return parseJsonlLenient<FileEntry>(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<FileEntry>(content);
const entries = parseJsonlLenient<FileEntry>(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<Record<string, unknown>>(content);
const entries = parseJsonlLenient<Record<string, unknown>>(content);
if (entries.length === 0) return;
const header = entries[0] as Record<string, unknown>;
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<Record<string, unknown>>(content);
const entries = parseJsonlLenient<Record<string, unknown>>(content);
if (entries.length === 0) return;
// Check first entry for valid session header
@@ -169,7 +169,7 @@ export class FileSessionStorage implements SessionStorage {
async readTextPrefix(path: string, maxBytes: number): Promise<string> {
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 {
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -8,7 +8,7 @@ const FILE_TYPE_SNIFF_BYTES = 4100;
export async function detectSupportedImageMimeTypeFromFile(filePath: string): Promise<string | null> {
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;
@@ -68,9 +68,11 @@ export async function convertWithMarkitdown(
}
}
const kEmptyBuffer = Buffer.alloc(0);
export async function fetchBinary(url: string, timeout: number, userSignal?: AbortSignal): Promise<BinaryFetchResult> {
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)}`,
+313 -145
View File
@@ -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<T>(options: {
newLine?: boolean;
mapFn: (chunk: Uint8Array) => T;
}): TransformStream<Uint8Array, T> {
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<Uint8Array, T>({
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<Uint8Array>, signal?: AbortSignal): AsyncGenerator<Uint8Array> {
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<Uint8Array, string> {
const dec = new TextDecoder("utf-8", { ignoreBOM: true, fatal: true });
if (sanitize) {
return createSplitterStream({ mapFn: chunk => sanitizeText(dec.decode(chunk)) });
export async function* readJsonl<T>(stream: ReadableStream<Uint8Array>, signal?: AbortSignal): AsyncGenerator<T> {
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<Uint8Array, string> {
// 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<T>(
stream: ReadableStream<Uint8Array>,
abortSignal?: AbortSignal,
): AsyncGenerator<T> {
const sink = new ArrayBufferSink();
sink.start({ asUint8Array: true, stream: true, highWaterMark: 4096 });
let pending = false;
export async function* readSseJson<T>(stream: ReadableStream<Uint8Array>, signal?: AbortSignal): AsyncGenerator<T> {
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<MyType>(fileContents);
* ```
*/
export function parseJsonlLenient<T>(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 ?? [];
}
+60 -26
View File
@@ -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<Uint8Array>({
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<Uint8Array>({
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<Uint8Array>({
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 = [