assistant-stream.ts 6.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194
  1. /** Web presentation fold joining transient Assistant frames to one durable v2 settlement. */
  2. import type {
  3. SessionAssistantStreamBaseline,
  4. SessionAssistantStreamFrame,
  5. } from '../../types.ts'
  6. import { expandAssistantStream } from '@deepseek-ai/dsh-llm/assistant-stream'
  7. import type { AssistantStreamRecord } from '@deepseek-ai/dsh-llm/assistant-stream'
  8. import type {
  9. SessionEventLikeEntry,
  10. SessionLiveEventEntry,
  11. SessionTransientEventEntry,
  12. } from '../contract/events.ts'
  13. interface ActiveAttempt {
  14. readonly attemptId: string
  15. readonly startedAfterSeq: number
  16. readonly turn: number
  17. readonly step: number
  18. nextIndex: number
  19. }
  20. /** One Web publication decision from the assistant stream fold. */
  21. export type ClientAssistantStreamResult =
  22. | { readonly type: 'publish'; readonly entry: SessionLiveEventEntry }
  23. | {
  24. readonly type: 'settlement'
  25. readonly attemptId: string
  26. readonly entry: SessionLiveEventEntry
  27. }
  28. | { readonly type: 'transient'; readonly entry: SessionTransientEventEntry }
  29. | { readonly type: 'rebaseline' }
  30. | undefined
  31. /** Keeps transient Assistant presentation behind one settlement-aware interface. */
  32. export class ClientAssistantStream {
  33. private activeAttempt: ActiveAttempt | undefined
  34. private readonly pending = new Map<number, SessionLiveEventEntry>()
  35. private publishedSeqs = new Set<number>()
  36. private durableCursor = -1
  37. private transientInGap = 0
  38. /**
  39. * Replace the durable Web window and adopt an optional reconnect baseline.
  40. * @param entries - durable entries in the replacement window.
  41. * @param baseline - compact prefix for an Assistant attempt that is still live.
  42. * @returns immediately visible durable entries plus reconstructed transient chunks.
  43. */
  44. replace(
  45. entries: readonly SessionEventLikeEntry[],
  46. baseline?: SessionAssistantStreamBaseline,
  47. ): readonly SessionEventLikeEntry[] {
  48. this.pending.clear()
  49. this.transientInGap = 0
  50. this.activeAttempt = undefined
  51. const opening = baseline?.activeAttempt
  52. if (opening !== undefined) {
  53. this.activeAttempt = {
  54. attemptId: String(opening.attemptId),
  55. startedAfterSeq: opening.startedAfterSeq,
  56. turn: opening.turn,
  57. step: opening.step,
  58. nextIndex: opening.nextIndex,
  59. }
  60. }
  61. const visible: SessionEventLikeEntry[] = [...entries]
  62. this.publishedSeqs = new Set(visible.map(entry => entry.event.seq))
  63. this.durableCursor = visible.reduce((cursor, entry) => Math.max(cursor, entry.event.seq), -1)
  64. if (opening !== undefined) {
  65. for (const [index, member] of expandAssistantStream(
  66. opening.stream as unknown as readonly AssistantStreamRecord[],
  67. ).entries()) {
  68. this.transientInGap += 1
  69. visible.push({
  70. type: 'transient',
  71. event: {
  72. type: 'assistant/live-chunk',
  73. seq: this.durableCursor + 1 - 1 / (this.transientInGap + 1),
  74. time: member.time,
  75. data: {
  76. attemptId: opening.attemptId,
  77. turn: opening.turn,
  78. step: opening.step,
  79. chunk: member.chunk,
  80. },
  81. },
  82. })
  83. if (index + 1 >= opening.nextIndex) break
  84. }
  85. }
  86. return visible
  87. }
  88. /**
  89. * Stage one durable v2 settlement while its matching live attempt is open.
  90. * @param entry - newly followed durable entry.
  91. * @returns a publication decision, or `undefined` when no entry becomes visible.
  92. */
  93. acceptDurable(entry: SessionLiveEventEntry): ClientAssistantStreamResult {
  94. const event = entry.event
  95. this.durableCursor = Math.max(this.durableCursor, event.seq)
  96. this.transientInGap = 0
  97. if (this.attemptForSettlement(event) !== undefined) {
  98. this.pending.set(event.seq, entry)
  99. return undefined
  100. }
  101. return this.publish(entry)
  102. }
  103. /**
  104. * Fold one dense transient frame and release its named durable settlement.
  105. * @param frame - next Assistant stream frame received by the follow connection.
  106. * @returns a transient, publication, or rebaseline decision, or `undefined` when no entry becomes visible.
  107. */
  108. acceptFrame(frame: SessionAssistantStreamFrame): ClientAssistantStreamResult {
  109. switch (frame.type) {
  110. case 'start':
  111. this.pending.clear()
  112. this.activeAttempt = {
  113. attemptId: String(frame.attemptId),
  114. startedAfterSeq: frame.startedAfterSeq,
  115. turn: frame.turn,
  116. step: frame.step,
  117. nextIndex: 0,
  118. }
  119. return undefined
  120. case 'chunk': {
  121. const attempt = this.activeAttempt
  122. if (attempt === undefined
  123. || attempt.attemptId !== String(frame.attemptId)
  124. || frame.index !== attempt.nextIndex) return { type: 'rebaseline' }
  125. attempt.nextIndex += 1
  126. this.transientInGap += 1
  127. return {
  128. type: 'transient',
  129. entry: {
  130. type: 'transient',
  131. event: {
  132. type: 'assistant/live-chunk',
  133. seq: this.durableCursor + 1 - 1 / (this.transientInGap + 1),
  134. time: frame.time,
  135. data: {
  136. attemptId: frame.attemptId,
  137. turn: attempt.turn,
  138. step: attempt.step,
  139. chunk: frame.chunk as never,
  140. },
  141. },
  142. },
  143. }
  144. }
  145. case 'end': {
  146. const attempt = this.activeAttempt
  147. this.activeAttempt = undefined
  148. if (attempt === undefined || attempt.attemptId !== String(frame.attemptId)) {
  149. return { type: 'rebaseline' }
  150. }
  151. if (frame.index !== attempt.nextIndex) return { type: 'rebaseline' }
  152. if (frame.outcome.kind === 'abandoned') {
  153. return this.pending.size === 0 ? undefined : { type: 'rebaseline' }
  154. }
  155. if (this.publishedSeqs.has(frame.outcome.seq)) return undefined
  156. const entry = this.pending.get(frame.outcome.seq)
  157. if (entry === undefined
  158. || entry.event.type !== frame.outcome.eventType
  159. || entry.event.data.turn !== attempt.turn
  160. || entry.event.data.step !== attempt.step) {
  161. return { type: 'rebaseline' }
  162. }
  163. this.pending.delete(frame.outcome.seq)
  164. this.publishedSeqs.add(entry.event.seq)
  165. return { type: 'settlement', attemptId: attempt.attemptId, entry }
  166. }
  167. }
  168. }
  169. private attemptForSettlement(
  170. event: SessionLiveEventEntry['event'],
  171. ): ActiveAttempt | undefined {
  172. const attempt = this.activeAttempt
  173. if (attempt === undefined
  174. || (event.type !== 'assistant/message' && event.type !== 'assistant/attempt')
  175. || (event.type === 'assistant/message' && event.surfaceOp !== 'append')
  176. || event.seq <= attempt.startedAfterSeq
  177. || attempt.turn !== event.data.turn
  178. || attempt.step !== event.data.step) return undefined
  179. return attempt
  180. }
  181. private publish(entry: SessionLiveEventEntry): ClientAssistantStreamResult {
  182. this.publishedSeqs.add(entry.event.seq)
  183. return { type: 'publish', entry }
  184. }
  185. }