history.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350
  1. /** Cold Session history pagination and live-event source. */
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import { isAppendSurfaceEvent } from '@deepseek-ai/dsh-session'
  4. import { isChunkRow, packChunkRuns, type ChunkRow } from '@deepseek-ai/dsh-session/chunk-rows'
  5. import type { SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  6. import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
  7. import type {} from '@deepseek-ai/dsh-subagent'
  8. import { TypertRemoteFailure } from '@deepseek-ai/dsh-typert-protocol'
  9. import type {
  10. SessionAddress,
  11. SessionChunkRun,
  12. SessionEventEntry,
  13. SessionFollowRequest,
  14. SessionFollowFrame,
  15. SessionHistoryRecord,
  16. SessionPage,
  17. SessionPageRequest,
  18. SessionProjectionBaseline,
  19. SessionProjectionValues,
  20. SessionWireEvent,
  21. } from './types.ts'
  22. const DEFAULT_MAX_MESSAGES = 50
  23. const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
  24. /** Implements cold-safe history operations delegated by the Session Controller. */
  25. export class SessionHistoryController {
  26. private readonly closeFollowers = new Set<() => void>()
  27. /**
  28. * @param ctx - Host context carrying Session query and projection services.
  29. * @param promote - starts ordinary Session activation after snapshot delivery.
  30. */
  31. constructor(
  32. private readonly ctx: Context,
  33. private readonly promote: (observation: SessionObservation) => void,
  34. ) {
  35. ctx.effect(() => () => {
  36. for (const close of this.closeFollowers) close()
  37. this.closeFollowers.clear()
  38. }, 'session-controller.history')
  39. }
  40. /**
  41. * Read one message-aligned history page without activating an Agent.
  42. * @param request - durable address and backwards-page cursor.
  43. * @param signal - caller cancellation for persistence reads.
  44. * @returns a contiguous event page.
  45. */
  46. async page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {
  47. validatePageRequest(request)
  48. using source = await this.sourceFor(request.address, signal, false)
  49. signal.throwIfAborted()
  50. const sourceLog = source.events
  51. const sourceCursor = sourceLog.at(-1)?.seq ?? -1
  52. if (request.throughSeq > sourceCursor) {
  53. reject(
  54. 'bad-request',
  55. `session page through seq ${String(request.throughSeq)} is past cursor ${String(sourceCursor)}`,
  56. {},
  57. )
  58. }
  59. /* v8 ignore next -- Session and persistence validation guarantee a dense zero-based event prefix. */
  60. if (request.throughSeq >= 0 && sourceLog[request.throughSeq]?.seq !== request.throughSeq) {
  61. reject('internal', `session log does not contain through seq ${String(request.throughSeq)}`, {})
  62. }
  63. const page = paginate(
  64. sourceLog,
  65. request.beforeSeq,
  66. request.maxMessages ?? DEFAULT_MAX_MESSAGES,
  67. request.throughSeq,
  68. )
  69. const records = pageRecords(page.events)
  70. return {
  71. records,
  72. hasMore: page.hasMore,
  73. }
  74. }
  75. /**
  76. * Follow events appended after an initial cursor on one durable address.
  77. * @param request - durable address and last committed sequence already held by the caller.
  78. * @param signal - stream cancellation owned by the Remote carrier.
  79. * @returns a complete opening snapshot followed by gap-free event frames.
  80. */
  81. async *follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {
  82. validateFollowRequest(request)
  83. const { address } = request
  84. const target = addressId(address)
  85. const buffered: SessionEvent[] = []
  86. let snapshotCursor: number | undefined
  87. let wake: (() => void) | undefined
  88. const notify = (): void => {
  89. const resume = wake
  90. wake = undefined
  91. resume?.()
  92. }
  93. const follower = { closed: false }
  94. const close = (): void => {
  95. follower.closed = true
  96. notify()
  97. }
  98. this.closeFollowers.add(close)
  99. const disposeEvent = this.ctx.on('session/event', (session, event) => {
  100. if (session.id !== target) return
  101. buffered.push(event)
  102. notify()
  103. }, { global: true })
  104. const disposeCreated = this.ctx.on('session/created', (session) => {
  105. if (session.id !== target) return
  106. // Constructor seed events have no session/event notification. Normally
  107. // only the end-seed suffix is new; if persistence advanced after the
  108. // opening observation, replay everything beyond that snapshot cursor.
  109. const suffix = session.events.slice(snapshotCursor === undefined
  110. ? session.firstLiveSeq
  111. : snapshotCursor + 1)
  112. buffered.unshift(...suffix)
  113. notify()
  114. }, { global: true })
  115. const onAbort = (): void => { notify() }
  116. signal.addEventListener('abort', onAbort, { once: true })
  117. try {
  118. using source = await this.sourceFor(address, signal, true)
  119. const events = source.events
  120. signal.throwIfAborted()
  121. const cursor = source.cursor
  122. snapshotCursor = cursor
  123. const page = paginate(events, undefined, request.maxMessages ?? DEFAULT_MAX_MESSAGES)
  124. yield {
  125. type: 'snapshot',
  126. header: source.header,
  127. cursor,
  128. records: pageRecords(page.events),
  129. hasMore: page.hasMore,
  130. projections: source.projections === undefined
  131. ? { asOfSeq: cursor, values: {} }
  132. : projectionBlock(source.projections),
  133. }
  134. if (address.kind === 'session' && source.source === 'prepared') {
  135. const promotion = source.retain()
  136. try {
  137. this.promote(promotion)
  138. } catch (error: unknown) {
  139. promotion[Symbol.dispose]()
  140. throw error
  141. }
  142. }
  143. let nextSeq = cursor + 1
  144. while (!follower.closed && !signal.aborted) {
  145. const item = buffered.shift()
  146. if (item === undefined) {
  147. await new Promise<void>((resolve) => { wake = resolve })
  148. continue
  149. }
  150. if (item.seq < nextSeq) continue
  151. if (item.seq !== nextSeq) {
  152. reject('internal', `session event stream skipped seq ${String(nextSeq)}`, {})
  153. }
  154. nextSeq++
  155. yield entryFor(item)
  156. }
  157. } finally {
  158. this.closeFollowers.delete(close)
  159. signal.removeEventListener('abort', onAbort)
  160. disposeCreated()
  161. disposeEvent()
  162. }
  163. }
  164. private async sourceFor(
  165. address: SessionAddress,
  166. signal: AbortSignal,
  167. withProjections: boolean,
  168. ): Promise<SessionObservation> {
  169. const sessionId = addressId(address)
  170. try {
  171. const observation = await this.ctx.sessionQuery.observeSession(sessionId, {
  172. signal,
  173. projectionMode: withProjections || address.kind === 'subagent' ? 'all' : 'none',
  174. })
  175. if (observation.header.cwd === undefined) {
  176. observation[Symbol.dispose]()
  177. rejectNotFound(address)
  178. }
  179. try {
  180. validateAddress(address, observation.header, observation.projections)
  181. } catch (error: unknown) {
  182. observation[Symbol.dispose]()
  183. throw error
  184. }
  185. return observation
  186. } catch (error: unknown) {
  187. if (error instanceof SessionQueryError
  188. && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') rejectNotFound(address)
  189. throw error
  190. }
  191. }
  192. }
  193. function projectionBlock(
  194. snapshot: NonNullable<SessionObservation['projections']>,
  195. ): SessionProjectionBaseline {
  196. return {
  197. asOfSeq: snapshot.asOfSeq,
  198. // Projection definitions validate whole JSON values before snapshot publication.
  199. values: snapshot.values as SessionProjectionValues,
  200. }
  201. }
  202. function validatePageRequest(request: SessionPageRequest): void {
  203. if (!Number.isSafeInteger(request.throughSeq) || request.throughSeq < -1) {
  204. reject('bad-request', 'throughSeq must be an integer greater than or equal to -1', {})
  205. }
  206. if (request.beforeSeq !== undefined
  207. && (!Number.isSafeInteger(request.beforeSeq) || request.beforeSeq < 0)) {
  208. reject('bad-request', 'beforeSeq must be a non-negative safe integer', {})
  209. }
  210. if (request.maxMessages !== undefined
  211. && (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {
  212. reject('bad-request', 'maxMessages must be a positive safe integer', {})
  213. }
  214. }
  215. function validateFollowRequest(request: SessionFollowRequest): void {
  216. if (request.maxMessages !== undefined
  217. && (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {
  218. reject('bad-request', 'maxMessages must be a positive safe integer', {})
  219. }
  220. }
  221. function addressId(address: SessionAddress): SessionId {
  222. return address.kind === 'session' ? address.sessionId : address.childSessionId
  223. }
  224. function validateAddress(
  225. address: SessionAddress,
  226. header: SessionHeader,
  227. projections: SessionObservation['projections'],
  228. ): void {
  229. if (address.kind === 'session') {
  230. if (header.origin === 'subagent') {
  231. reject('agent-busy', 'subagent Sessions require their durable parent address', {
  232. reason: 'use subagent delivery for this child session',
  233. })
  234. }
  235. return
  236. }
  237. if (header.origin !== 'subagent' || header.parentSession !== address.parentSessionId) {
  238. reject('subagent-unauthorized', 'subagent does not belong to the supplied parent', {
  239. childSessionId: address.childSessionId,
  240. })
  241. }
  242. const identity = projections?.values.subagent
  243. if (identity === null) {
  244. reject('subagent-catalog-diagnostic', 'subagent descriptor is corrupt', {
  245. parentSessionId: address.parentSessionId,
  246. childSessionId: address.childSessionId,
  247. reason: 'corrupt',
  248. })
  249. }
  250. if (identity === undefined || identity.seq < (header.seedLength ?? 0)) {
  251. reject('subagent-catalog-diagnostic', 'subagent descriptor is unavailable', {
  252. parentSessionId: address.parentSessionId,
  253. childSessionId: address.childSessionId,
  254. reason: 'unsupported',
  255. })
  256. }
  257. if (identity.mode !== address.mode) {
  258. reject('subagent-unauthorized', 'subagent mode does not match the supplied address', {
  259. childSessionId: address.childSessionId,
  260. })
  261. }
  262. }
  263. function rejectNotFound(address: SessionAddress): never {
  264. if (address.kind === 'session') {
  265. reject('session-not-found', `session "${address.sessionId}" not found`, { sessionId: address.sessionId })
  266. }
  267. reject('subagent-not-found', 'subagent is unavailable', {
  268. parentSessionId: address.parentSessionId,
  269. childSessionId: address.childSessionId,
  270. })
  271. }
  272. function reject(code: string, message: string, details: object): never {
  273. throw new TypertRemoteFailure({ code, message, details })
  274. }
  275. function paginate(
  276. events: readonly SessionEvent[],
  277. beforeSeq: number | undefined,
  278. maxMessages: number,
  279. throughSeq = events.at(-1)?.seq ?? -1,
  280. ): { readonly events: SessionEvent[]; readonly hasMore: boolean } {
  281. const end = Math.min(throughSeq + 1, beforeSeq ?? throughSeq + 1)
  282. let count = 0
  283. let cut = 0
  284. for (let index = end - 1; index >= 0; index--) {
  285. const event = events[index] as SessionEvent
  286. if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
  287. count++
  288. const sources = (event as { readonly sourceEventSeqs?: readonly number[] }).sourceEventSeqs
  289. let groupStart = event.seq
  290. if (sources !== undefined) {
  291. for (const source of sources) groupStart = Math.min(groupStart, source)
  292. }
  293. if (count >= maxMessages) {
  294. cut = groupStart
  295. break
  296. }
  297. }
  298. return { events: events.slice(cut, end), hasMore: cut > 0 }
  299. }
  300. function entryFor(event: SessionEvent): SessionEventEntry {
  301. return {
  302. type: 'event',
  303. // Session.append validates and freezes event data as JSON before publication.
  304. event: event as unknown as SessionWireEvent,
  305. }
  306. }
  307. function chunkEntryFor(row: ChunkRow): SessionChunkRun {
  308. switch (row.type) {
  309. case 'text-chunks':
  310. return {
  311. type: 'chunks',
  312. event: { type: 'chunkrow/text-chunks', seq: row.seq0, time: row.time0, data: row.data },
  313. }
  314. case 'reasoning-chunks':
  315. return {
  316. type: 'chunks',
  317. event: { type: 'chunkrow/reasoning-chunks', seq: row.seq0, time: row.time0, data: row.data },
  318. }
  319. case 'tool-call-chunks':
  320. return {
  321. type: 'chunks',
  322. event: { type: 'chunkrow/tool-call-chunks', seq: row.seq0, time: row.time0, data: row.data },
  323. }
  324. }
  325. }
  326. /** Encode one bounded logical page without changing its pagination cut. */
  327. function pageRecords(events: readonly SessionEvent[]): SessionHistoryRecord[] {
  328. return packChunkRuns(events).map(record => isChunkRow(record)
  329. ? chunkEntryFor(record)
  330. : entryFor(record))
  331. }