/** * Translate DeepSeek wire chunks into the harness `StreamChunk` protocol. * * A small state machine over the SSE payload stream: * - `delta.content` / `delta.reasoning_content` / `delta.tool_calls[i]` each * own one harness block (index allocated on first sight). The first * thinking-mode chunk carries `reasoning_content: ""` — that must NOT open * a reasoning block. * - `finish_reason` and `usage` are DEFERRED: emitted only at the `[DONE]` * sentinel, so the wire's two usage shapes (attached to the finish chunk, * or a trailing usage-only chunk) both work and nothing ever follows * `finish`. Last usage wins. * * @module dsh-llm-deepseek/translate */ import { CallId, LlmError } from '@deepseek-ai/dsh-llm' import type { ContentBlock, FinishReason, StreamChunk, TokenUsage } from '@deepseek-ai/dsh-llm' import { DONE } from './sse.ts' import type { WireChunk, WireUsage } from './types.ts' /** One open block under assembly. */ interface OpenBlock { index: number kind: 'text' | 'reasoning' | 'tool-call' text: string /** tool-call only */ callId?: string name?: string } /** Map the wire finish_reason vocabulary to the harness FinishReason. */ export function mapFinishReason(reason: string): FinishReason { switch (reason) { case 'stop': return { kind: 'stop' } case 'tool_calls': return { kind: 'tool-calls' } case 'length': return { kind: 'max-tokens' } default: // content_filter, insufficient_system_resource, future additions. return { kind: 'error', message: `model stopped: ${reason}`, code: reason.toUpperCase() } } } /** * Map wire usage fields. DeepSeek's `prompt_tokens` INCLUDES cache hits * (`prompt_tokens = prompt_cache_hit_tokens + prompt_cache_miss_tokens`, * api/create-chat-completion); the harness TokenUsage convention is * DISJOINT counts, so cache reads are subtracted out of `inputTokens`. */ export function mapUsage(usage: WireUsage): TokenUsage { const cacheRead = usage.prompt_tokens_details?.cached_tokens ?? usage.prompt_cache_hit_tokens const reasoning = usage.completion_tokens_details?.reasoning_tokens return { inputTokens: usage.prompt_tokens - (cacheRead ?? 0), outputTokens: usage.completion_tokens, ...cacheRead !== undefined ? { cacheReadTokens: cacheRead } : {}, ...reasoning !== undefined ? { reasoningTokens: reasoning } : {}, } } /** Assemble the final ContentBlock for one open block. */ function closeBlock(block: OpenBlock): ContentBlock { switch (block.kind) { case 'text': return { type: 'text', text: block.text } case 'reasoning': return { type: 'reasoning', text: block.text } case 'tool-call': return { type: 'tool-call', id: CallId(block.callId ?? ''), name: block.name ?? '', arguments: block.text, } } } /** * Consume SSE data payloads (ending with `[DONE]`) and yield StreamChunks. * Malformed JSON payloads abort the stream with `MALFORMED_RESPONSE`. */ export async function* translate(payloads: AsyncIterable): AsyncGenerator { let nextIndex = 0 let textBlock: OpenBlock | undefined let reasoningBlock: OpenBlock | undefined const toolBlocks = new Map() const order: OpenBlock[] = [] let pendingFinish: FinishReason | undefined let pendingUsage: TokenUsage | undefined function open(kind: OpenBlock['kind']): OpenBlock { const block: OpenBlock = { index: nextIndex++, kind, text: '' } order.push(block) return block } for await (const payload of payloads) { if (payload === DONE) { for (const block of order) { yield { type: 'block-end', index: block.index, block: closeBlock(block) } } if (pendingUsage) yield { type: 'usage', usage: pendingUsage } yield { type: 'finish', reason: pendingFinish ?? { kind: 'stop' } } return } let chunk: WireChunk try { chunk = JSON.parse(payload) as WireChunk } catch { throw new LlmError(`malformed SSE payload: ${payload.slice(0, 120)}`, 'MALFORMED_RESPONSE') } for (const choice of chunk.choices ?? []) { const delta = choice.delta // Reasoning first: thinking mode interleaves it before text. The // empty-string first chunk must not open a block. const reasoning = delta?.reasoning_content if (typeof reasoning === 'string' && reasoning.length > 0) { if (!reasoningBlock) { reasoningBlock = open('reasoning') yield { type: 'block-start', index: reasoningBlock.index, blockType: 'reasoning' } } reasoningBlock.text += reasoning yield { type: 'reasoning-delta', index: reasoningBlock.index, text: reasoning } } const content = delta?.content if (typeof content === 'string' && content.length > 0) { if (!textBlock) { textBlock = open('text') yield { type: 'block-start', index: textBlock.index, blockType: 'text' } } textBlock.text += content yield { type: 'text-delta', index: textBlock.index, text: content } } for (const call of delta?.tool_calls ?? []) { let block = toolBlocks.get(call.index) if (!block) { block = open('tool-call') toolBlocks.set(call.index, block) yield { type: 'block-start', index: block.index, blockType: 'tool-call' } } if (call.id !== undefined) block.callId = call.id if (call.function?.name !== undefined) block.name = call.function.name const fragment = call.function?.arguments ?? '' block.text += fragment yield { type: 'tool-call-delta', index: block.index, id: CallId(block.callId ?? ''), ...block.name !== undefined ? { name: block.name } : {}, argumentsDelta: fragment, } } if (typeof choice.finish_reason === 'string') { pendingFinish = mapFinishReason(choice.finish_reason) } } // Usage may arrive attached to the finish chunk or as a trailing // usage-only chunk — keep the latest. if (chunk.usage) pendingUsage = mapUsage(chunk.usage) } // parseSse guarantees the [DONE] sentinel (or throws); reaching here means // the payload source violated that contract. throw new LlmError('SSE payload stream ended without [DONE]', 'STREAM_CLOSED') }