region.ts 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226
  1. /**
  2. * Surface retention selection and the log-recorded compaction transaction.
  3. *
  4. * @module @deepseek-ai/dsh-compact-basic/region
  5. */
  6. import { isDeepStrictEqual } from 'node:util'
  7. import {
  8. COMPACT_CHECKPOINT_SOURCE,
  9. toolPairingBalancedAfter,
  10. toolPairingBalancedBefore,
  11. } from '@deepseek-ai/dsh-compact'
  12. import type { CompactionResult } from '@deepseek-ai/dsh-compact'
  13. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  14. import type { Message } from '@deepseek-ai/dsh-llm'
  15. import type { TokenMeasurement, TokenMeterService } from '@deepseek-ai/dsh-token-meter'
  16. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  17. import type { Agent } from '@deepseek-ai/dsh-agent'
  18. import { frameSummary } from './summarizer.ts'
  19. import type { SummarizationInput, SummaryResult } from './summarizer.ts'
  20. interface RegionDependencies {
  21. readonly meter: TokenMeterService
  22. summarize(input: SummarizationInput, agent: Agent, signal?: AbortSignal): Promise<SummaryResult>
  23. }
  24. /**
  25. * Resolve the next head-anchored range while retaining a priced recent tail
  26. * and never splitting an assistant tool-call/result pair.
  27. * @param session - session supplying authoritative current surface positions.
  28. * @param measurement - unified pressure and surface measurement from the conversation meter.
  29. * @param retainTokens - minimum recent tail budget retained verbatim.
  30. * @returns the inclusive positional seq range to compact, or `null`.
  31. */
  32. export function selectCompactableRange(
  33. session: Session,
  34. measurement: TokenMeasurement,
  35. retainTokens: number,
  36. ): { start: number; end: number } | null {
  37. const pricedNodes = measurement.nodes
  38. if (pricedNodes.length === 0) return null
  39. const surfaceNodes = session.surface.nodes
  40. if (surfaceNodes.length !== pricedNodes.length
  41. || surfaceNodes.some((seq, index) => seq !== pricedNodes[index]?.seq)) {
  42. throw new Error('compaction: token-meter surface does not match the current session surface')
  43. }
  44. let accumulated = 0
  45. let keepFromIdx = pricedNodes.length
  46. for (let index = pricedNodes.length - 1; index >= 0; index -= 1) {
  47. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  48. accumulated += pricedNodes[index]!.tokens
  49. keepFromIdx = index
  50. if (accumulated >= retainTokens) break
  51. }
  52. if (keepFromIdx === 0) return null
  53. while (keepFromIdx > 0) {
  54. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  55. if (toolPairingBalancedBefore(session, surfaceNodes[keepFromIdx]!)) break
  56. keepFromIdx -= 1
  57. }
  58. if (keepFromIdx === 0) return null
  59. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  60. const first = surfaceNodes[0]!
  61. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  62. const cutoff = surfaceNodes[keepFromIdx - 1]!
  63. return { start: first, end: cutoff }
  64. }
  65. /**
  66. * Validate and compact one positional surface span.
  67. * @param dependencies - conversation meter and dynamically dispatched summarizer hook.
  68. * @param session - session whose surface is mutated.
  69. * @param start - inclusive first surface-node seq.
  70. * @param end - inclusive last surface-node seq.
  71. * @param agent - agent used by the summarizer.
  72. * @param signal - optional summarization cancellation signal.
  73. * @returns the successful durable compaction result.
  74. */
  75. export async function compactSurfaceRegion(
  76. dependencies: RegionDependencies,
  77. session: Session,
  78. start: number,
  79. end: number,
  80. agent: Agent,
  81. signal?: AbortSignal,
  82. ): Promise<CompactionResult> {
  83. const nodes = session.surface.nodes
  84. const startIdx = nodes.indexOf(start)
  85. const endIdx = nodes.indexOf(end)
  86. if (startIdx === -1) throw new Error(`compactRegion: start seq ${start} not found in surface`)
  87. if (endIdx === -1) throw new Error(`compactRegion: end seq ${end} not found in surface`)
  88. if (startIdx > endIdx) {
  89. throw new Error(
  90. `compactRegion: start seq ${start} (position ${startIdx}) is after end seq ${end} (position ${endIdx}) on the surface`,
  91. )
  92. }
  93. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  94. if (!toolPairingBalancedBefore(session, nodes[startIdx]!)) {
  95. throw new Error(`compactRegion: start seq ${start} is not a balanced boundary (would split a step's tool-call/result pair)`)
  96. }
  97. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  98. if (!toolPairingBalancedAfter(session, nodes[endIdx]!)) {
  99. throw new Error(`compactRegion: end seq ${end} is not a balanced boundary (would split a step, or the step is still open)`)
  100. }
  101. const tail = inspectTurnTail(session.events)
  102. if (tail.compactionInProgress) throw new Error('compaction already in progress')
  103. if (tail.turn === null) {
  104. throw new Error('compactRegion: no open turn — compaction events must be enclosed in a turn')
  105. }
  106. const shadowedSeqs = nodes.slice(startIdx, endIdx + 1)
  107. const startEvent = session.append('compact/start', { turn: tail.turn })
  108. try {
  109. // Capture after the lock event so a later surface mutation invalidates the
  110. // async selection before replacement. Unrelated log-only facts may append.
  111. const lockedMeasurement = dependencies.meter.measure(session)
  112. const selected = lockedMeasurement.nodes.slice(startIdx, endIdx + 1)
  113. if (selected.length !== shadowedSeqs.length
  114. || selected.some((node, index) => node.seq !== shadowedSeqs[index])) {
  115. throw new Error('compaction: selected surface changed before summarization began')
  116. }
  117. const shadowedTokenCount = selected.reduce((total, node) => total + node.tokens, 0)
  118. const summarizationInput = buildSummarizationInput(session, shadowedSeqs)
  119. const { summary, provider, model, maxTokens } = await dependencies.summarize(summarizationInput, agent, signal)
  120. const currentMeasurement = dependencies.meter.measure(session)
  121. if (!isDeepStrictEqual(currentMeasurement.nodes, lockedMeasurement.nodes)) {
  122. throw new Error('compaction: session surface changed during summarization')
  123. }
  124. const framedSummary = frameSummary(summary)
  125. const checkpointMessage = createUserMessage({
  126. content: framedSummary,
  127. source: COMPACT_CHECKPOINT_SOURCE,
  128. })
  129. const framedSummaryTokenCount = dependencies.meter.estimateMessage(checkpointMessage)
  130. if (framedSummaryTokenCount >= shadowedTokenCount) {
  131. throw new Error(
  132. `summary is not smaller than the shadowed content (${framedSummaryTokenCount} estimated framed tokens >= ${shadowedTokenCount})`,
  133. )
  134. }
  135. const summaryEvent = session.append('compact/summary', {
  136. summary,
  137. shadowedRange: { start, end },
  138. shadowedSeqs,
  139. shadowedTokenCount,
  140. provider,
  141. model,
  142. ...maxTokens === undefined ? {} : { maxTokens },
  143. })
  144. session.append('user/message', checkpointMessage, {
  145. surfaceOp: { op: 'replace', start, end },
  146. sourceEventSeqs: [startEvent.seq, summaryEvent.seq, ...shadowedSeqs],
  147. })
  148. const endEvent = session.append('compact/end', { turn: tail.turn })
  149. return {
  150. startSeq: startEvent.seq,
  151. summarySeq: summaryEvent.seq,
  152. endSeq: endEvent.seq,
  153. summary,
  154. shadowedRange: { start, end },
  155. shadowedSeqs,
  156. shadowedTokenCount,
  157. }
  158. } catch (error: unknown) {
  159. const message = error instanceof Error ? error.message : String(error)
  160. session.append('compact/end', { turn: tail.turn, error: message })
  161. throw error
  162. }
  163. }
  164. /**
  165. * Reconstruct the last routed request's cacheable prefix for the shadowed
  166. * region: its system prompt and tool schemas, then the region's own derived
  167. * messages in surface order. The summarizer appends only the compaction
  168. * instruction after this, so the call is a genuine prefix of the conversation
  169. * and reuses the provider's KV cache.
  170. * @param session - session supplying the request header and per-node projection.
  171. * @param shadowedSeqs - the surface-node seqs, in order, being compacted.
  172. * @returns the replayed conversation prefix to condense.
  173. */
  174. function buildSummarizationInput(
  175. session: Session,
  176. shadowedSeqs: readonly number[],
  177. ): SummarizationInput {
  178. const header = session.requestHeader()
  179. const events = session.events
  180. const regionMessages = shadowedSeqs
  181. // shadowedSeqs are current surface seqs, so each is a valid log index.
  182. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  183. .map(seq => session.deriveEventMessage(events[seq]!))
  184. .filter((message): message is Message => message !== null)
  185. return {
  186. ...header?.system === undefined ? {} : { system: header.system },
  187. ...header?.tools === undefined ? {} : { tools: header.tools },
  188. messages: regionMessages,
  189. }
  190. }
  191. /** Inspect the current turn boundary and latest compaction bracket once. */
  192. function inspectTurnTail(
  193. events: readonly SessionEvent[],
  194. ): { turn: number | null; compactionInProgress: boolean } {
  195. let compactionInProgress = false
  196. let compactionStateKnown = false
  197. for (let index = events.length - 1; index >= 0; index -= 1) {
  198. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  199. const event = events[index]!
  200. if (!compactionStateKnown) {
  201. if (event.type === 'compact/start') {
  202. compactionInProgress = true
  203. compactionStateKnown = true
  204. } else if (event.type === 'compact/end') {
  205. compactionStateKnown = true
  206. }
  207. }
  208. if (event.type === 'turn/start') return { turn: event.data.turn, compactionInProgress }
  209. if (event.type === 'turn/end') return { turn: null, compactionInProgress }
  210. }
  211. return { turn: null, compactionInProgress }
  212. }