From da7b1e2b214f380285ddaace6202fb251bfa43cf Mon Sep 17 00:00:00 2001 From: oyi77 Date: Tue, 14 Jul 2026 01:20:26 +0700 Subject: [PATCH] perf: hoist Date.now, hoist hasActiveDeltaValue, avoid per-chunk buffer.split in SSE stream - Hoist to transform entry, remove 2 inner declarations that shadowed the outer one (stream.ts) - Hoist from inline closure to module-level function to avoid allocation per chunk (stream.ts) - Replace unconditional per chunk with -gated split to avoid allocating array when no embedded newline (stream.ts streamHelpers.ts) - Guard to avoid allocating when already at/above limit --- open-sse/utils/stream.ts | 44 +++++++++++++++++++-------------- open-sse/utils/streamHelpers.ts | 4 +-- 2 files changed, 27 insertions(+), 21 deletions(-) diff --git a/open-sse/utils/stream.ts b/open-sse/utils/stream.ts index 581f1f6342..fb4bd85456 100644 --- a/open-sse/utils/stream.ts +++ b/open-sse/utils/stream.ts @@ -181,6 +181,12 @@ const STREAM_SUMMARY_TEXT_LIMIT = 64 * 1024; function appendBoundedText(current: string, next: string): string { if (!next) return current; + // Avoid allocating `current + next` when already at/above limit — slide the window instead. + if (current.length >= STREAM_SUMMARY_TEXT_LIMIT) { + const keep = + STREAM_SUMMARY_TEXT_LIMIT > next.length ? STREAM_SUMMARY_TEXT_LIMIT - next.length : 0; + return current.slice(-keep) + next; + } const combined = current + next; if (combined.length <= STREAM_SUMMARY_TEXT_LIMIT) return combined; return combined.slice(-STREAM_SUMMARY_TEXT_LIMIT); @@ -191,6 +197,16 @@ function parseTextualToolCallFromContent(text: unknown): { name: string; args: u return candidate?.kind === "complete" ? { name: candidate.name, args: candidate.args } : null; } +/** Per-chunk recursive check for meaningful delta content. Hoisted to avoid closure re-allocation in hot-path. */ +function hasActiveDeltaValue(value: unknown): boolean { + if (typeof value === "string") return value.length > 0; + if (Array.isArray(value)) return value.some((entry) => hasActiveDeltaValue(entry)); + if (value && typeof value === "object") { + return Object.values(value).some((entry) => hasActiveDeltaValue(entry)); + } + return value !== null && value !== undefined; +} + function containsTextualToolCallCandidate(text: unknown): boolean { return parseTextualToolCallCandidate(text) !== null; } @@ -1153,13 +1169,14 @@ export function createSSEStream(options: StreamOptions = {}) { transform(chunk, controller) { if (streamTimedOut) return; - lastChunkTime = Date.now(); + const now = Date.now(); + lastChunkTime = now; const text = decoder.decode(chunk, { stream: true }); buffer += text; reqLogger?.appendProviderChunk?.(text); - - const lines = buffer.split("\n"); - buffer = lines.pop() || ""; + const nlIdx = buffer.lastIndexOf("\n"); + const lines = nlIdx >= 0 ? buffer.slice(0, nlIdx).split("\n") : []; + if (nlIdx >= 0) buffer = buffer.slice(nlIdx + 1); for (const line of multilineSseDataLineNormalizer.normalize(lines)) { const trimmed = line.trim(); @@ -1262,15 +1279,6 @@ export function createSSEStream(options: StreamOptions = {}) { // bootstrap chunk (assistant role + empty content) before emitting proper // `response.*` events. That chunk is invalid on /v1/responses and breaks strict // clients like OpenCode, so drop it only for Responses-native consumers. - const hasActiveDeltaValue = (value: unknown): boolean => { - if (typeof value === "string") return value.length > 0; - if (Array.isArray(value)) - return value.some((entry) => hasActiveDeltaValue(entry)); - if (value && typeof value === "object") { - return Object.values(value).some((entry) => hasActiveDeltaValue(entry)); - } - return value !== null && value !== undefined; - }; const isEmptyAssistantBootstrapChunkForResponsesClient = clientExpectsResponsesStream && @@ -1738,7 +1746,7 @@ export function createSSEStream(options: StreamOptions = {}) { // T18: Track if we saw tool calls & accumulate for call log if (delta?.tool_calls && delta.tool_calls.length > 0) { passthroughHasToolCalls = true; - lastToolCallChunkTime = Date.now(); + lastToolCallChunkTime = now; for (const tc of delta.tool_calls) { // Note: sanitizeStreamingChunk above already coerces non-string // tool_call IDs, but this defensive check catches edge cases @@ -1787,7 +1795,6 @@ export function createSSEStream(options: StreamOptions = {}) { const lastChunkTs = lastToolCallChunkTime; if (toolTs || lastChunkTs) { contentAfterToolSeen = true; - const now = Date.now(); try { recordToolLatency( provider || "unknown", @@ -1824,7 +1831,7 @@ export function createSSEStream(options: StreamOptions = {}) { const isFinishChunk = parsed.choices?.[0]?.finish_reason; if (isFinishChunk && passthroughHasToolCalls) { - toolFinishTime = Date.now(); + toolFinishTime = now; try { markToolFinish(sessionId); } catch {} @@ -1949,10 +1956,10 @@ export function createSSEStream(options: StreamOptions = {}) { } if (parsed.choices?.[0]?.delta?.tool_calls) { - lastToolCallChunkTime = Date.now(); + lastToolCallChunkTime = now; } if (parsed.choices?.[0]?.finish_reason === "tool_calls") { - toolFinishTime = Date.now(); + toolFinishTime = now; try { markToolFinish(sessionId); } catch {} @@ -2063,7 +2070,6 @@ export function createSSEStream(options: StreamOptions = {}) { const lastChunkTs = lastToolCallChunkTime; if (toolTs || lastChunkTs) { contentAfterToolSeen = true; - const now = Date.now(); try { recordToolLatency( provider || "unknown", diff --git a/open-sse/utils/streamHelpers.ts b/open-sse/utils/streamHelpers.ts index 1f4b5f9198..3b7b95d0b6 100644 --- a/open-sse/utils/streamHelpers.ts +++ b/open-sse/utils/streamHelpers.ts @@ -432,8 +432,8 @@ export function fixInvalidId(parsed: Record): boolean { // Remove null perf_metrics from usage (common across formats) function cleanPerfMetrics(data: unknown): unknown { if (isRecord(data) && isRecord(data.usage) && data.usage.perf_metrics === null) { - const { perf_metrics, ...usageWithoutPerf } = data.usage; - return { ...data, usage: usageWithoutPerf }; + // Mutate in-place to avoid spread copy per chunk — data is ephemeral, used only for serialization. + delete data.usage.perf_metrics; } return data; }