refactor(llm-deepseek): replace hand-rolled SSE parser with eventsource-parser
Implements the approved simplification Agent Note: sse.ts now pipes the
response body through TextDecoderStream and EventSourceParserStream
(eventsource-parser/stream) and keeps only the DeepSeek protocol shim —
yield each event's data, terminate on [DONE], throw
LlmError('STREAM_CLOSED') on EOF without the sentinel. The SSE
spec-conformance tests are deleted; sse.spec.ts pins only the
[DONE]/STREAM_CLOSED/EOF contract, including the new spec-strict verdict
that an unterminated trailing event is truncation (the old parser
flushed it — a robustness nicety no real provider shape needs).
eventsource-parser@^3.1.0 becomes llm-deepseek's second runtime
dependency (already in the lockfile transitively via the MCP SDK).
Docs: the Agent Note moves proposed/ → implemented/ and is rewritten per
the lifecycle contract; the rejected NIH roll-up note's inbound links
follow. The twin-adapters note, dsh-llm LlmAdapter JSDoc (and its
type-equiv fences), cookbook, group/package READMEs, root AGENTS.md
layout line, sdk-helper comments, and the regenerated config catalog
drop the "hand-rolled fetch + SSE" claim in both languages; all eight
touched pairs re-recorded.
This commit is contained in:
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write
|
||||
README.md: e191f3fcd265a6ca9cec3a8dae5f730ce27accf1
|
||||
README.zh.md: 268096e5f1a145e8d5cf6469524d36fe48984617
|
||||
README.md: dc373f32ccb5e4ade6dde3128fee75c913fd99b2
|
||||
README.zh.md: 22edefb91dc78a10961f0bc9fa676614ace26a33
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
English | [中文](README.zh.md)
|
||||
|
||||
DeepSeek chat-completions adapter for the harness LLM seam: hand-rolled `fetch` + SSE translation from the official wire format (source of truth: the API docs — guides/thinking_mode, guides/tool_calls, api/create-chat-completion) into the `StreamChunk` protocol.
|
||||
DeepSeek chat-completions adapter for the harness LLM seam: direct `fetch` + SSE (framed by `eventsource-parser`) translating the official wire format (source of truth: the API docs — guides/thinking_mode, guides/tool_calls, api/create-chat-completion) into the `StreamChunk` protocol.
|
||||
|
||||
A second, library-backed implementation of the same seam exists in `@deepseek-ai/dsh-llm-pi-ai`. This package always owns the `deepseek` provider route; mounting a pi-ai profile with `provider: deepseek` in the same context throws `LlmError('DUPLICATE_ADAPTER')` by design.
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
[English](README.md) | 中文
|
||||
|
||||
harness LLM seam 的 DeepSeek chat-completions 适配器:手写 `fetch` + SSE,将官方协议格式(真源:API 文档 guides/thinking_mode、guides/tool_calls、api/create-chat-completion)转换为 `StreamChunk` 协议。
|
||||
harness LLM seam 的 DeepSeek chat-completions 适配器:直接 `fetch` + SSE(由 `eventsource-parser` 分帧),将官方协议格式(真源:API 文档 guides/thinking_mode、guides/tool_calls、api/create-chat-completion)转换为 `StreamChunk` 协议。
|
||||
|
||||
同一 seam 的第二个库支持实现位于 `@deepseek-ai/dsh-llm-pi-ai`。本包始终拥有 `deepseek` 提供方路由;在同一上下文中装载 `provider: deepseek` 的 pi-ai profile 会按设计抛出 `LlmError('DUPLICATE_ADAPTER')`。
|
||||
|
||||
|
||||
@@ -33,6 +33,7 @@
|
||||
"cordis": "^4.0.0-rc.7"
|
||||
},
|
||||
"dependencies": {
|
||||
"eventsource-parser": "^3.1.0",
|
||||
"schemastery": "^3.18.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
|
||||
@@ -20,7 +20,7 @@ import { parseSse } from './sse.ts'
|
||||
import { translate } from './translate.ts'
|
||||
import type { WireError } from './types.ts'
|
||||
|
||||
/** One optional model entry advertised by the hand-written adapter. */
|
||||
/** One optional model entry advertised by the direct-fetch adapter. */
|
||||
export interface DeepSeekCatalogModel {
|
||||
/** Wire model id accepted by the configured endpoint. */
|
||||
id: string
|
||||
|
||||
@@ -1,65 +1,33 @@
|
||||
/**
|
||||
* Decode an SSE byte stream into event `data` payloads. Network reads may split UTF-8 or lines;
|
||||
* CRLF, comments, non-data fields, and multi-data events are handled per SSE rules. The literal
|
||||
* `[DONE]` is yielded so the caller owns final flushing, and EOF before it raises {@link LlmError}.
|
||||
* 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.
|
||||
*
|
||||
* Minimal SSE (text/event-stream) parser for the chat-completions stream.
|
||||
* @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]'
|
||||
|
||||
/** Extract the joined data payload from one raw SSE event block. */
|
||||
function eventData(block: string): string | undefined {
|
||||
const data: string[] = []
|
||||
for (const rawLine of block.split('\n')) {
|
||||
const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine
|
||||
if (line.startsWith('data:')) {
|
||||
// The spec strips ONE leading space after the colon.
|
||||
data.push(line.startsWith('data: ') ? line.slice(6) : line.slice(5))
|
||||
}
|
||||
// Comments (':…') and other fields (event:, id:, retry:) are ignored.
|
||||
}
|
||||
if (data.length === 0) return undefined
|
||||
return data.join('\n')
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a byte stream into SSE data payloads. Yields `[DONE]` as the final
|
||||
* 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: AsyncIterable<Uint8Array>): AsyncGenerator<string> {
|
||||
const decoder = new TextDecoder()
|
||||
let buffer = ''
|
||||
|
||||
for await (const bytes of stream) {
|
||||
buffer += decoder.decode(bytes, { stream: true })
|
||||
// Events are separated by a blank line (\n\n; tolerate \r\n\r\n via the
|
||||
// per-line \r strip in eventData and a normalized split here).
|
||||
let boundary: number
|
||||
while ((boundary = buffer.search(/\r?\n\r?\n/)) !== -1) {
|
||||
const matched = /\r?\n\r?\n/.exec(buffer.slice(boundary))
|
||||
const block = buffer.slice(0, boundary)
|
||||
// matched cannot be null: search() just found the same pattern at 0.
|
||||
buffer = buffer.slice(boundary + (matched as RegExpExecArray)[0].length)
|
||||
const data = eventData(block)
|
||||
if (data === undefined) continue
|
||||
yield data
|
||||
if (data === DONE) return
|
||||
}
|
||||
}
|
||||
|
||||
// Flush any final un-terminated event (servers usually end with \n\n, but
|
||||
// a trailing block without one is still parseable).
|
||||
buffer += decoder.decode()
|
||||
const data = eventData(buffer)
|
||||
if (data !== undefined) {
|
||||
export async function* parseSse(stream: ReadableStream<BufferSource>): AsyncGenerator<string> {
|
||||
const events = stream
|
||||
.pipeThrough(new TextDecoderStream())
|
||||
.pipeThrough(new EventSourceParserStream())
|
||||
for await (const { data } of events) {
|
||||
yield data
|
||||
if (data === DONE) return
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import type { Config } from '@deepseek-ai/dsh-llm-deepseek'
|
||||
import { assemble, type AssembledResult } from './assemble.ts'
|
||||
|
||||
/**
|
||||
* Real-API e2e for the hand-rolled adapter: V4 Flash + V4 Pro across
|
||||
* Real-API e2e for the direct-fetch adapter: V4 Flash + V4 Pro across
|
||||
* thinking modes and both official effort levels. Key-gated — skips
|
||||
* entirely without $DEEPSEEK_API_KEY (see vitest.e2e.config.ts).
|
||||
*/
|
||||
|
||||
@@ -2,12 +2,21 @@ import { describe, expect, it } from 'vitest'
|
||||
import { LlmError } from '@deepseek-ai/dsh-llm'
|
||||
import { DONE, parseSse } from '../src/sse.ts'
|
||||
|
||||
/** Build a byte stream from string fragments (fragments = network reads). */
|
||||
async function* bytes(...fragments: (string | Uint8Array)[]): AsyncGenerator<Uint8Array> {
|
||||
/**
|
||||
* DeepSeek protocol contract only: the [DONE] sentinel and STREAM_CLOSED on
|
||||
* EOF without it. SSE framing (chunk splits, CRLF, multi-data joins, comments)
|
||||
* is eventsource-parser's contract, not re-proven here.
|
||||
*/
|
||||
|
||||
/** Build an SSE byte stream from string fragments (fragments = network reads). */
|
||||
function bytes(...fragments: string[]): ReadableStream<Uint8Array<ArrayBuffer>> {
|
||||
const encoder = new TextEncoder()
|
||||
for (const fragment of fragments) {
|
||||
yield typeof fragment === 'string' ? encoder.encode(fragment) : fragment
|
||||
}
|
||||
return new ReadableStream({
|
||||
start(controller) {
|
||||
for (const fragment of fragments) controller.enqueue(encoder.encode(fragment))
|
||||
controller.close()
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
async function collect(stream: AsyncIterable<string>): Promise<string[]> {
|
||||
@@ -17,57 +26,14 @@ async function collect(stream: AsyncIterable<string>): Promise<string[]> {
|
||||
}
|
||||
|
||||
describe('parseSse', () => {
|
||||
it('parses simple events and the DONE sentinel', async () => {
|
||||
it('yields event payloads and the DONE sentinel', async () => {
|
||||
const events = await collect(parseSse(bytes('data: {"a":1}\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('handles events split across reads at arbitrary positions', async () => {
|
||||
const events = await collect(parseSse(bytes('da', 'ta: {"a"', ':1}\n', '\ndata: [DO', 'NE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('handles multi-byte UTF-8 split across reads', async () => {
|
||||
const encoded = new TextEncoder().encode('data: {"text":"日本語"}\n\ndata: [DONE]\n\n')
|
||||
// Split inside the 3-byte sequence for 日.
|
||||
const splitAt = 16
|
||||
const events = await collect(parseSse(bytes(encoded.slice(0, splitAt), encoded.slice(splitAt))))
|
||||
expect(events).toEqual(['{"text":"日本語"}', DONE])
|
||||
})
|
||||
|
||||
it('tolerates CRLF line endings', async () => {
|
||||
const events = await collect(parseSse(bytes('data: {"a":1}\r\n\r\ndata: [DONE]\r\n\r\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('joins multi-data events with newlines (SSE spec)', async () => {
|
||||
const events = await collect(parseSse(bytes('data: line1\ndata: line2\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['line1\nline2', DONE])
|
||||
})
|
||||
|
||||
it('ignores comments and non-data fields', async () => {
|
||||
const events = await collect(parseSse(bytes(': keepalive\nevent: chunk\nid: 7\ndata: {"a":1}\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('skips blocks without data fields', async () => {
|
||||
const events = await collect(parseSse(bytes(': ping\n\ndata: {"a":1}\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('preserves data lines without the optional space', async () => {
|
||||
const events = await collect(parseSse(bytes('data:{"a":1}\n\ndata:[DONE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('parses several events from one read', async () => {
|
||||
const events = await collect(parseSse(bytes('data: 1\n\ndata: 2\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['1', '2', DONE])
|
||||
})
|
||||
|
||||
it('flushes a final un-terminated DONE at stream end', async () => {
|
||||
const events = await collect(parseSse(bytes('data: {"a":1}\n\ndata: [DONE]')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
it('stops yielding after DONE even when more data follows', async () => {
|
||||
const events = await collect(parseSse(bytes('data: [DONE]\n\ndata: {"late":1}\n\n')))
|
||||
expect(events).toEqual([DONE])
|
||||
})
|
||||
|
||||
it('throws STREAM_CLOSED when the stream ends without DONE', async () => {
|
||||
@@ -83,26 +49,10 @@ describe('parseSse', () => {
|
||||
await expect(collect(parseSse(bytes('data: {"a"')))).rejects.toThrow(/without \[DONE\]/)
|
||||
})
|
||||
|
||||
it('stops yielding after DONE even when more data follows', async () => {
|
||||
const events = await collect(parseSse(bytes('data: [DONE]\n\ndata: {"late":1}\n\n')))
|
||||
expect(events).toEqual([DONE])
|
||||
})
|
||||
})
|
||||
|
||||
describe('parseSse edge branches', () => {
|
||||
it('handles a lone CR-terminated data line', async () => {
|
||||
// Exercises the \r-strip branch on a line that is ONLY "data:…\r".
|
||||
const events = await collect(parseSse(bytes('data: {"a":1}\r\n\r\ndata:[DONE]\r\n\r\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('strips CR from non-data field lines too', async () => {
|
||||
const events = await collect(parseSse(bytes('event: chunk\r\ndata: {"a":1}\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['{"a":1}', DONE])
|
||||
})
|
||||
|
||||
it('treats bare "data:" lines as empty payload entries', async () => {
|
||||
const events = await collect(parseSse(bytes('data:\ndata: x\n\ndata: [DONE]\n\n')))
|
||||
expect(events).toEqual(['\nx', DONE])
|
||||
it('treats a final DONE missing its blank-line terminator as truncation', async () => {
|
||||
// Spec-strict framing: an event dispatches only on its blank-line
|
||||
// terminator, so an unterminated tail at EOF is STREAM_CLOSED — real
|
||||
// providers always terminate events, so a missing terminator is truncation.
|
||||
await expect(collect(parseSse(bytes('data: {"a":1}\n\ndata: [DONE]')))).rejects.toThrow(/without \[DONE\]/)
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user