| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160 |
- /**
- * Incremental chunk-to-message assembler. This is the single canonical assembly
- * algorithm used by the agent loop to build an assistant message from a chunk
- * stream while logging the raw chunks for replay fidelity.
- *
- * @module @deepseek-ai/dsh-llm/assembler
- */
- import { CallId } from './brand.ts'
- import { assertNever } from './never.ts'
- import { createMessage } from './message.ts'
- import type { Message, MessageSource } from './message.ts'
- import type { ContentBlock, FinishReason, StreamChunk, TokenUsage } from './types.ts'
- interface PartialBlock {
- blockType: string
- text: string
- toolCallId?: CallId
- toolCallName?: string
- toolCallArguments: string
- /** Set by `block-end` — authoritative, and freezes the partial. */
- block?: ContentBlock
- }
- /**
- * Incrementally assembles raw {@link StreamChunk}s into complete
- * {@link ContentBlock}s and a final assistant {@link Message}.
- *
- * The agent loop feeds it while logging raw chunks for replay fidelity, then
- * reads `blocks()` / `message()` / `usage` / `finish` once the stream ends.
- *
- * Tolerant of delta-only protocols (no block-start/end); deltas arriving for
- * an index already closed by `block-end` are ignored (malformed stream) so a
- * misbehaving adapter cannot grow memory or corrupt a completed block.
- */
- export class BlockAssembler {
- private partials = new Map<number, PartialBlock>()
- private order: number[] = []
- private _usage: TokenUsage | undefined
- private _finish: FinishReason | undefined
- private _replayState: unknown = undefined
- /**
- * Feed one chunk into the assembly state.
- * @param chunk - the next raw chunk, in stream order.
- */
- push(chunk: StreamChunk): void {
- switch (chunk.type) {
- case 'block-start': {
- if (!this.partials.has(chunk.index)) {
- this.order.push(chunk.index)
- this.partials.set(chunk.index, {
- blockType: chunk.blockType,
- text: '',
- toolCallArguments: '',
- })
- }
- return
- }
- case 'text-delta':
- case 'reasoning-delta': {
- const partial = this.ensure(chunk.index, chunk.type === 'text-delta' ? 'text' : 'reasoning')
- if (partial.block) return // closed by block-end; ignore stragglers
- partial.text += chunk.text
- return
- }
- case 'tool-call-delta': {
- const partial = this.ensure(chunk.index, 'tool-call')
- if (partial.block) return // closed by block-end; ignore stragglers
- partial.toolCallId = chunk.id
- if (chunk.name) partial.toolCallName = chunk.name
- partial.toolCallArguments += chunk.argumentsDelta
- return
- }
- case 'block-end': {
- const partial = this.ensure(chunk.index, chunk.block.type)
- // First close wins; ignoring re-close stragglers keeps streamed output
- // and the final assembled block in agreement.
- if (partial.block) return
- partial.block = chunk.block
- return
- }
- case 'usage': {
- this._usage = chunk.usage
- return
- }
- case 'finish': {
- this._finish = chunk.reason
- this._replayState = chunk.replayState
- return
- }
- default: return assertNever(chunk, 'BlockAssembler.push')
- }
- }
- private ensure(index: number, blockType: string): PartialBlock {
- let partial = this.partials.get(index)
- if (!partial) {
- partial = { blockType, text: '', toolCallArguments: '' }
- this.partials.set(index, partial)
- this.order.push(index)
- }
- return partial
- }
- private assemble(partial: PartialBlock, index: number): ContentBlock {
- if (partial.block) return partial.block
- switch (partial.blockType) {
- case 'text': return { type: 'text', text: partial.text }
- case 'reasoning': return { type: 'reasoning', text: partial.text }
- case 'tool-call': return {
- type: 'tool-call',
- id: partial.toolCallId ?? CallId(`call-${index}`),
- name: partial.toolCallName ?? '',
- arguments: partial.toolCallArguments,
- }
- default: throw new Error(`cannot assemble incomplete block of type "${partial.blockType}"`)
- }
- }
- /** Invariant accessor: every index in `order` has a partial. */
- private mustGet(index: number): PartialBlock {
- const partial = this.partials.get(index)
- if (!partial) throw new Error(`BlockAssembler invariant violated: no partial for index ${index}`)
- return partial
- }
- /**
- * Assemble all blocks seen so far, in stream order.
- * @returns one block per seen index; an open block assembles from its
- * accumulated deltas (an unknown block type never closed by `block-end` throws).
- */
- blocks(): ContentBlock[] {
- return this.order.map(index => this.assemble(this.mustGet(index), index))
- }
- /** Usage from the `usage` chunk; undefined until one arrives. */
- get usage(): TokenUsage | undefined {
- return this._usage
- }
- /** Finish reason from the `finish` chunk; `{kind: 'stop'}` when the stream ended without one. */
- get finish(): FinishReason {
- return this._finish ?? { kind: 'stop' }
- }
- /** Adapter-private replay state from the terminal finish chunk, if any. */
- get replayState(): unknown {
- return this._replayState
- }
- /**
- * The assembled assistant message.
- * @param source - producer attribution for the assembled message.
- * @returns a frozen assistant-role message over `blocks()` (same open-block assembly rules).
- */
- message(source: MessageSource = { kind: 'plugin', plugin: 'dsh-llm/assembler' }): Message {
- return createMessage({ role: 'assistant', content: this.blocks(), source })
- }
- }
|