assistant-stream.ts 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140
  1. /** Process-local assistant attempt framing and durable stream accumulation. */
  2. import {
  3. AssistantStreamAccumulator,
  4. BlockAssembler,
  5. LlmAttemptId,
  6. type AssistantStreamRecord,
  7. type ContentBlock,
  8. type FinishReason,
  9. type ReplayEnvelope,
  10. type StreamChunk,
  11. type TokenUsage,
  12. } from '@deepseek-ai/dsh-llm'
  13. import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'
  14. import type { SessionEventMap, SessionId, SessionSeq } from '@deepseek-ai/dsh-session'
  15. /** Folds one model attempt into one compact stream plus ordered transient frames. */
  16. export class AssistantStreamAttempt {
  17. private readonly accumulator = new AssistantStreamAccumulator()
  18. private readonly assembler = new BlockAssembler()
  19. private index = 0
  20. private terminal = false
  21. /** Attempt identity unique within this Agent lifecycle. */
  22. readonly attemptId: LlmAttemptId
  23. /** Whether this started attempt has emitted its terminal frame. */
  24. get ended(): boolean { return this.terminal }
  25. /**
  26. * @param sessionId - identity embedded only in the Agent-lifecycle-local attempt id.
  27. * @param attempt - attached-Session-local attempt counter.
  28. * @param nextRevision - allocates the next emitted frame revision.
  29. * @param turn - durable turn owning the request.
  30. * @param step - durable step owning the request.
  31. * @param emit - agent-scoped notification publisher.
  32. */
  33. constructor(
  34. sessionId: SessionId,
  35. attempt: number,
  36. private readonly nextRevision: () => number,
  37. readonly turn: number,
  38. readonly step: number,
  39. private readonly emit: (frame: AssistantStreamFrame) => void,
  40. ) {
  41. this.attemptId = LlmAttemptId(`${sessionId}:${attempt}`)
  42. }
  43. /** Publish the opening marker before the first delivered chunk. */
  44. start(): void {
  45. this.emit({
  46. type: 'start',
  47. attemptId: this.attemptId,
  48. revision: this.nextRevision(),
  49. turn: this.turn,
  50. step: this.step,
  51. })
  52. }
  53. /** Snapshot one chunk once, then feed durable compaction, assembly, and live publication. */
  54. push(chunk: StreamChunk): void {
  55. const timed = this.accumulator.push({ time: Date.now(), chunk })
  56. this.assembler.push(timed.chunk)
  57. this.emit({
  58. type: 'chunk',
  59. attemptId: this.attemptId,
  60. revision: this.nextRevision(),
  61. index: this.index++,
  62. time: timed.time,
  63. chunk: timed.chunk,
  64. })
  65. }
  66. /**
  67. * Publish terminal settlement after the matching durable event commits.
  68. * @param eventType - durable settlement type.
  69. * @param append - synchronous durable append returning its committed seq.
  70. */
  71. settle(
  72. eventType: 'assistant/message' | 'assistant/attempt',
  73. append: () => SessionSeq,
  74. ): void {
  75. let seq: SessionSeq
  76. try {
  77. seq = append()
  78. } catch (error: unknown) {
  79. this.abandon()
  80. throw error
  81. }
  82. this.terminal = true
  83. this.emit({
  84. type: 'end',
  85. attemptId: this.attemptId,
  86. revision: this.nextRevision(),
  87. index: this.index,
  88. outcome: { kind: 'committed', eventType, seq },
  89. })
  90. }
  91. /** Publish abandonment when no durable attempt event can be committed. */
  92. abandon(): void {
  93. this.terminal = true
  94. this.emit({
  95. type: 'end',
  96. attemptId: this.attemptId,
  97. revision: this.nextRevision(),
  98. index: this.index,
  99. outcome: { kind: 'abandoned' },
  100. })
  101. }
  102. /** Exact compact stream for the final durable event. */
  103. get stream(): SessionEventMap['assistant/attempt']['stream'] {
  104. return [...this.accumulator.snapshot()] as AssistantStreamRecord[]
  105. }
  106. /** Canonical completed-message blocks from the same chunks. */
  107. blocks(): ContentBlock[] {
  108. return this.assembler.blocks()
  109. }
  110. /** Safe visible prefix when cancellation interrupts the attempt. */
  111. interruptedBlocks(): ContentBlock[] {
  112. return this.assembler.interruptedBlocks()
  113. }
  114. /** Latest adapter-reported usage in the stream. */
  115. get usage(): TokenUsage | undefined {
  116. return this.assembler.usage
  117. }
  118. /** Terminal reason, defaulting to stop when the stream omitted one. */
  119. get finish(): FinishReason {
  120. return this.assembler.finish
  121. }
  122. /** Replay state carried by the terminal finish record. */
  123. get replayState(): ReplayEnvelope | undefined {
  124. return this.assembler.replayState
  125. }
  126. }