region.ts 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560
  1. /**
  2. * Surface retention selection and the shared log-recorded compaction
  3. * transaction for automatic open-turn and manual idle-session compaction.
  4. *
  5. * @module @deepseek-ai/dsh-compaction-basic/region
  6. */
  7. import { randomUUID } from 'node:crypto'
  8. import { isDeepStrictEqual } from 'node:util'
  9. import {
  10. CompactionId,
  11. ManualCompactionError,
  12. compactCheckpointSource,
  13. toolPairingBalancedAfter,
  14. toolPairingBalancedBefore,
  15. } from '@deepseek-ai/dsh-compaction'
  16. import type { CompactionResult } from '@deepseek-ai/dsh-compaction'
  17. import type { CommandId } from '@deepseek-ai/dsh-commands/brand'
  18. import { createUserMessage, errorChain } from '@deepseek-ai/dsh-llm'
  19. import type { Message, UserMessage } from '@deepseek-ai/dsh-llm'
  20. import type { TokenMeasurement, TokenMeter } from '@deepseek-ai/dsh-token-meter'
  21. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  22. import type { Agent } from '@deepseek-ai/dsh-agent'
  23. import { frameSummary } from './summarizer.ts'
  24. import type { SummarizationInput, SummaryResult } from './summarizer.ts'
  25. interface RegionDependencies {
  26. readonly meter: TokenMeter
  27. summarize(input: SummarizationInput, agent: Agent, signal?: AbortSignal): Promise<SummaryResult>
  28. }
  29. /** One validated inclusive span of current surface positions. */
  30. interface SurfaceSelection {
  31. readonly start: number
  32. readonly end: number
  33. readonly startIdx: number
  34. readonly endIdx: number
  35. readonly shadowedSeqs: readonly number[]
  36. }
  37. /** A selection with its priced snapshot and the replay input built from it. */
  38. interface PreparedCompaction extends SurfaceSelection {
  39. readonly measurement: TokenMeasurement
  40. readonly selectedNodes: TokenMeasurement['nodes']
  41. readonly shadowedTokenCount: number
  42. /** Route-priced total of the selected span; the shrink comparison's unit. */
  43. readonly shadowedRouteTokenCount: number
  44. readonly input: SummarizationInput
  45. }
  46. type SummarizedCompaction = PreparedCompaction & SummaryResult & {
  47. readonly checkpointMessage: UserMessage
  48. }
  49. interface CompactionTransactionOptions {
  50. /** `current-turn` derives a numbered owner; `null` writes a standalone bracket. */
  51. readonly owner: 'current-turn' | null
  52. /** Surface relationship that must survive asynchronous summarization. */
  53. readonly stability: 'whole-surface' | 'selected-span'
  54. /** Optional durability checkpoint after a successfully closed bracket. */
  55. readonly flush?: () => Promise<void>
  56. /** Manual command that initiated this transaction, when present. */
  57. readonly sourceCommandId?: CommandId
  58. }
  59. interface CompactionEntryState {
  60. readonly openTurn: number | null
  61. readonly unmatchedCompactionStart: SessionEvent<'compaction/start'> | undefined
  62. readonly latestEndSeedSeq: number | undefined
  63. }
  64. /**
  65. * Rejects a summary whose replacement boundaries are no longer the ones it was
  66. * built from, distinguished from summarizer and shrink failures so a manual
  67. * caller can report the two causes differently.
  68. */
  69. class SurfaceChangedError extends Error {}
  70. /** Whether the summary may still replace the span it was built from. */
  71. type StabilityCheck = (
  72. dependencies: RegionDependencies,
  73. session: Session,
  74. prepared: PreparedCompaction,
  75. ) => void
  76. /** Failure captured after `compaction/start` has committed. */
  77. interface TransactionFailure {
  78. readonly error: unknown
  79. readonly stage: 'summary' | 'commit'
  80. }
  81. /**
  82. * Resolve the next head-anchored range while retaining a priced recent tail
  83. * and never splitting an assistant tool-call/result pair.
  84. * @param session - session supplying authoritative current surface positions.
  85. * @param measurement - unified pressure and surface measurement from the conversation meter.
  86. * @param retainTokens - minimum recent tail budget retained verbatim.
  87. * @returns the inclusive positional seq range to compact, or `null`.
  88. */
  89. export function selectCompactableRange(
  90. session: Session,
  91. measurement: TokenMeasurement,
  92. retainTokens: number,
  93. ): { start: number; end: number } | null {
  94. const pricedNodes = measurement.nodes
  95. if (pricedNodes.length === 0) return null
  96. const surfaceNodes = session.surface.nodes
  97. if (surfaceNodes.length !== pricedNodes.length
  98. || surfaceNodes.some((seq, index) => seq !== pricedNodes[index]?.seq)) {
  99. throw new Error('compaction: token-meter surface does not match the current session surface')
  100. }
  101. let accumulated = 0
  102. let keepFromIdx = pricedNodes.length
  103. for (let index = pricedNodes.length - 1; index >= 0; index -= 1) {
  104. // oxlint-disable-next-line typescript/no-non-null-assertion
  105. accumulated += pricedNodes[index]!.tokens
  106. keepFromIdx = index
  107. if (accumulated >= retainTokens) break
  108. }
  109. if (keepFromIdx === 0) return null
  110. while (keepFromIdx > 0) {
  111. // oxlint-disable-next-line typescript/no-non-null-assertion
  112. if (toolPairingBalancedBefore(session, surfaceNodes[keepFromIdx]!)) break
  113. keepFromIdx -= 1
  114. }
  115. if (keepFromIdx === 0) return null
  116. // oxlint-disable-next-line typescript/no-non-null-assertion
  117. const first = surfaceNodes[0]!
  118. // oxlint-disable-next-line typescript/no-non-null-assertion
  119. const cutoff = surfaceNodes[keepFromIdx - 1]!
  120. return { start: first, end: cutoff }
  121. }
  122. /**
  123. * Run the single compaction transaction over one selected positional span.
  124. * Selection and validation are read-only. Idle/log validation and
  125. * `compaction/start` are synchronously adjacent, so the durable opening marker is
  126. * the compaction lock before summarization yields. Every later failure makes
  127. * exactly one `compaction/end` attempt; a failed close deliberately leaves the
  128. * unmatched start detectable.
  129. * @param dependencies - conversation meter and dynamically dispatched summarizer hook.
  130. * @param session - session whose surface is mutated.
  131. * @param start - inclusive first surface-node seq.
  132. * @param end - inclusive last surface-node seq.
  133. * @param agent - agent used by the summarizer.
  134. * @param options - bracket owner, stability rule, and optional durability checkpoint.
  135. * @param signal - optional summarization cancellation signal.
  136. * @returns the successful durable compaction result.
  137. */
  138. export async function compactSurfaceRegion(
  139. dependencies: RegionDependencies,
  140. session: Session,
  141. start: number,
  142. end: number,
  143. agent: Agent,
  144. options: CompactionTransactionOptions,
  145. signal?: AbortSignal,
  146. ): Promise<CompactionResult> {
  147. if (options.owner === null) signal?.throwIfAborted()
  148. const selection = validateSurfaceRegion(session, start, end)
  149. const entryState = inspectCompactionEntryState(session.events)
  150. assertCompactionInactive(
  151. entryState.unmatchedCompactionStart,
  152. entryState.latestEndSeedSeq,
  153. 'compaction',
  154. )
  155. let owner: number | null
  156. if (options.owner === null) {
  157. if (entryState.openTurn !== null) {
  158. throw new ManualCompactionError('busy', 'manual compaction: the session already has an open turn')
  159. }
  160. owner = null
  161. } else {
  162. if (entryState.openTurn === null) {
  163. throw new Error('compactRegion: no open turn — automatic compaction events must be enclosed in a turn')
  164. }
  165. owner = entryState.openTurn
  166. }
  167. const compactionId = CompactionId(randomUUID())
  168. const lifecycle = {
  169. compactionId,
  170. ...options.sourceCommandId === undefined ? {} : { sourceCommandId: options.sourceCommandId },
  171. turn: owner,
  172. }
  173. const startEvent = session.append('compaction/start', lifecycle)
  174. const assertStable: StabilityCheck = options.stability === 'whole-surface'
  175. ? assertWholeSurfaceUnchanged
  176. : assertSelectedSpanStable
  177. let failure: TransactionFailure | undefined
  178. let flushFailure: unknown
  179. let result: CompactionResult | undefined
  180. let closed = false
  181. let closing = false
  182. let stage: TransactionFailure['stage'] = 'summary'
  183. try {
  184. const prepared = prepareCompaction(dependencies, session, selection)
  185. const summarized = await summarizeCompaction(
  186. dependencies,
  187. prepared,
  188. agent,
  189. compactionId,
  190. options.sourceCommandId,
  191. signal,
  192. )
  193. if (options.owner === null) signal?.throwIfAborted()
  194. assertStable(dependencies, session, summarized)
  195. stage = 'commit'
  196. const pending = commitCompactionBody(session, startEvent, summarized)
  197. closing = true
  198. const endEvent = session.append('compaction/end', lifecycle)
  199. closed = true
  200. result = completeCompaction(pending, endEvent)
  201. } catch (error: unknown) {
  202. failure = { error, stage: closing ? 'commit' : stage }
  203. if (!closing) {
  204. closing = true
  205. try {
  206. session.append('compaction/end', { ...lifecycle, error: errorChain(error) })
  207. closed = true
  208. } catch (closeError: unknown) {
  209. failure = { error: closeError, stage: 'commit' }
  210. }
  211. }
  212. }
  213. if (closed && options.flush !== undefined) {
  214. try {
  215. await options.flush()
  216. } catch (error: unknown) {
  217. flushFailure = error
  218. }
  219. }
  220. if (options.owner === null) signal?.throwIfAborted()
  221. if (failure !== undefined) {
  222. if (options.owner === null) throwManualFailure(failure)
  223. throw failure.error
  224. }
  225. if (flushFailure !== undefined) {
  226. throw new ManualCompactionError(
  227. 'persistence',
  228. 'manual compaction durability checkpoint failed',
  229. { cause: flushFailure },
  230. )
  231. }
  232. /* v8 ignore next -- every path without a result records and throws a failure above. */
  233. if (result === undefined) throw new Error('compaction committed without a result')
  234. return result
  235. }
  236. /** Classify one closed manual attempt without weakening cancellation precedence. */
  237. function throwManualFailure(failure: TransactionFailure): never {
  238. if (failure.stage === 'commit') {
  239. throw new ManualCompactionError(
  240. 'commit',
  241. 'manual compaction did not commit cleanly',
  242. { cause: failure.error },
  243. )
  244. }
  245. if (failure.error instanceof SurfaceChangedError) {
  246. throw new ManualCompactionError(
  247. 'changed',
  248. 'the compacted history changed during manual compaction',
  249. { cause: failure.error },
  250. )
  251. }
  252. throw new ManualCompactionError(
  253. 'summary',
  254. 'manual compaction could not produce a smaller summary',
  255. { cause: failure.error },
  256. )
  257. }
  258. /**
  259. * Reject a durable unmatched compaction marker unless a later constructor-seed
  260. * boundary proves that its owner belongs to an earlier session lifecycle.
  261. * @param unmatchedCompactionStart - latest unmatched opening marker, if any.
  262. * @param latestEndSeedSeq - newest constructor-seed boundary, if any.
  263. * @param stage - operation label included in the busy diagnostic.
  264. */
  265. function assertCompactionInactive(
  266. unmatchedCompactionStart: SessionEvent<'compaction/start'> | undefined,
  267. latestEndSeedSeq: number | undefined,
  268. stage: string,
  269. ): void {
  270. if (unmatchedCompactionStart === undefined
  271. || (latestEndSeedSeq !== undefined
  272. && latestEndSeedSeq > unmatchedCompactionStart.seq)) return
  273. throw new ManualCompactionError(
  274. 'busy',
  275. `${stage}: compaction already in progress; the session compaction lock is already active`,
  276. )
  277. }
  278. /**
  279. * Recheck the durable compaction lock after an asynchronous policy decision.
  280. * @param session - session whose latest marker state is inspected.
  281. * @param stage - operation label included in the busy diagnostic.
  282. */
  283. export function assertNoActiveCompaction(session: Session, stage: string): void {
  284. const entryState = inspectCompactionEntryState(session.events)
  285. assertCompactionInactive(
  286. entryState.unmatchedCompactionStart,
  287. entryState.latestEndSeedSeq,
  288. stage,
  289. )
  290. }
  291. /** Validate one requested surface-position span before asynchronous work begins. */
  292. function validateSurfaceRegion(session: Session, start: number, end: number): SurfaceSelection {
  293. const nodes = session.surface.nodes
  294. const startIdx = nodes.indexOf(start)
  295. const endIdx = nodes.indexOf(end)
  296. if (startIdx === -1) throw new Error(`compactRegion: start seq ${start} not found in surface`)
  297. if (endIdx === -1) throw new Error(`compactRegion: end seq ${end} not found in surface`)
  298. if (startIdx > endIdx) {
  299. throw new Error(
  300. `compactRegion: start seq ${start} (position ${startIdx}) is after end seq ${end} (position ${endIdx}) on the surface`,
  301. )
  302. }
  303. // oxlint-disable-next-line typescript/no-non-null-assertion
  304. if (!toolPairingBalancedBefore(session, nodes[startIdx]!)) {
  305. throw new Error(`compactRegion: start seq ${start} is not a balanced boundary (would split a step's tool-call/result pair)`)
  306. }
  307. // oxlint-disable-next-line typescript/no-non-null-assertion
  308. if (!toolPairingBalancedAfter(session, nodes[endIdx]!)) {
  309. throw new Error(`compactRegion: end seq ${end} is not a balanced boundary (would split a step, or the step is still open)`)
  310. }
  311. return { start, end, startIdx, endIdx, shadowedSeqs: nodes.slice(startIdx, endIdx + 1) }
  312. }
  313. /** Snapshot pricing and replay input for a validated surface range. */
  314. function prepareCompaction(
  315. dependencies: RegionDependencies,
  316. session: Session,
  317. selection: SurfaceSelection,
  318. ): PreparedCompaction {
  319. const measurement = dependencies.meter.measure(session)
  320. const selectedNodes = measurement.nodes.slice(selection.startIdx, selection.endIdx + 1)
  321. if (selectedNodes.length !== selection.shadowedSeqs.length
  322. || selectedNodes.some((node, index) => node.seq !== selection.shadowedSeqs[index])) {
  323. throw new SurfaceChangedError('compaction: selected surface changed before summarization began')
  324. }
  325. return {
  326. ...selection,
  327. measurement,
  328. selectedNodes,
  329. // The shadow-price protocol prices replacements with the fixed heuristic
  330. // so the O(1) projection fold stays in agreement with its own appends;
  331. // retention, range selection, and the shrink comparison read the
  332. // route-priced `tokens` instead.
  333. shadowedTokenCount: selectedNodes.reduce((total, node) => total + node.heuristicTokens, 0),
  334. shadowedRouteTokenCount: selectedNodes.reduce((total, node) => total + node.tokens, 0),
  335. input: buildSummarizationInput(session, selection.shadowedSeqs),
  336. }
  337. }
  338. /** Run the summarizer and frame its replacement checkpoint. */
  339. async function summarizeCompaction(
  340. dependencies: RegionDependencies,
  341. prepared: PreparedCompaction,
  342. agent: Agent,
  343. compactionId: CompactionResult['compactionId'],
  344. sourceCommandId: CommandId | undefined,
  345. signal?: AbortSignal,
  346. ): Promise<SummarizedCompaction> {
  347. const summaryResult = await dependencies.summarize(prepared.input, agent, signal)
  348. const checkpointMessage = createUserMessage({
  349. content: frameSummary(summaryResult.summary),
  350. source: compactCheckpointSource(compactionId, sourceCommandId),
  351. })
  352. // The checkpoint is text-only, so its fixed-heuristic price IS its route
  353. // price; comparing it against the span's route price asks the real
  354. // question — does the replacement lower the next request's pressure.
  355. const framedSummaryTokenCount = dependencies.meter.estimateMessage(checkpointMessage)
  356. if (framedSummaryTokenCount >= prepared.shadowedRouteTokenCount) {
  357. throw new Error(
  358. `summary is not smaller than the shadowed content (${framedSummaryTokenCount} estimated framed tokens >= ${prepared.shadowedRouteTokenCount})`,
  359. )
  360. }
  361. return {
  362. ...prepared,
  363. ...summaryResult,
  364. checkpointMessage,
  365. }
  366. }
  367. /** Reject a summary prepared against any earlier surface generation. */
  368. function assertWholeSurfaceUnchanged(
  369. dependencies: RegionDependencies,
  370. session: Session,
  371. prepared: PreparedCompaction,
  372. ): void {
  373. const current = dependencies.meter.measure(session)
  374. if (!isDeepStrictEqual(current.nodes, prepared.measurement.nodes)) {
  375. throw new SurfaceChangedError('compaction: session surface changed during summarization')
  376. }
  377. }
  378. /**
  379. * Require only that the selected span remain the same present, contiguous,
  380. * equally priced, balanced replacement target. Nodes added outside it remain
  381. * visible and do not invalidate the summary.
  382. */
  383. function assertSelectedSpanStable(
  384. dependencies: RegionDependencies,
  385. session: Session,
  386. prepared: PreparedCompaction,
  387. ): void {
  388. let current: SurfaceSelection
  389. try {
  390. current = validateSurfaceRegion(session, prepared.start, prepared.end)
  391. } catch (error: unknown) {
  392. throw new SurfaceChangedError(
  393. 'compaction: the selected span is no longer a valid replacement target',
  394. { cause: error },
  395. )
  396. }
  397. if (!isDeepStrictEqual([...current.shadowedSeqs], [...prepared.shadowedSeqs])) {
  398. throw new SurfaceChangedError('compaction: the selected span changed during summarization')
  399. }
  400. const measured = dependencies.meter.measure(session).nodes.slice(current.startIdx, current.endIdx + 1)
  401. if (!isDeepStrictEqual(measured, prepared.selectedNodes)) {
  402. throw new SurfaceChangedError('compaction: the selected span was rewritten during summarization')
  403. }
  404. }
  405. /** Append one completed summary record and replacement body without yielding. */
  406. function commitCompactionBody(
  407. session: Session,
  408. startEvent: SessionEvent<'compaction/start'>,
  409. summarized: SummarizedCompaction,
  410. ): Omit<CompactionResult, 'endSeq'> {
  411. const {
  412. start,
  413. end,
  414. shadowedSeqs,
  415. shadowedTokenCount,
  416. summary,
  417. provider,
  418. model,
  419. maxTokens,
  420. usage,
  421. checkpointMessage,
  422. } = summarized
  423. const callProvenance = summarized.llmStreamCall === true
  424. ? { rawOutput: summarized.rawOutput, llmStreamCall: true as const }
  425. : summarized.rawOutput === undefined ? {} : { rawOutput: summarized.rawOutput }
  426. const summaryEvent = session.append('compaction/summary', {
  427. compactionId: startEvent.data.compactionId,
  428. ...startEvent.data.sourceCommandId === undefined
  429. ? {}
  430. : { sourceCommandId: startEvent.data.sourceCommandId },
  431. summary,
  432. ...callProvenance,
  433. shadowedRange: { start, end },
  434. shadowedSeqs: [...shadowedSeqs],
  435. shadowedTokenCount,
  436. provider,
  437. model,
  438. ...maxTokens === undefined ? {} : { maxTokens },
  439. ...usage === undefined ? {} : { usage },
  440. })
  441. session.append('user/message', checkpointMessage, {
  442. surfaceOp: { op: 'replace', start, end },
  443. sourceEventSeqs: [startEvent.seq, summaryEvent.seq, ...shadowedSeqs],
  444. })
  445. return {
  446. compactionId: startEvent.data.compactionId,
  447. ...startEvent.data.sourceCommandId === undefined
  448. ? {}
  449. : { sourceCommandId: startEvent.data.sourceCommandId },
  450. startSeq: startEvent.seq,
  451. summarySeq: summaryEvent.seq,
  452. summary,
  453. shadowedRange: { start, end },
  454. shadowedSeqs: [...shadowedSeqs],
  455. shadowedTokenCount,
  456. }
  457. }
  458. /** Attach the successfully appended close event to a pending result. */
  459. function completeCompaction(
  460. pending: Omit<CompactionResult, 'endSeq'>,
  461. endEvent: SessionEvent<'compaction/end'>,
  462. ): CompactionResult {
  463. return { ...pending, endSeq: endEvent.seq }
  464. }
  465. /**
  466. * Reconstruct the last routed request's cacheable prefix for the shadowed
  467. * region: its system prompt and tool schemas, then the region's own derived
  468. * messages in surface order. The summarizer appends only the compaction
  469. * instruction after this, so the call is a genuine prefix of the conversation
  470. * and reuses the provider's KV cache.
  471. * @param session - session supplying the request header and per-node projection.
  472. * @param shadowedSeqs - the surface-node seqs, in order, being compacted.
  473. * @returns the replayed conversation prefix to condense.
  474. */
  475. function buildSummarizationInput(
  476. session: Session,
  477. shadowedSeqs: readonly number[],
  478. ): SummarizationInput {
  479. const header = session.requestHeader()
  480. const events = session.events
  481. const regionMessages = shadowedSeqs
  482. // shadowedSeqs are current surface seqs, so each is a valid log index.
  483. // oxlint-disable-next-line typescript/no-non-null-assertion
  484. .map(seq => session.deriveEventMessage(events[seq]!))
  485. .filter((message): message is Message => message !== null)
  486. return {
  487. ...header?.system === undefined ? {} : { system: header.system },
  488. ...header?.tools === undefined ? {} : { tools: header.tools },
  489. messages: regionMessages,
  490. }
  491. }
  492. /** Inspect open-turn, unmatched-compaction, and latest seed-boundary state independently. */
  493. function inspectCompactionEntryState(events: readonly SessionEvent[]): CompactionEntryState {
  494. let openTurn: number | null = null
  495. let openTurnStateKnown = false
  496. let unmatchedCompactionStart: SessionEvent<'compaction/start'> | undefined
  497. let compactionEntryStateKnown = false
  498. let latestEndSeedSeq: number | undefined
  499. for (let index = events.length - 1; index >= 0; index -= 1) {
  500. // oxlint-disable-next-line typescript/no-non-null-assertion
  501. const event = events[index]!
  502. if (latestEndSeedSeq === undefined && event.type === 'session/end-seed') {
  503. latestEndSeedSeq = event.seq
  504. }
  505. if (!compactionEntryStateKnown) {
  506. if (event.type === 'compaction/start') {
  507. unmatchedCompactionStart = event
  508. compactionEntryStateKnown = true
  509. } else if (event.type === 'compaction/end') {
  510. compactionEntryStateKnown = true
  511. }
  512. }
  513. if (!openTurnStateKnown) {
  514. if (event.type === 'turn/start') {
  515. openTurn = event.data.turn
  516. openTurnStateKnown = true
  517. } else if (event.type === 'turn/end') {
  518. openTurnStateKnown = true
  519. }
  520. }
  521. if (openTurnStateKnown
  522. && compactionEntryStateKnown
  523. && latestEndSeedSeq !== undefined) break
  524. }
  525. return { openTurn, unmatchedCompactionStart, latestEndSeedSeq }
  526. }