1
/** SSE framing delegated to eventsource-parser; JSON errors remain provider failures. */3
import { EventSourceParserStream } from 'eventsource-parser/stream'4
import { LlmError } from '@deepseek-ai/dsh-llm'5
import { object } from './replay.ts'6
import { providerError } from './transport.ts'8
/** Decode complete SSE frames without treating an unterminated tail as an event.9
* @param body - provider response bytes.10
* @param activity - pulse the idle watchdog for events and heartbeat comments.11
* @returns JSON events, including message_stop; the translator owns completion.12
*/13
export async function* parseSse(body: ReadableStream<BufferSource>, activity: () => void): AsyncGenerator<Record<string, unknown>> {14
const events = body.pipeThrough(new TextDecoderStream()).pipeThrough(new EventSourceParserStream({ onComment: activity }))15
for await (const frame of events) {16
activity()17
let raw: unknown18
try { raw = JSON.parse(frame.data) } catch (_invalidSseJson) {19
throw new LlmError('DeepSeek Messages SSE contains invalid JSON', 'MALFORMED_RESPONSE')20
}21
const event = object(raw)22
if (typeof event.type !== 'string' || (frame.event !== undefined && frame.event !== event.type)) {23
throw new LlmError('DeepSeek Messages SSE event type mismatch', 'MALFORMED_RESPONSE')24
}25
if (event.type === 'error') throw providerError(event, undefined)26
yield event27
}28
}