history.ts 15 KB

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