diff options
Diffstat (limited to 'packages/openai-stream/src/parse-sse.ts')
| -rw-r--r-- | packages/openai-stream/src/parse-sse.ts | 222 |
1 files changed, 111 insertions, 111 deletions
diff --git a/packages/openai-stream/src/parse-sse.ts b/packages/openai-stream/src/parse-sse.ts index cfeb5b0..112a418 100644 --- a/packages/openai-stream/src/parse-sse.ts +++ b/packages/openai-stream/src/parse-sse.ts @@ -1,130 +1,130 @@ import type { ProviderEvent } from "@dispatch/kernel"; interface ToolCallAccumulator { - id: string; - name: string; - arguments: string; + id: string; + name: string; + arguments: string; } interface SSEChunkDelta { - content?: string; - reasoning_content?: string; - tool_calls?: Array<{ - index: number; - id?: string; - function?: { name?: string; arguments?: string }; - }>; + content?: string; + reasoning_content?: string; + tool_calls?: Array<{ + index: number; + id?: string; + function?: { name?: string; arguments?: string }; + }>; } interface SSEChunkChoice { - delta: SSEChunkDelta; - finish_reason?: string | null; - index: number; + delta: SSEChunkDelta; + finish_reason?: string | null; + index: number; } interface SSEChunkUsageDetails { - cached_tokens?: number; + cached_tokens?: number; } interface SSEChunk { - id?: string; - choices?: SSEChunkChoice[]; - usage?: { - prompt_tokens?: number; - completion_tokens?: number; - cache_read_tokens?: number; - cache_write_tokens?: number; - prompt_tokens_details?: SSEChunkUsageDetails; - completion_tokens_details?: Record<string, unknown>; - }; + id?: string; + choices?: SSEChunkChoice[]; + usage?: { + prompt_tokens?: number; + completion_tokens?: number; + cache_read_tokens?: number; + cache_write_tokens?: number; + prompt_tokens_details?: SSEChunkUsageDetails; + completion_tokens_details?: Record<string, unknown>; + }; } export function parseSSELines(lines: readonly string[]): ProviderEvent[] { - const events: ProviderEvent[] = []; - const toolCalls = new Map<number, ToolCallAccumulator>(); - - for (const line of lines) { - const trimmed = line.trim(); - if (!trimmed.startsWith("data:")) continue; - - const data = trimmed.slice(5).trim(); - if (data === "[DONE]") break; - - let chunk: SSEChunk; - try { - chunk = JSON.parse(data) as SSEChunk; - } catch { - events.push({ type: "error", message: `Invalid JSON in SSE data: ${data}` }); - continue; - } - - if (chunk.choices) { - for (const choice of chunk.choices) { - const delta = choice.delta; - - if (delta.content) { - events.push({ type: "text-delta", delta: delta.content }); - } - - if (delta.reasoning_content) { - events.push({ type: "reasoning-delta", delta: delta.reasoning_content }); - } - - if (delta.tool_calls) { - for (const tc of delta.tool_calls) { - const existing = toolCalls.get(tc.index); - if (existing) { - if (tc.function?.arguments) { - existing.arguments += tc.function.arguments; - } - } else { - toolCalls.set(tc.index, { - id: tc.id ?? "", - name: tc.function?.name ?? "", - arguments: tc.function?.arguments ?? "", - }); - } - } - } - - if (choice.finish_reason) { - const sortedIndices = [...toolCalls.keys()].sort((a, b) => a - b); - for (const idx of sortedIndices) { - const acc = toolCalls.get(idx); - if (!acc) continue; - let input: unknown; - try { - input = JSON.parse(acc.arguments); - } catch { - input = acc.arguments; - } - events.push({ - type: "tool-call", - toolCallId: acc.id, - toolName: acc.name, - input, - }); - } - events.push({ type: "finish", reason: choice.finish_reason }); - } - } - } - - if (chunk.usage) { - const cacheRead = - chunk.usage.cache_read_tokens ?? chunk.usage.prompt_tokens_details?.cached_tokens; - const cacheWrite = chunk.usage.cache_write_tokens; - events.push({ - type: "usage", - usage: { - inputTokens: chunk.usage.prompt_tokens ?? 0, - outputTokens: chunk.usage.completion_tokens ?? 0, - ...(cacheRead !== undefined ? { cacheReadTokens: cacheRead } : {}), - ...(cacheWrite !== undefined ? { cacheWriteTokens: cacheWrite } : {}), - }, - }); - } - } - - return events; + const events: ProviderEvent[] = []; + const toolCalls = new Map<number, ToolCallAccumulator>(); + + for (const line of lines) { + const trimmed = line.trim(); + if (!trimmed.startsWith("data:")) continue; + + const data = trimmed.slice(5).trim(); + if (data === "[DONE]") break; + + let chunk: SSEChunk; + try { + chunk = JSON.parse(data) as SSEChunk; + } catch { + events.push({ type: "error", message: `Invalid JSON in SSE data: ${data}` }); + continue; + } + + if (chunk.choices) { + for (const choice of chunk.choices) { + const delta = choice.delta; + + if (delta.content) { + events.push({ type: "text-delta", delta: delta.content }); + } + + if (delta.reasoning_content) { + events.push({ type: "reasoning-delta", delta: delta.reasoning_content }); + } + + if (delta.tool_calls) { + for (const tc of delta.tool_calls) { + const existing = toolCalls.get(tc.index); + if (existing) { + if (tc.function?.arguments) { + existing.arguments += tc.function.arguments; + } + } else { + toolCalls.set(tc.index, { + id: tc.id ?? "", + name: tc.function?.name ?? "", + arguments: tc.function?.arguments ?? "", + }); + } + } + } + + if (choice.finish_reason) { + const sortedIndices = [...toolCalls.keys()].sort((a, b) => a - b); + for (const idx of sortedIndices) { + const acc = toolCalls.get(idx); + if (!acc) continue; + let input: unknown; + try { + input = JSON.parse(acc.arguments); + } catch { + input = acc.arguments; + } + events.push({ + type: "tool-call", + toolCallId: acc.id, + toolName: acc.name, + input, + }); + } + events.push({ type: "finish", reason: choice.finish_reason }); + } + } + } + + if (chunk.usage) { + const cacheRead = + chunk.usage.cache_read_tokens ?? chunk.usage.prompt_tokens_details?.cached_tokens; + const cacheWrite = chunk.usage.cache_write_tokens; + events.push({ + type: "usage", + usage: { + inputTokens: chunk.usage.prompt_tokens ?? 0, + outputTokens: chunk.usage.completion_tokens ?? 0, + ...(cacheRead !== undefined ? { cacheReadTokens: cacheRead } : {}), + ...(cacheWrite !== undefined ? { cacheWriteTokens: cacheWrite } : {}), + }, + }); + } + } + + return events; } |
