返回源码地图

packages/core/agent-loop/src/assistant-stream.ts

main snapshot · da00f7f5358f · 正文引用章节 02 / 05;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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