返回源码地图

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

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

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

1/** Cold Session history pagination and live-event source. */
2
3import type { Context } from '@deepseek-ai/cordis'
4import { Deque } from '@deepseek-ai/dsh-deque'
5import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'
6import {
7 isAppendSurfaceEvent,
8 SessionLogOffset,
9 SessionSeq,
10} from '@deepseek-ai/dsh-session'
11import type {
12 SessionEvent,
13 SessionHeader,
14 SessionId,
15 SessionLogOffset as SessionLogOffsetType,
16 SessionSeqCursor,
17} from '@deepseek-ai/dsh-session'
18import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
19import type {} from '@deepseek-ai/dsh-subagent'
20import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
21import type { JsonValue } from '@deepseek-ai/dsh-util-values'
22import type {
23 SessionAddress,
24 SessionAssistantStreamFrame,
25 SessionEventEntry,
26 SessionFollowRequest,
27 SessionFollowFrame,
28 SessionHistoryRecord,
29 SessionPage,
30 SessionPageRequest,
31 SessionProjectionBaseline,
32 SessionProjectionValues,
33 SessionWireHeader,
34 SessionWireEvent,
35} from './types.ts'
36import { SessionAssistantStreamAccumulator } from './assistant-stream.ts'
37
38const DEFAULT_MAX_MESSAGES = 50
39const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
40
41/** Implements cold-safe history operations delegated by the Session Controller. */
42export class SessionHistoryController {
43 private readonly closeFollowers = new Set<() => void>()
44 private readonly assistantStreams = new Map<SessionId, SessionAssistantStreamAccumulator>()
45
46 /**
47 * @param ctx - Host context carrying Session query and projection services.
48 * @param promote - starts ordinary Session activation after snapshot delivery.
49 */
50 constructor(
51 private readonly ctx: Context,
52 private readonly promote: (observation: SessionObservation) => void,
53 ) {
54 ctx.on('agent/assistant-stream', ({ agent, frame }) => {
55 let stream = this.assistantStreams.get(agent.session.id)
56 if (stream === undefined) {
57 stream = new SessionAssistantStreamAccumulator()
58 this.assistantStreams.set(agent.session.id, stream)
59 }
60 stream.accept(frame, cursorBeforeNext(agent.session.seq))
61 }, { global: true })
62 ctx.on('agent/disposed', ({ agent }) => {
63 this.assistantStreams.delete(agent.session.id)
64 }, { global: true })
65 ctx.effect(() => () => {
66 for (const close of this.closeFollowers) close()
67 this.closeFollowers.clear()
68 }, 'session-controller.history')
69 }
70
71 /**
72 * Read one message-aligned history page without activating an Agent.
73 * @param request - durable address and backwards-page cursor.
74 * @param signal - caller cancellation for persistence reads.
75 * @returns a contiguous event page.
76 */
77 async page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {
78 validatePageRequest(request)
79 const throughSeq: SessionSeqCursor = request.throughSeq === -1
80 ? -1
81 : SessionSeq(request.throughSeq)
82 const beforeSeq = request.beforeSeq === undefined
83 ? undefined
84 : SessionLogOffset(request.beforeSeq)
85 using source = await this.sourceFor(request.address, signal, false)
86 signal.throwIfAborted()
87 const sourceLog = source.events
88 const sourceCursor: SessionSeqCursor = sourceLog.at(-1)?.seq ?? -1
89 if (throughSeq > sourceCursor) {
90 throw new RemoteError(
91 'gateway/bad-request',
92 `session page through seq ${String(throughSeq)} is past cursor ${String(sourceCursor)}`,
93 {},
94 )
95 }
96 /* v8 ignore next -- Session and persistence validation guarantee a dense zero-based event prefix. */
97 if (throughSeq >= 0 && sourceLog[throughSeq]?.seq !== throughSeq) {
98 throw new RemoteError('gateway/internal', `session log does not contain through seq ${String(throughSeq)}`, {})
99 }
100 const page = paginate(
101 sourceLog,
102 beforeSeq,
103 request.maxMessages ?? DEFAULT_MAX_MESSAGES,
104 throughSeq,
105 request.turnWindow,
106 )
107 const records = pageRecords(page.events)
108 return {
109 records,
110 hasMore: page.hasMore,
111 }
112 }
113
114 /**
115 * Follow events appended after an initial cursor on one durable address.
116 * @param request - durable address and last committed sequence already held by the caller.
117 * @param signal - stream cancellation owned by the Remote carrier.
118 * @returns a complete opening snapshot followed by gap-free durable events and opted-in assistant frames.
119 */
120 async *follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {
121 validateHistoryWindow(request)
122 const { address } = request
123 const target = addressId(address)
124 const buffered = new Deque<
125 | { readonly type: 'event'; readonly event: SessionEvent }
126 | {
127 readonly type: 'assistant-stream'
128 readonly frame: SessionAssistantStreamFrame
129 readonly ordinal: number
130 }
131 >()
132 let snapshotCursor: SessionSeqCursor | undefined
133 let assistantStreamOrdinal = 0
134 let wake: (() => void) | undefined
135 const notify = (): void => {
136 const resume = wake
137 wake = undefined
138 resume?.()
139 }
140 const follower = { closed: false }
141 const close = (): void => {
142 follower.closed = true
143 notify()
144 }
145 this.closeFollowers.add(close)
146 const disposeEvent = this.ctx.on('session/event', (session, event) => {
147 if (session.id !== target) return
148 buffered.pushBack({ type: 'event', event })
149 notify()
150 }, { global: true })
151 const disposeCreated = this.ctx.on('session/created', (session) => {
152 if (session.id !== target) return
153 // Constructor seed events have no session/event notification. Normally
154 // only the end-seed suffix is new; if persistence advanced after the
155 // opening observation, replay everything beyond that snapshot cursor.
156 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
157 const suffix = session.snapshotEvents(snapshotCursor === undefined
158 ? session.firstLiveSeq
159 : SessionLogOffset(snapshotCursor + 1))
160 for (let index = suffix.length - 1; index >= 0; index -= 1) {
161 buffered.pushFront({ type: 'event', event: suffix[index] as SessionEvent })
162 }
163 notify()
164 }, { global: true })
165 const disposeAssistantStream = request.assistantStream !== true
166 ? undefined
167 : this.ctx.on('agent/assistant-stream', ({ agent, frame }) => {
168 if (agent.session.id !== target) return
169 buffered.pushBack({
170 type: 'assistant-stream',
171 frame: wireAssistantStreamFrame(frame, cursorBeforeNext(agent.session.seq)),
172 ordinal: ++assistantStreamOrdinal,
173 })
174 notify()
175 }, { global: true })
176 const onAbort = (): void => { notify() }
177 signal.addEventListener('abort', onAbort, { once: true })
178 try {
179 using source = await this.sourceFor(address, signal, true)
180 const events = source.events
181 signal.throwIfAborted()
182 const cursor = source.cursor
183 snapshotCursor = cursor
184 const page = paginate(events, undefined, request.maxMessages ?? DEFAULT_MAX_MESSAGES, cursor, request.turnWindow)
185 const assistantStream = request.assistantStream === true
186 ? this.assistantStreams.get(target)?.snapshot() ?? { revision: 0 }
187 : undefined
188 // The accumulator snapshot and this watermark are synchronous. Frames
189 // through the cut are represented or superseded by that baseline,
190 // including larger revisions from a retired Agent; later revision
191 // resets reach Client continuity validation.
192 const assistantStreamOrdinalCut = assistantStreamOrdinal
193 yield {
194 type: 'snapshot',
195 header: wireHeader(source.header),
196 cursor,
197 records: pageRecords(page.events),
198 hasMore: page.hasMore,
199 projections: source.projections === undefined
200 ? { asOfSeq: cursor, values: {} }
201 : projectionBlock(source.projections),
202 ...assistantStream === undefined ? {} : { assistantStream },
203 }
204 if (address.kind === 'session' && source.source === 'prepared') {
205 const promotion = source.retain()
206 try {
207 this.promote(promotion)
208 } catch (error: unknown) {
209 promotion[Symbol.dispose]()
210 throw error
211 }
212 }
213 let nextOffset = SessionLogOffset(cursor + 1)
214 while (!follower.closed && !signal.aborted) {
215 const item = buffered.popFront()
216 if (item === undefined) {
217 await new Promise<void>((resolve) => { wake = resolve })
218 continue
219 }
220 if (item.type === 'assistant-stream') {
221 if (item.ordinal > assistantStreamOrdinalCut) {
222 yield { type: 'assistant-stream', frame: item.frame }
223 }
224 continue
225 }
226 const expectedSeq = SessionSeq(nextOffset)
227 if (item.event.seq < expectedSeq) continue
228 if (item.event.seq !== expectedSeq) {
229 throw new RemoteError('gateway/internal', `session event stream skipped seq ${String(expectedSeq)}`, {})
230 }
231 nextOffset = SessionLogOffset(nextOffset + 1)
232 yield entryFor(item.event)
233 }
234 } finally {
235 this.closeFollowers.delete(close)
236 signal.removeEventListener('abort', onAbort)
237 disposeCreated()
238 disposeEvent()
239 disposeAssistantStream?.()
240 }
241 }
242
243 private async sourceFor(
244 address: SessionAddress,
245 signal: AbortSignal,
246 withProjections: boolean,
247 ): Promise<SessionObservation> {
248 const sessionId = addressId(address)
249 try {
250 const observation = await this.ctx.sessionQuery.observeSession(sessionId, {
251 signal,
252 projectionMode: withProjections || address.kind === 'subagent' ? 'all' : 'none',
253 })
254 if (observation.header.cwd === undefined) {
255 observation[Symbol.dispose]()
256 rejectNotFound(address)
257 }
258 try {
259 validateAddress(
260 address,
261 observation.header,
262 observation.inheritedEventCount,
263 observation.projections,
264 )
265 } catch (error: unknown) {
266 observation[Symbol.dispose]()
267 throw error
268 }
269 return observation
270 } catch (error: unknown) {
271 if (error instanceof SessionQueryError
272 && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') rejectNotFound(address)
273 throw error
274 }
275 }
276
277}
278
279function cursorBeforeNext(nextSeq: SessionLogOffsetType): SessionSeqCursor {
280 return nextSeq === 0 ? -1 : SessionSeq(nextSeq - 1)
281}
282
283function wireAssistantStreamFrame(
284 frame: AssistantStreamFrame,
285 durableCursor: SessionSeqCursor,
286): SessionAssistantStreamFrame {
287 if (frame.type === 'start') return { ...frame, startedAfterSeq: durableCursor }
288 if (frame.type === 'end') return frame
289 return {
290 ...frame,
291 chunk: frame.chunk as JsonValue,
292 }
293}
294
295function projectionBlock(
296 snapshot: NonNullable<SessionObservation['projections']>,
297): SessionProjectionBaseline {
298 return {
299 asOfSeq: snapshot.asOfSeq,
300 // Projection definitions validate whole JSON values before snapshot publication.
301 values: snapshot.values as SessionProjectionValues,
302 }
303}
304
305function validatePageRequest(request: SessionPageRequest): void {
306 if (!Number.isSafeInteger(request.throughSeq)
307 || request.throughSeq < -1
308 || Object.is(request.throughSeq, -0)) {
309 throw new RemoteError('gateway/bad-request', 'throughSeq must be an integer greater than or equal to -1', {})
310 }
311 if (request.beforeSeq !== undefined
312 && (!Number.isSafeInteger(request.beforeSeq)
313 || request.beforeSeq < 0
314 || Object.is(request.beforeSeq, -0))) {
315 throw new RemoteError('gateway/bad-request', 'beforeSeq must be a non-negative safe integer', {})
316 }
317 validateHistoryWindow(request)
318}
319
320function validateHistoryWindow(request: Pick<SessionPageRequest, 'maxMessages' | 'turnWindow'>): void {
321 if (request.maxMessages !== undefined
322 && (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {
323 throw new RemoteError('gateway/bad-request', 'maxMessages must be a positive safe integer', {})
324 }
325 const window = request.turnWindow
326 if (window !== undefined) {
327 if (!Number.isSafeInteger(window.minMessages) || window.minMessages <= 0
328 || window.minMessages > (request.maxMessages ?? DEFAULT_MAX_MESSAGES)) {
329 throw new RemoteError('gateway/bad-request', 'turnWindow.minMessages must be a positive safe integer no greater than maxMessages', {})
330 }
331 if (!Number.isSafeInteger(window.minTurns) || window.minTurns <= 0) {
332 throw new RemoteError('gateway/bad-request', 'turnWindow.minTurns must be a positive safe integer', {})
333 }
334 }
335}
336
337function addressId(address: SessionAddress): SessionId {
338 return address.kind === 'session' ? address.sessionId : address.childSessionId
339}
340
341function validateAddress(
342 address: SessionAddress,
343 header: SessionHeader,
344 inheritedEventCount: SessionLogOffsetType,
345 projections: SessionObservation['projections'],
346): void {
347 if (address.kind === 'session') {
348 if (header.origin === 'subagent') {
349 throw new RemoteError('session/agent-busy', 'subagent Sessions require their durable parent address', {
350 reason: 'use subagent delivery for this child session',
351 })
352 }
353 return
354 }
355 if (header.origin !== 'subagent' || header.parentSession !== address.parentSessionId) {
356 throw new RemoteError('subagent/unauthorized', 'subagent does not belong to the supplied parent', {
357 childSessionId: address.childSessionId,
358 })
359 }
360 const identity = projections?.values.subagent
361 if (identity === null) {
362 throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is corrupt', {
363 parentSessionId: address.parentSessionId,
364 childSessionId: address.childSessionId,
365 reason: 'corrupt',
366 })
367 }
368 if (identity === undefined || identity.seq < inheritedEventCount) {
369 throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is unavailable', {
370 parentSessionId: address.parentSessionId,
371 childSessionId: address.childSessionId,
372 reason: 'unsupported',
373 })
374 }
375 if (address.mode !== 'unknown' && identity.mode !== address.mode) {
376 throw new RemoteError('subagent/unauthorized', 'subagent mode does not match the supplied address', {
377 childSessionId: address.childSessionId,
378 })
379 }
380}
381
382function rejectNotFound(address: SessionAddress): never {
383 if (address.kind === 'session') {
384 throw new RemoteError('session/not-found', `session "${address.sessionId}" not found`, { sessionId: address.sessionId })
385 }
386 throw new RemoteError('subagent/not-found', 'subagent is unavailable', {
387 parentSessionId: address.parentSessionId,
388 childSessionId: address.childSessionId,
389 })
390}
391
392function paginate(
393 events: readonly SessionEvent[],
394 beforeSeq: SessionLogOffsetType | undefined,
395 maxMessages: number,
396 throughSeq: SessionSeqCursor,
397 turnWindow?: SessionPageRequest['turnWindow'],
398): { readonly events: SessionEvent[]; readonly hasMore: boolean } {
399 const end = SessionLogOffset(Math.min(throughSeq + 1, beforeSeq ?? throughSeq + 1))
400 let count = 0
401 let turns = 0
402 let cut = SessionLogOffset(0)
403 for (let index = end - 1; index >= 0; index--) {
404 const event = events[index] as SessionEvent
405 if (turnWindow !== undefined && event.type === 'turn/start') {
406 turns++
407 if (count >= turnWindow.minMessages && turns >= turnWindow.minTurns) {
408 cut = SessionLogOffset(index)
409 break
410 }
411 }
412 if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
413 count++
414 const sources = event.sourceEventSeqs
415 let groupStart = event.seq
416 if (sources !== undefined) {
417 for (const source of sources) {
418 if (source < groupStart) groupStart = source
419 }
420 }
421 if (count >= maxMessages) {
422 cut = SessionLogOffset(groupStart)
423 break
424 }
425 }
426 return { events: events.slice(cut, end), hasMore: cut > 0 }
427}
428
429/** Translate current logical Session metadata to the browser wire. */
430function wireHeader(header: SessionHeader): SessionWireHeader {
431 return { ...header }
432}
433
434function entryFor(event: SessionEvent): SessionEventEntry {
435 return {
436 type: 'event',
437 // Session.append validates and freezes event data as JSON before publication.
438 event: event as unknown as SessionWireEvent,
439 }
440}
441
442/** Encode one bounded logical page without changing its pagination cut. */
443function pageRecords(events: readonly SessionEvent[]): SessionHistoryRecord[] {
444 return events.map(entryFor)
445}