assembler.ts 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160
  1. /**
  2. * Incremental chunk-to-message assembler. This is the single canonical assembly
  3. * algorithm used by the agent loop to build an assistant message from a chunk
  4. * stream while logging the raw chunks for replay fidelity.
  5. *
  6. * @module @deepseek-ai/dsh-llm/assembler
  7. */
  8. import { CallId } from './brand.ts'
  9. import { assertNever } from './never.ts'
  10. import { createMessage } from './message.ts'
  11. import type { Message, MessageSource } from './message.ts'
  12. import type { ContentBlock, FinishReason, StreamChunk, TokenUsage } from './types.ts'
  13. interface PartialBlock {
  14. blockType: string
  15. text: string
  16. toolCallId?: CallId
  17. toolCallName?: string
  18. toolCallArguments: string
  19. /** Set by `block-end` — authoritative, and freezes the partial. */
  20. block?: ContentBlock
  21. }
  22. /**
  23. * Incrementally assembles raw {@link StreamChunk}s into complete
  24. * {@link ContentBlock}s and a final assistant {@link Message}.
  25. *
  26. * The agent loop feeds it while logging raw chunks for replay fidelity, then
  27. * reads `blocks()` / `message()` / `usage` / `finish` once the stream ends.
  28. *
  29. * Tolerant of delta-only protocols (no block-start/end); deltas arriving for
  30. * an index already closed by `block-end` are ignored (malformed stream) so a
  31. * misbehaving adapter cannot grow memory or corrupt a completed block.
  32. */
  33. export class BlockAssembler {
  34. private partials = new Map<number, PartialBlock>()
  35. private order: number[] = []
  36. private _usage: TokenUsage | undefined
  37. private _finish: FinishReason | undefined
  38. private _replayState: unknown = undefined
  39. /**
  40. * Feed one chunk into the assembly state.
  41. * @param chunk - the next raw chunk, in stream order.
  42. */
  43. push(chunk: StreamChunk): void {
  44. switch (chunk.type) {
  45. case 'block-start': {
  46. if (!this.partials.has(chunk.index)) {
  47. this.order.push(chunk.index)
  48. this.partials.set(chunk.index, {
  49. blockType: chunk.blockType,
  50. text: '',
  51. toolCallArguments: '',
  52. })
  53. }
  54. return
  55. }
  56. case 'text-delta':
  57. case 'reasoning-delta': {
  58. const partial = this.ensure(chunk.index, chunk.type === 'text-delta' ? 'text' : 'reasoning')
  59. if (partial.block) return // closed by block-end; ignore stragglers
  60. partial.text += chunk.text
  61. return
  62. }
  63. case 'tool-call-delta': {
  64. const partial = this.ensure(chunk.index, 'tool-call')
  65. if (partial.block) return // closed by block-end; ignore stragglers
  66. partial.toolCallId = chunk.id
  67. if (chunk.name) partial.toolCallName = chunk.name
  68. partial.toolCallArguments += chunk.argumentsDelta
  69. return
  70. }
  71. case 'block-end': {
  72. const partial = this.ensure(chunk.index, chunk.block.type)
  73. // First close wins; ignoring re-close stragglers keeps streamed output
  74. // and the final assembled block in agreement.
  75. if (partial.block) return
  76. partial.block = chunk.block
  77. return
  78. }
  79. case 'usage': {
  80. this._usage = chunk.usage
  81. return
  82. }
  83. case 'finish': {
  84. this._finish = chunk.reason
  85. this._replayState = chunk.replayState
  86. return
  87. }
  88. default: return assertNever(chunk, 'BlockAssembler.push')
  89. }
  90. }
  91. private ensure(index: number, blockType: string): PartialBlock {
  92. let partial = this.partials.get(index)
  93. if (!partial) {
  94. partial = { blockType, text: '', toolCallArguments: '' }
  95. this.partials.set(index, partial)
  96. this.order.push(index)
  97. }
  98. return partial
  99. }
  100. private assemble(partial: PartialBlock, index: number): ContentBlock {
  101. if (partial.block) return partial.block
  102. switch (partial.blockType) {
  103. case 'text': return { type: 'text', text: partial.text }
  104. case 'reasoning': return { type: 'reasoning', text: partial.text }
  105. case 'tool-call': return {
  106. type: 'tool-call',
  107. id: partial.toolCallId ?? CallId(`call-${index}`),
  108. name: partial.toolCallName ?? '',
  109. arguments: partial.toolCallArguments,
  110. }
  111. default: throw new Error(`cannot assemble incomplete block of type "${partial.blockType}"`)
  112. }
  113. }
  114. /** Invariant accessor: every index in `order` has a partial. */
  115. private mustGet(index: number): PartialBlock {
  116. const partial = this.partials.get(index)
  117. if (!partial) throw new Error(`BlockAssembler invariant violated: no partial for index ${index}`)
  118. return partial
  119. }
  120. /**
  121. * Assemble all blocks seen so far, in stream order.
  122. * @returns one block per seen index; an open block assembles from its
  123. * accumulated deltas (an unknown block type never closed by `block-end` throws).
  124. */
  125. blocks(): ContentBlock[] {
  126. return this.order.map(index => this.assemble(this.mustGet(index), index))
  127. }
  128. /** Usage from the `usage` chunk; undefined until one arrives. */
  129. get usage(): TokenUsage | undefined {
  130. return this._usage
  131. }
  132. /** Finish reason from the `finish` chunk; `{kind: 'stop'}` when the stream ended without one. */
  133. get finish(): FinishReason {
  134. return this._finish ?? { kind: 'stop' }
  135. }
  136. /** Adapter-private replay state from the terminal finish chunk, if any. */
  137. get replayState(): unknown {
  138. return this._replayState
  139. }
  140. /**
  141. * The assembled assistant message.
  142. * @param source - producer attribution for the assembled message.
  143. * @returns a frozen assistant-role message over `blocks()` (same open-block assembly rules).
  144. */
  145. message(source: MessageSource = { kind: 'plugin', plugin: 'dsh-llm/assembler' }): Message {
  146. return createMessage({ role: 'assistant', content: this.blocks(), source })
  147. }
  148. }