返回源码地图

packages/api/session-controller/src/assistant-stream.ts

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

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

1/** Process-local assistant state retained for reconnecting Web followers. */
2
3import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'
4import { AssistantStreamAccumulator } from '@deepseek-ai/dsh-llm'
5import type { SessionSeqCursor } from '@deepseek-ai/dsh-session'
6import type { JsonValue } from '@deepseek-ai/dsh-util-values'
7import type {
8 SessionAssistantStreamAttempt,
9 SessionAssistantStreamBaseline,
10} from './types.ts'
11
12interface MutableAttempt {
13 readonly attemptId: SessionAssistantStreamAttempt['attemptId']
14 readonly startedAfterSeq: SessionSeqCursor
15 readonly turn: number
16 readonly step: number
17 readonly stream: AssistantStreamAccumulator
18 nextIndex: number
19}
20
21const EMPTY_BASELINE: SessionAssistantStreamBaseline = { revision: 0 }
22
23/**
24 * Folds dense Agent frames and materializes one shared immutable reconnect
25 * baseline per accepted revision.
26 */
27export class SessionAssistantStreamAccumulator {
28 private activeAttempt: MutableAttempt | undefined
29 private revision = 0
30 private snapshotValue: SessionAssistantStreamBaseline = EMPTY_BASELINE
31 private dirty = false
32
33 /**
34 * Fold one trusted frame from the current attached Agent lifecycle.
35 * @param frame - next dense process-local Assistant frame.
36 * @param durableCursor - last committed Session seq when this frame was observed.
37 */
38 accept(frame: AssistantStreamFrame, durableCursor: SessionSeqCursor): void {
39 if (frame.type === 'start' && frame.revision === 1 && this.revision !== 0) {
40 this.activeAttempt = undefined
41 this.revision = 0
42 }
43 if (frame.revision !== this.revision + 1) {
44 this.activeAttempt = undefined
45 this.revision = frame.revision
46 this.dirty = true
47 return
48 }
49 this.revision = frame.revision
50 switch (frame.type) {
51 case 'start':
52 this.activeAttempt = {
53 attemptId: frame.attemptId,
54 startedAfterSeq: durableCursor,
55 turn: frame.turn,
56 step: frame.step,
57 stream: new AssistantStreamAccumulator(),
58 nextIndex: 0,
59 }
60 break
61 case 'chunk': {
62 const attempt = this.activeAttempt
63 if (attempt === undefined
64 || attempt.attemptId !== frame.attemptId
65 || frame.index !== attempt.nextIndex) {
66 this.activeAttempt = undefined
67 break
68 }
69 attempt.stream.push({ time: frame.time, chunk: frame.chunk })
70 attempt.nextIndex += 1
71 break
72 }
73 case 'end':
74 this.activeAttempt = undefined
75 break
76 }
77 this.dirty = true
78 }
79
80 /**
81 * Read the cached reconnect baseline, materializing it after a state change.
82 * @returns the identity-stable baseline for the latest accepted revision.
83 */
84 snapshot(): SessionAssistantStreamBaseline {
85 if (!this.dirty) return this.snapshotValue
86 this.snapshotValue = {
87 revision: this.revision,
88 ...this.activeAttempt === undefined ? {} : {
89 activeAttempt: {
90 attemptId: this.activeAttempt.attemptId,
91 startedAfterSeq: this.activeAttempt.startedAfterSeq,
92 turn: this.activeAttempt.turn,
93 step: this.activeAttempt.step,
94 nextIndex: this.activeAttempt.nextIndex,
95 stream: this.activeAttempt.stream.snapshot() as unknown as readonly JsonValue[],
96 },
97 },
98 }
99 this.dirty = false
100 return this.snapshotValue
101 }
102}