返回源码地图

packages/api/session-controller/src/control.ts

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

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

1/** Live Session projection state with reconnect baselines. */
2
3import type { Context } from '@deepseek-ai/cordis'
4import { Deque } from '@deepseek-ai/dsh-deque'
5import type {
6 Session, SessionId,
7} from '@deepseek-ai/dsh-session'
8import type { JsonValue } from '@deepseek-ai/dsh-util-values'
9import type {
10 SessionControlBaseline,
11 SessionControlFrame,
12 SessionProjectionBaseline,
13 SessionProjectionValues,
14} from './types.ts'
15
16/** Owns the Host-wide Session control stream. */
17export class SessionControlController {
18 private readonly streams = new Set<ControlQueue>()
19
20 /** @param ctx - Host context carrying live Agent and projection services. */
21 constructor(private readonly ctx: Context) {
22 ctx.sessionProjections.onChanged((session, key, value, seq) => {
23 this.broadcast({
24 type: 'projection',
25 sessionId: session.id,
26 key,
27 value: value as JsonValue,
28 seq,
29 })
30 })
31 ctx.effect(() => () => {
32 for (const stream of this.streams) stream.end()
33 this.streams.clear()
34 }, 'session-controller.control')
35 }
36
37 /**
38 * Open one generation of Host-wide live control state.
39 * @param signal - Remote stream cancellation.
40 * @returns one complete baseline followed by live replacement frames.
41 */
42 async *control(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
43 signal.throwIfAborted()
44 const queue = new ControlQueue()
45 this.streams.add(queue)
46 try {
47 yield { type: 'baseline', value: this.baseline() }
48 yield* queue.iterate(signal)
49 } finally {
50 this.streams.delete(queue)
51 queue.end()
52 }
53 }
54
55 private baseline(): SessionControlBaseline {
56 const sessions = this.ctx.sessions.list()
57 return {
58 projections: this.projectionBaseline(sessions),
59 }
60 }
61
62 private projectionBaseline(
63 sessions: readonly Session[],
64 ): Readonly<Record<SessionId, SessionProjectionBaseline>> {
65 const blocks = Object.create(null) as Record<SessionId, SessionProjectionBaseline>
66 for (const session of sessions) {
67 const snapshot = this.ctx.sessionProjections.snapshot(session)
68 blocks[session.id] = {
69 asOfSeq: snapshot.asOfSeq,
70 // Every projection definition validates its value before snapshot publication.
71 values: snapshot.values as SessionProjectionValues,
72 }
73 }
74 return blocks
75 }
76
77 private broadcast(frame: SessionControlFrame): void {
78 for (const stream of this.streams) stream.push(frame)
79 }
80}
81
82class ControlQueue {
83 private readonly buffer = new Deque<SessionControlFrame>()
84 private wake: (() => void) | undefined
85 private done = false
86
87 push(frame: SessionControlFrame): void {
88 if (this.done) return
89 this.buffer.pushBack(frame)
90 const wake = this.wake
91 this.wake = undefined
92 wake?.()
93 }
94
95 end(): void {
96 if (this.done) return
97 this.done = true
98 const wake = this.wake
99 this.wake = undefined
100 wake?.()
101 }
102
103 async *iterate(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
104 const onAbort = (): void => { this.end() }
105 signal.addEventListener('abort', onAbort, { once: true })
106 try {
107 while (!this.done && !signal.aborted) {
108 const frame = this.buffer.popFront()
109 if (frame !== undefined) {
110 yield frame
111 continue
112 }
113 await new Promise<void>((resolve) => { this.wake = resolve })
114 }
115 while (this.buffer.size > 0 && !signal.aborted) yield this.buffer.popFront() as SessionControlFrame
116 } finally {
117 signal.removeEventListener('abort', onAbort)
118 this.end()
119 }
120 }
121}