refactor: consolidated json parsing and stream utilities
- Centralized JSON parsing and stream processing logic by moving utilities from `packages/ai` to the shared `@oh-my-pi/pi-utils` package. - Standardized import paths for JSON parsing and streaming across the agent, ai, and coding-agent packages. - Refactored SSE stream handling to use consolidated `parseStreamingJson` logic and introduced robust error recovery for malformed container-shaped tail events. - Cleaned up legacy bundled registry references and updated related module exports and tests to reflect the new utility structure.
This commit is contained in:
@@ -13,9 +13,8 @@ import {
|
||||
type StopReason,
|
||||
type ToolCall,
|
||||
} from "@oh-my-pi/pi-ai";
|
||||
import { parseStreamingJson } from "@oh-my-pi/pi-ai/utils/json-parse";
|
||||
import { calculateCost } from "@oh-my-pi/pi-catalog/models";
|
||||
import { readSseJson } from "@oh-my-pi/pi-utils";
|
||||
import { parseStreamingJson, readSseJson } from "@oh-my-pi/pi-utils";
|
||||
|
||||
// Event stream adapter for proxy SSE events
|
||||
export class ProxyMessageEventStream extends EventStream<AssistantMessageEvent, AssistantMessage> {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair } from "../utils/json-parse";
|
||||
import dialectPrompt from "./anthropic.md" with { type: "text" };
|
||||
import { buildArgShapes, buildStringArgsResolver, mintToolCallId, type ToolArgShape } from "./coercion";
|
||||
import {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair } from "../utils/json-parse";
|
||||
import { asRecord, mintToolCallId, partialSuffixOverlapAny } from "./coercion";
|
||||
import dialectPrompt from "./deepseek.md" with { type: "text" };
|
||||
import { assistantTranscriptParts, collectToolResultRun, messageContentText, stringifyJson } from "./rendering";
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair } from "../utils/json-parse";
|
||||
import { asRecord, mintToolCallId, partialSuffixOverlapAny } from "./coercion";
|
||||
import dialectPrompt from "./harmony.md" with { type: "text" };
|
||||
import {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair, parseStreamingJson } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair, parseStreamingJson } from "../utils/json-parse";
|
||||
import { asRecord, mintToolCallId, partialSuffixOverlapAny } from "./coercion";
|
||||
import dialectPrompt from "./hermes.md" with { type: "text" };
|
||||
import { renderChatMlTranscript, renderDelimitedThinking, renderToolResponseResults, stringifyJson } from "./rendering";
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair } from "../utils/json-parse";
|
||||
import { asRecord, normalizeKimiFunctionName, partialSuffixOverlapAny } from "./coercion";
|
||||
import dialectPrompt from "./kimi.md" with { type: "text" };
|
||||
import { assistantTranscriptParts, collectToolResultRun, messageContentText, stringifyJson } from "./rendering";
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { parseJsonWithRepair } from "@oh-my-pi/pi-utils";
|
||||
import type { Message, ToolCall } from "../types";
|
||||
import { parseJsonWithRepair } from "../utils/json-parse";
|
||||
import { asRecord, mintToolCallId, partialSuffixOverlapAny } from "./coercion";
|
||||
import dialectPrompt from "./qwen3.md" with { type: "text" };
|
||||
import { renderChatMlTranscript, renderToolResponseResults, stringifyJson } from "./rendering";
|
||||
|
||||
@@ -10,7 +10,14 @@
|
||||
import type { Effort } from "@oh-my-pi/pi-catalog/effort";
|
||||
import { mapEffortToAnthropicAdaptiveEffort, requireSupportedEffort } from "@oh-my-pi/pi-catalog/model-thinking";
|
||||
import { calculateCost } from "@oh-my-pi/pi-catalog/models";
|
||||
import { $env, $flag, extractHttpStatusFromError, fetchWithRetry } from "@oh-my-pi/pi-utils";
|
||||
import {
|
||||
$env,
|
||||
$flag,
|
||||
extractHttpStatusFromError,
|
||||
fetchWithRetry,
|
||||
parseStreamingJson,
|
||||
parseStreamingJsonThrottled,
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
import { ProviderHttpError } from "../errors";
|
||||
import type {
|
||||
Api,
|
||||
@@ -32,7 +39,6 @@ import { normalizeToolCallId, resolveCacheRetention } from "../utils";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector";
|
||||
import { armPreResponseTimeout, getStreamFirstEventTimeoutMs } from "../utils/idle-iterator";
|
||||
import { parseStreamingJson, parseStreamingJsonThrottled } from "../utils/json-parse";
|
||||
import { toolWireSchema } from "../utils/schema/wire";
|
||||
import { invalidateAwsCredentialCache, resolveAwsCredentials } from "./aws-credentials";
|
||||
import { decodeEventStream } from "./aws-eventstream";
|
||||
|
||||
@@ -15,6 +15,8 @@ import {
|
||||
isRetryableError,
|
||||
isUnexpectedSocketCloseMessage,
|
||||
logger,
|
||||
parseJsonWithRepair,
|
||||
parseStreamingJsonThrottled,
|
||||
readSseEvents,
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
import { isUsageLimitError } from "../rate-limit-utils";
|
||||
@@ -51,7 +53,6 @@ import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { isFoundryEnabled } from "../utils/foundry";
|
||||
import { finalizeErrorMessage, type RawHttpRequestDump, rewriteCopilotError } from "../utils/http-inspector";
|
||||
import { getStreamFirstEventTimeoutMs, getStreamIdleTimeoutMs, iterateWithIdleTimeout } from "../utils/idle-iterator";
|
||||
import { parseJsonWithRepair, parseStreamingJsonThrottled } from "../utils/json-parse";
|
||||
import { notifyProviderResponse } from "../utils/provider-response";
|
||||
import { isCopilotTransientModelError } from "../utils/retry";
|
||||
import { COMBINATOR_KEYS, NO_STRICT, toolWireSchema } from "../utils/schema";
|
||||
|
||||
@@ -102,7 +102,13 @@ import {
|
||||
WriteSuccessSchema,
|
||||
} from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb";
|
||||
import { calculateCost } from "@oh-my-pi/pi-catalog/models";
|
||||
import { $env, extractHttpStatusFromError, sanitizeText } from "@oh-my-pi/pi-utils";
|
||||
import {
|
||||
$env,
|
||||
extractHttpStatusFromError,
|
||||
parseJsonWithRepair,
|
||||
parseStreamingJson,
|
||||
sanitizeText,
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
import type {
|
||||
Api,
|
||||
AssistantMessage,
|
||||
@@ -126,7 +132,6 @@ import type {
|
||||
import { normalizeSystemPrompts } from "../utils";
|
||||
import { deterministicUuid } from "../utils/deterministic-id";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { parseJsonWithRepair, parseStreamingJson } from "../utils/json-parse";
|
||||
import { connectProxiedSocket, getProxyForProvider, shouldBypassProxy } from "../utils/proxy";
|
||||
import { createRequestDebugSession, isRequestDebugEnabled, type RequestDebugResponseLog } from "../utils/request-debug";
|
||||
import { formatErrorMessageWithRetryAfter } from "../utils/retry-after";
|
||||
|
||||
@@ -28,7 +28,7 @@ import {
|
||||
StopReason,
|
||||
} from "@oh-my-pi/pi-catalog/discovery/devin-gen/exa/codeium_common_pb/codeium_common_pb";
|
||||
import { calculateCost } from "@oh-my-pi/pi-catalog/models";
|
||||
import { extractHttpStatusFromError, logger } from "@oh-my-pi/pi-utils";
|
||||
import { extractHttpStatusFromError, logger, parseStreamingJson } from "@oh-my-pi/pi-utils";
|
||||
import type {
|
||||
Api,
|
||||
AssistantMessage,
|
||||
@@ -44,7 +44,6 @@ import type {
|
||||
} from "../types";
|
||||
import { deterministicUuid } from "../utils/deterministic-id";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { parseStreamingJson } from "../utils/json-parse";
|
||||
import { formatErrorMessageWithRetryAfter } from "../utils/retry-after";
|
||||
import { toolWireSchema } from "../utils/schema/wire";
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { extractHttpStatusFromError, fetchWithRetry } from "@oh-my-pi/pi-utils";
|
||||
import { extractHttpStatusFromError, fetchWithRetry, parseStreamingJson } from "@oh-my-pi/pi-utils";
|
||||
import { ProviderHttpError } from "../errors";
|
||||
import { getEnvApiKey } from "../stream";
|
||||
import type {
|
||||
@@ -22,7 +22,6 @@ import {
|
||||
getOpenAIStreamFirstEventTimeoutMs,
|
||||
getOpenAIStreamIdleTimeoutMs,
|
||||
} from "../utils/idle-iterator";
|
||||
import { parseStreamingJson } from "../utils/json-parse";
|
||||
import { toolWireSchema } from "../utils/schema/wire";
|
||||
import {
|
||||
getStreamMarkupHealingPattern,
|
||||
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
extractHttpStatusFromError,
|
||||
fetchWithRetry,
|
||||
logger,
|
||||
parseStreamingJson,
|
||||
readSseJson,
|
||||
structuredCloneJSON,
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
@@ -51,7 +52,6 @@ import {
|
||||
getOpenAIStreamIdleTimeoutMs,
|
||||
iterateWithIdleTimeout,
|
||||
} from "../utils/idle-iterator";
|
||||
import { parseStreamingJson } from "../utils/json-parse";
|
||||
import { createRequestDebugSession, isRequestDebugEnabled, type RequestDebugResponseLog } from "../utils/request-debug";
|
||||
import { adaptSchemaForStrict, NO_STRICT, sanitizeSchemaForOpenAIResponses, toolWireSchema } from "../utils/schema";
|
||||
import { notifyRawSseEvent } from "../utils/sse-debug";
|
||||
|
||||
@@ -3,7 +3,7 @@ import { isKimiModelId } from "@oh-my-pi/pi-catalog/identity";
|
||||
import { resolveWireModelId } from "@oh-my-pi/pi-catalog/model-thinking";
|
||||
import { calculateCost } from "@oh-my-pi/pi-catalog/models";
|
||||
import type { ResolvedOpenAICompat } from "@oh-my-pi/pi-catalog/types";
|
||||
import { $env, extractHttpStatusFromError } from "@oh-my-pi/pi-utils";
|
||||
import { $env, extractHttpStatusFromError, parseStreamingJson, parseStreamingJsonThrottled } from "@oh-my-pi/pi-utils";
|
||||
import { getKimiCommonHeaders } from "../registry/oauth/kimi";
|
||||
import { getEnvApiKey } from "../stream";
|
||||
import type {
|
||||
@@ -36,7 +36,6 @@ import {
|
||||
iterateWithIdleTimeout,
|
||||
iterateWithTerminalGrace,
|
||||
} from "../utils/idle-iterator";
|
||||
import { parseStreamingJson, parseStreamingJsonThrottled } from "../utils/json-parse";
|
||||
import { OpenAIHttpError, postOpenAIStream } from "../utils/openai-http";
|
||||
import { notifyProviderResponse } from "../utils/provider-response";
|
||||
import { callWithCopilotModelRetry } from "../utils/retry";
|
||||
|
||||
@@ -19,7 +19,14 @@ import {
|
||||
hasCoreWeaveProjectHeader,
|
||||
} from "@oh-my-pi/pi-catalog/wire/coreweave";
|
||||
import { parseGitHubCopilotApiKey } from "@oh-my-pi/pi-catalog/wire/github-copilot";
|
||||
import { $env, extractHttpStatusFromError, logger, structuredCloneJSON } from "@oh-my-pi/pi-utils";
|
||||
import {
|
||||
$env,
|
||||
extractHttpStatusFromError,
|
||||
logger,
|
||||
parseStreamingJson,
|
||||
parseStreamingJsonThrottled,
|
||||
structuredCloneJSON,
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
import {
|
||||
type Api,
|
||||
type AssistantMessage,
|
||||
@@ -54,7 +61,6 @@ import {
|
||||
} from "../utils";
|
||||
import type { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import type { CapturedHttpErrorResponse } from "../utils/http-inspector";
|
||||
import { parseStreamingJson, parseStreamingJsonThrottled } from "../utils/json-parse";
|
||||
import { getOpenRouterHeaders } from "../utils/openrouter-headers";
|
||||
import { isForcedToolChoice } from "../utils/tool-choice";
|
||||
import {
|
||||
|
||||
@@ -107,7 +107,6 @@ export const BUNDLED_PI_REGISTRY_KEYS: ReadonlySet<string> = new Set([
|
||||
"@oh-my-pi/pi-ai/utils/google-validation",
|
||||
"@oh-my-pi/pi-ai/utils/http-inspector",
|
||||
"@oh-my-pi/pi-ai/utils/idle-iterator",
|
||||
"@oh-my-pi/pi-ai/utils/json-parse",
|
||||
"@oh-my-pi/pi-ai/utils/openai-http",
|
||||
"@oh-my-pi/pi-ai/utils/openrouter-headers",
|
||||
"@oh-my-pi/pi-ai/utils/overflow",
|
||||
|
||||
@@ -131,7 +131,6 @@ import * as bundledPiAiUtilsGoogleValidation from "@oh-my-pi/pi-ai/utils/google-
|
||||
import * as bundledPiAiUtilsHarmonyLeak from "@oh-my-pi/pi-ai/utils/harmony-leak";
|
||||
import * as bundledPiAiUtilsHttpInspector from "@oh-my-pi/pi-ai/utils/http-inspector";
|
||||
import * as bundledPiAiUtilsIdleIterator from "@oh-my-pi/pi-ai/utils/idle-iterator";
|
||||
import * as bundledPiAiUtilsJsonParse from "@oh-my-pi/pi-ai/utils/json-parse";
|
||||
import * as bundledPiAiUtilsOpenaiHttp from "@oh-my-pi/pi-ai/utils/openai-http";
|
||||
import * as bundledPiAiUtilsOpenrouterHeaders from "@oh-my-pi/pi-ai/utils/openrouter-headers";
|
||||
import * as bundledPiAiUtilsOverflow from "@oh-my-pi/pi-ai/utils/overflow";
|
||||
@@ -1207,7 +1206,6 @@ export const BUNDLED_PI_REGISTRY: Readonly<Record<string, Readonly<Record<string
|
||||
Record<string, unknown>
|
||||
>,
|
||||
"@oh-my-pi/pi-ai/utils/idle-iterator": bundledPiAiUtilsIdleIterator as unknown as Readonly<Record<string, unknown>>,
|
||||
"@oh-my-pi/pi-ai/utils/json-parse": bundledPiAiUtilsJsonParse as unknown as Readonly<Record<string, unknown>>,
|
||||
"@oh-my-pi/pi-ai/utils/openai-http": bundledPiAiUtilsOpenaiHttp as unknown as Readonly<Record<string, unknown>>,
|
||||
"@oh-my-pi/pi-ai/utils/openrouter-headers": bundledPiAiUtilsOpenrouterHeaders as unknown as Readonly<
|
||||
Record<string, unknown>
|
||||
|
||||
@@ -1690,7 +1690,8 @@ export class InputController {
|
||||
// that state — inform the user instead of silently flipping the
|
||||
// persisted value. When thinking is on, the toggle works normally
|
||||
// even if blocks are already hidden (user may want to show them).
|
||||
const thinkingOff = ((this.ctx.viewSession ?? this.ctx.session)?.thinkingLevel ?? ThinkingLevel.Off) === ThinkingLevel.Off;
|
||||
const thinkingOff =
|
||||
((this.ctx.viewSession ?? this.ctx.session)?.thinkingLevel ?? ThinkingLevel.Off) === ThinkingLevel.Off;
|
||||
if (thinkingOff) {
|
||||
this.ctx.showStatus("Thinking is off — enable thinking to show blocks");
|
||||
return;
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { parseStreamingJson } from "@oh-my-pi/pi-ai/utils/json-parse";
|
||||
import { parseStreamingJson } from "@oh-my-pi/pi-utils";
|
||||
import { nextStep, STREAMING_REVEAL_FRAME_MS } from "./streaming-reveal";
|
||||
|
||||
/** Minimal component surface the reveal pushes frames into. */
|
||||
|
||||
@@ -9,12 +9,7 @@ interface RenderableBlock {
|
||||
}
|
||||
|
||||
function isRenderableBlock(value: unknown): value is RenderableBlock {
|
||||
return (
|
||||
value !== null &&
|
||||
typeof value === "object" &&
|
||||
"render" in value &&
|
||||
typeof value.render === "function"
|
||||
);
|
||||
return value !== null && typeof value === "object" && "render" in value && typeof value.render === "function";
|
||||
}
|
||||
|
||||
function renderPresentedBlocks(value: unknown): string {
|
||||
|
||||
@@ -1,9 +1,15 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Added
|
||||
|
||||
- Added a relaxed JSON parser that supports single-quoted strings, unquoted keys, and comments
|
||||
- Added `parseStreamingJson` for robust parsing of truncated or malformed streaming JSON
|
||||
- Added `parseStreamingJsonThrottled` for efficient processing of incremental streaming updates
|
||||
|
||||
### Changed
|
||||
|
||||
- Improved streaming SSE JSON processing to gracefully handle common malformed tail events
|
||||
- Increased the EBUSY retry delay from 25ms to 50ms (40 retries × 50ms = 2s total window, up from 1s). Windows can hold file locks on SQLite databases for up to ~1.5s after `close()`, and the previous 1-second window was too short for some test cleanup scenarios — `settings-manager.test.ts` and `sdk-credential-disabled-bridge.test.ts` still failed with EBUSY even when using `removeSyncWithRetries`.
|
||||
|
||||
## [16.1.8] - 2026-06-20
|
||||
@@ -196,4 +202,4 @@
|
||||
|
||||
### Added
|
||||
|
||||
- Added an XDG-aware tiny-title model cache directory helper for coding-agent local title models.
|
||||
- Added an XDG-aware tiny-title model cache directory helper for coding-agent local title models.
|
||||
@@ -9,6 +9,7 @@ export * from "./frontmatter";
|
||||
export * from "./fs-error";
|
||||
export * from "./glob";
|
||||
export * from "./json";
|
||||
export * from "./json-parse";
|
||||
export * as logger from "./logger";
|
||||
export * from "./loop-phase";
|
||||
export * from "./mermaid-ascii";
|
||||
|
||||
+12
-245
@@ -1,6 +1,7 @@
|
||||
const trailingEvents = new WeakSet<ServerSentEvent>();
|
||||
|
||||
import { abortableSource } from "./abortable";
|
||||
import { parseStreamingJson } from "./json-parse";
|
||||
|
||||
const LF = 0x0a;
|
||||
type JsonlChunkResult = {
|
||||
@@ -212,250 +213,16 @@ function notifySseEventObserver(observer: SseEventObserver | undefined, event: S
|
||||
}
|
||||
}
|
||||
|
||||
interface JsonLevel {
|
||||
type: "object" | "array";
|
||||
state?: "expect_key" | "expect_colon" | "expect_value" | "expect_comma_or_close";
|
||||
}
|
||||
|
||||
interface ScanResult {
|
||||
stack: JsonLevel[];
|
||||
inString: boolean;
|
||||
escaped: boolean;
|
||||
}
|
||||
|
||||
function scanJsonState(trimmed: string): ScanResult | null {
|
||||
const stack: JsonLevel[] = [];
|
||||
let inString = false;
|
||||
let escaped = false;
|
||||
|
||||
for (let i = 0; i < trimmed.length; i++) {
|
||||
const char = trimmed[i];
|
||||
if (escaped) {
|
||||
escaped = false;
|
||||
continue;
|
||||
}
|
||||
if (char === "\\") {
|
||||
escaped = true;
|
||||
continue;
|
||||
}
|
||||
if (char === '"') {
|
||||
inString = !inString;
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (inString) {
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
} else {
|
||||
if (top.type === "object" && top.state === "expect_key") {
|
||||
top.state = "expect_colon";
|
||||
}
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (!inString) {
|
||||
if (char === "{") {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
stack.push({ type: "object", state: "expect_key" });
|
||||
} else if (char === "[") {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
stack.push({ type: "array", state: "expect_value" });
|
||||
} else if (char === "}") {
|
||||
if (stack.length === 0 || stack[stack.length - 1].type !== "object") {
|
||||
return null;
|
||||
}
|
||||
stack.pop();
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
} else if (char === "]") {
|
||||
if (stack.length === 0 || stack[stack.length - 1].type !== "array") {
|
||||
return null;
|
||||
}
|
||||
stack.pop();
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
} else if (char === ":") {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.type === "object" && top.state === "expect_colon") {
|
||||
top.state = "expect_value";
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
} else if (char === ",") {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_comma_or_close") {
|
||||
top.state = top.type === "object" ? "expect_key" : "expect_value";
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
} else if (char !== " " && char !== "\t" && char !== "\n" && char !== "\r") {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (inString) {
|
||||
if (stack.length > 0) {
|
||||
const top = stack[stack.length - 1];
|
||||
if (top.type === "object") {
|
||||
if (top.state === "expect_key") {
|
||||
top.state = "expect_colon";
|
||||
} else if (top.state === "expect_value") {
|
||||
top.state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return { stack, inString, escaped };
|
||||
}
|
||||
|
||||
function hasValidPredecessorForComma(trimmed: string): boolean {
|
||||
if (!trimmed.endsWith(",")) return false;
|
||||
const beforeComma = trimmed.slice(0, -1).trim();
|
||||
const scan = scanJsonState(beforeComma);
|
||||
if (!scan || scan.stack.length === 0 || scan.inString) return false;
|
||||
return scan.stack[scan.stack.length - 1].state === "expect_comma_or_close";
|
||||
}
|
||||
|
||||
function getClosingSuffix(data: string): string | null {
|
||||
let trimmed = data.trim();
|
||||
if (!(trimmed.startsWith("{") || trimmed.startsWith("["))) {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (hasValidPredecessorForComma(trimmed)) {
|
||||
trimmed = trimmed.slice(0, -1).trim();
|
||||
}
|
||||
|
||||
const scan = scanJsonState(trimmed);
|
||||
if (!scan) return null;
|
||||
|
||||
const { stack, inString, escaped } = scan;
|
||||
let suffix = "";
|
||||
|
||||
if (inString) {
|
||||
if (escaped) {
|
||||
suffix += "\\";
|
||||
} else {
|
||||
const match = trimmed.match(/(\\+)u([0-9a-fA-F]{0,3})$/);
|
||||
if (match) {
|
||||
const backslashes = match[1].length;
|
||||
if (backslashes % 2 === 1) {
|
||||
const digits = match[2].length;
|
||||
suffix += "0".repeat(4 - digits);
|
||||
}
|
||||
}
|
||||
}
|
||||
suffix += '"';
|
||||
}
|
||||
|
||||
let literalCompletion = "";
|
||||
if (!inString) {
|
||||
const match = trimmed.match(/(t|tr|tru|f|fa|fal|fals|n|nu|nul)$/i);
|
||||
if (match) {
|
||||
const partial = match[0].toLowerCase();
|
||||
const completions: Record<string, string> = {
|
||||
t: "rue",
|
||||
tr: "ue",
|
||||
tru: "e",
|
||||
f: "alse",
|
||||
fa: "lse",
|
||||
fal: "se",
|
||||
fals: "e",
|
||||
n: "ull",
|
||||
nu: "ll",
|
||||
nul: "l",
|
||||
};
|
||||
literalCompletion = completions[partial] ?? "";
|
||||
}
|
||||
}
|
||||
|
||||
suffix += literalCompletion;
|
||||
|
||||
for (let i = stack.length - 1; i >= 0; i--) {
|
||||
const level = stack[i];
|
||||
if (level.type === "array") {
|
||||
suffix += "]";
|
||||
} else {
|
||||
let state = level.state;
|
||||
if (state === "expect_value") {
|
||||
const lower = trimmed.toLowerCase();
|
||||
if (
|
||||
literalCompletion !== "" ||
|
||||
lower.endsWith("true") ||
|
||||
lower.endsWith("false") ||
|
||||
lower.endsWith("null") ||
|
||||
/\d$/.test(trimmed) ||
|
||||
trimmed.endsWith('"') ||
|
||||
trimmed.endsWith("}") ||
|
||||
trimmed.endsWith("]")
|
||||
) {
|
||||
state = "expect_comma_or_close";
|
||||
}
|
||||
}
|
||||
|
||||
if (state === "expect_key") {
|
||||
suffix += "}";
|
||||
} else if (state === "expect_colon") {
|
||||
suffix += ":null}";
|
||||
} else if (state === "expect_value") {
|
||||
suffix += "null}";
|
||||
} else if (state === "expect_comma_or_close") {
|
||||
suffix += "}";
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return suffix;
|
||||
}
|
||||
|
||||
function isJsonTruncated(data: string): boolean {
|
||||
let trimmed = data.trim();
|
||||
if (hasValidPredecessorForComma(trimmed)) {
|
||||
trimmed = trimmed.slice(0, -1).trim();
|
||||
}
|
||||
const suffix = getClosingSuffix(data);
|
||||
if (suffix === null || suffix === "") return false;
|
||||
|
||||
try {
|
||||
JSON.parse(trimmed + suffix);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
function isRecoverableTrailingJson(data: string): boolean {
|
||||
const first = data.trimStart()[0];
|
||||
if (first !== "{" && first !== "[") return false;
|
||||
// Best-effort relaxed recovery via the shared streaming JSON parser: a
|
||||
// container-shaped final event that fails strict `JSON.parse` is treated as a
|
||||
// cut-off (or lightly malformed) stream tail and ends iteration cleanly instead
|
||||
// of throwing. Non-container final events (plain-text errors, bare scalars) are
|
||||
// not recoverable and still surface as a SyntaxError.
|
||||
const recovered = parseStreamingJson<unknown>(data);
|
||||
return typeof recovered === "object" && recovered !== null;
|
||||
}
|
||||
|
||||
export async function* readSseJson<T>(
|
||||
@@ -474,7 +241,7 @@ export async function* readSseJson<T>(
|
||||
try {
|
||||
yield JSON.parse(data) as T;
|
||||
} catch (err) {
|
||||
if (err instanceof SyntaxError && isTrailing && isJsonTruncated(data)) {
|
||||
if (err instanceof SyntaxError && isTrailing && isRecoverableTrailingJson(data)) {
|
||||
return;
|
||||
}
|
||||
throw err;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { parseJsonWithRepair, parseStreamingJson, repairJson } from "@oh-my-pi/pi-ai/utils/json-parse";
|
||||
import { parseJsonWithRepair, parseStreamingJson, repairJson } from "@oh-my-pi/pi-utils/json-parse";
|
||||
|
||||
describe("JSON repair", () => {
|
||||
it("leaves valid string escapes unchanged", () => {
|
||||
+1
-1
@@ -3,7 +3,7 @@ import {
|
||||
parseStreamingJson,
|
||||
parseStreamingJsonThrottled,
|
||||
STREAMING_JSON_PARSE_MIN_GROWTH,
|
||||
} from "@oh-my-pi/pi-ai/utils/json-parse";
|
||||
} from "@oh-my-pi/pi-utils/json-parse";
|
||||
|
||||
describe("parseStreamingJsonThrottled (F5)", () => {
|
||||
it("parses the first non-empty buffer even when growth is below the threshold", () => {
|
||||
@@ -296,14 +296,10 @@ describe("readSseJson", () => {
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event is complete but malformed JSON", async () => {
|
||||
const testCases = [
|
||||
'data: {"b":2,}', // balanced but malformed trailing comma
|
||||
'data: {"b":,', // invalid predecessor before comma
|
||||
"data: [,,", // invalid predecessor before comma
|
||||
'data: {"b",', // invalid predecessor after key
|
||||
'data: {"s":"\\u12,', // comma inside unterminated string
|
||||
];
|
||||
it("throws SyntaxError when a final event is not JSON-container-shaped", async () => {
|
||||
// Non-object/array final events are not recoverable as a truncated stream tail
|
||||
// and still surface as errors (e.g. provider error text, bare scalars).
|
||||
const testCases = ["data: Internal Server Error", 'data: "an unterminated string', "data: 42 then junk"];
|
||||
for (const dataChunk of testCases) {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode(dataChunk)];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
@@ -317,109 +313,32 @@ describe("readSseJson", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event is plain text", async () => {
|
||||
const chunks = [
|
||||
encoder.encode('data: {"a":1}\n\n'),
|
||||
encoder.encode("data: Internal Server Error"), // plain text (not JSON)
|
||||
it("stops cleanly on a container-shaped final event that fails strict parse", async () => {
|
||||
// Lenient recovery: any object/array-shaped final event JSON.parse rejects is
|
||||
// treated as a cut-off or lightly malformed stream tail and ends iteration after
|
||||
// the last valid event, rather than throwing.
|
||||
const testCases = [
|
||||
'data: {"b":2,}', // trailing comma
|
||||
"data: [{]", // mismatched closer
|
||||
'data: {"b" 2}', // missing colon
|
||||
"data: {unterminated}", // bareword body
|
||||
'data: {"b": true garbage', // trailing garbage after a value
|
||||
'data: {"b":1 "c":2', // missing comma
|
||||
'data: {"b": ]', // mismatched closer
|
||||
'data: {"b": @', // invalid character
|
||||
];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
for (const dataChunk of testCases) {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode(dataChunk)];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event is malformed with missing colons but balanced braces", async () => {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b" 2}')];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event is malformed with mismatched brackets/braces", async () => {
|
||||
const chunks = [
|
||||
encoder.encode('data: {"a":1}\n\n'),
|
||||
encoder.encode("data: [{]"), // mismatched closer
|
||||
];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event is malformed with syntax errors after values", async () => {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b": true garbage')];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event has invalid characters", async () => {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b": @')];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event has mismatched closers", async () => {
|
||||
const chunks = [
|
||||
encoder.encode('data: {"a":1}\n\n'),
|
||||
encoder.encode('data: {"b": ]'), // mismatched closer
|
||||
];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event has unbalanced braces but internal syntax errors", async () => {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b":1 "c":2')];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
});
|
||||
|
||||
it("throws SyntaxError when a final event has unbalanced braces but has complete invalid text inside", async () => {
|
||||
const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode("data: {unterminated}")];
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
for (const chunk of chunks) controller.enqueue(chunk);
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError);
|
||||
const output = await collectAsync(readSseJson(stream));
|
||||
expect(output).toEqual([{ a: 1 }]);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -31,7 +31,8 @@ test("removeSyncWithRetries outlasts a transient 1.5s Windows EBUSY lock", async
|
||||
|
||||
try {
|
||||
// Dynamic import is required so the node:fs mock is installed before temp.ts binds it.
|
||||
const { removeSyncWithRetries } = await import("../src/temp?ebusy-window-test");
|
||||
const tempModulePath = "../src/temp?ebusy-window-test";
|
||||
const { removeSyncWithRetries } = (await import(tempModulePath)) as typeof import("../src/temp");
|
||||
removeSyncWithRetries("/tmp/pr3348");
|
||||
expect(calls).toBe(31);
|
||||
expect(elapsedMs).toBe(1500);
|
||||
|
||||
Reference in New Issue
Block a user