region.ts 7.8 KB

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