/** * Decode an SSE byte stream into event `data` payloads. Framing — chunk * reassembly, UTF-8/CRLF/BOM handling, comment and non-data field skipping, * multi-`data:` joining — is `eventsource-parser`'s; this module keeps only * the DeepSeek protocol: the literal `[DONE]` is yielded so the caller owns * final flushing, and EOF before it raises {@link LlmError}. Framing is * spec-strict: an event dispatches only on its blank-line terminator, so an * unterminated tail at EOF is truncation, not a flushable payload. * * @module dsh-llm-deepseek/sse */ import { EventSourceParserStream } from 'eventsource-parser/stream' import { LlmError } from '@deepseek-ai/dsh-llm' /** The terminal payload DeepSeek (and OpenAI) send after the last chunk. */ export const DONE = '[DONE]' /** * Parse an SSE byte stream into data payloads. Yields `[DONE]` as the final * value and returns; throws `LlmError('STREAM_CLOSED')` when the stream ends * without it (truncated response — the model call cannot be trusted). * @param stream - raw SSE bytes; reads may split anywhere, including mid-UTF-8 sequence. * @returns each event's data payload in arrival order, the `[DONE]` sentinel last. */ export async function* parseSse(stream: ReadableStream): AsyncGenerator { const events = stream .pipeThrough(new TextDecoderStream()) .pipeThrough(new EventSourceParserStream()) for await (const { data } of events) { yield data if (data === DONE) return } throw new LlmError('SSE stream ended without [DONE]', 'STREAM_CLOSED') }