1
0

history.ts 12 KB

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