index.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244
  1. /**
  2. * Basic replay-aware compaction backend.
  3. *
  4. * @module @deepseek-ai/dsh-compact-basic
  5. */
  6. import { Context } from 'cordis'
  7. import z from 'schemastery'
  8. import { CompactService } from '@deepseek-ai/dsh-compact'
  9. import type { CompactionResult, CompactionTrigger } from '@deepseek-ai/dsh-compact'
  10. import type { Session } from '@deepseek-ai/dsh-session'
  11. import { CONTEXT_WINDOW_EXCEEDED_CODE, assertNever } from '@deepseek-ai/dsh-llm'
  12. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  13. import type { Agent } from '@deepseek-ai/dsh-agent'
  14. // Type-only: makes the optional sibling service available to `ctx.get()`.
  15. import type {} from '@deepseek-ai/dsh-compact-tool-result-prune'
  16. import { resolveConfig } from './config.ts'
  17. import { compactSurfaceRegion, selectCompactableRange } from './region.ts'
  18. import { summarizeWithLlm } from './summarizer.ts'
  19. import type {
  20. BasicCompactConfig,
  21. ResolvedConfig,
  22. } from './types.ts'
  23. export type {
  24. BasicCompactConfig,
  25. ResolvedConfig,
  26. } from './types.ts'
  27. /** Resolve the exact model durably routed for the latest provider request. */
  28. function routedModel(session: Session): string | undefined {
  29. const model = session.requestHeader()?.config.model
  30. return model === undefined || model.length === 0 ? undefined : model
  31. }
  32. /**
  33. * Dependency-light compaction backend using `ctx.tokenMeter` for pressure,
  34. * retention, provenance, and summary-convergence pricing.
  35. *
  36. * `summarize()` is the sole subclass customization hook; the replay and durable
  37. * mutation strategy stays fixed so every pricing decision uses the singleton
  38. * token meter.
  39. */
  40. export class BasicCompactService extends CompactService {
  41. static inject = ['llm', 'tokenMeter']
  42. static Config: z<BasicCompactConfig> = z.object({
  43. thresholdRatio: z.number().default(0.8),
  44. retainTokens: z.number().step(1),
  45. summarizationProvider: z.string().default(''),
  46. summarizationModel: z.string().default(''),
  47. maxTokens: z.number().step(1).min(1).default(8192),
  48. compactionRetries: z.number().step(1).min(0).default(1),
  49. maxOverflowRetries: z.number().step(1).min(0).default(1),
  50. auto: z.boolean().default(true),
  51. })
  52. /** Resolved and validated compaction configuration. */
  53. readonly config: ResolvedConfig
  54. constructor(ctx: Context, config: BasicCompactConfig = {}) {
  55. super(ctx)
  56. this.config = resolveConfig(config, ctx.tokenMeter)
  57. if (this.config.auto) this._registerAutomaticCompaction()
  58. }
  59. /**
  60. * Register the automatic post-step pressure and context-overflow recovery
  61. * listeners. `compactIfNeeded` stays dynamically dispatched so subclass
  62. * overrides are honored at event time.
  63. */
  64. private _registerAutomaticCompaction(): void {
  65. const { ctx } = this
  66. const logResult = (result: CompactionResult, trigger: string): void => {
  67. ctx.logger.info(
  68. `compaction (${trigger}): shadowed ${result.shadowedSeqs.length} surface nodes `
  69. + `(seqs ${result.shadowedRange.start}-${result.shadowedRange.end}, `
  70. + `~${result.shadowedTokenCount} tokens)`,
  71. )
  72. }
  73. ctx.on('agent/post-step', async (
  74. agent: Agent,
  75. _turn: number,
  76. _step: number,
  77. signal: AbortSignal,
  78. ) => {
  79. if (signal.aborted) return
  80. try {
  81. const result = await this.compactIfNeeded(agent, 'pressure', signal)
  82. if (result !== null) logResult(result, 'post-step pressure')
  83. } catch (error: unknown) {
  84. const message = error instanceof Error ? error.message : String(error)
  85. ctx.logger.warn(`post-step compaction failed: ${message}; continuing the turn`)
  86. }
  87. })
  88. ctx.on('agent/request-error', async (agent, _turn, _step, _error, failure, priorFailures, signal, next) => {
  89. const priorOverflowFailures = priorFailures.filter(item => item.code === CONTEXT_WINDOW_EXCEEDED_CODE).length
  90. if (failure.code !== CONTEXT_WINDOW_EXCEEDED_CODE
  91. || priorOverflowFailures >= this.config.maxOverflowRetries
  92. || signal.aborted) return next()
  93. const generation = agent.session.surface.replaceGeneration
  94. let result: CompactionResult | null
  95. try {
  96. result = await this.compactIfNeeded(agent, 'context-overflow', signal)
  97. } catch (recoveryError: unknown) {
  98. const message = recoveryError instanceof Error ? recoveryError.message : String(recoveryError)
  99. // A model-free prune can land before later summary work fails. That
  100. // durable reduction is sufficient retry proof; do not discard it just
  101. // because the optional second phase threw. Cancellation still wins.
  102. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition -- signal can abort while recovery is awaited.
  103. if (!signal.aborted && agent.session.surface.replaceGeneration > generation) {
  104. ctx.logger.warn(
  105. `context-overflow compaction failed after durable surface progress: ${message}; `
  106. + 'retrying from the replacement surface',
  107. )
  108. return { action: 'retry' }
  109. }
  110. ctx.logger.warn(
  111. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition -- signal can abort while recovery is awaited.
  112. `context-overflow compaction failed: ${message}; ${signal.aborted
  113. ? 'cancellation prevents retry'
  114. : 'preserving the original request error'}`,
  115. )
  116. return next()
  117. }
  118. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition -- signal can abort while compaction is awaited.
  119. if (signal.aborted
  120. || agent.session.surface.replaceGeneration <= generation) return next()
  121. if (result !== null) logResult(result, 'context overflow recovery')
  122. return { action: 'retry' }
  123. })
  124. }
  125. /**
  126. * Summarize a rendered region through a direct one-shot `ctx.llm.stream()`
  127. * call. Override this sole hook for a template or remote summarizer.
  128. * @param text - plain-text conversation region to condense.
  129. * @param agent - supplies routed-model history, fallback model, and session id.
  130. * @param signal - optional cancellation forwarded to the adapter.
  131. * @returns safe text summary blocks and exact auxiliary-call provenance.
  132. */
  133. protected async summarize(
  134. text: string,
  135. agent: Agent,
  136. signal?: AbortSignal,
  137. ): Promise<{ summary: ContentBlock[]; provider: string; model: string; maxTokens?: number }> {
  138. return summarizeWithLlm(this.ctx, this.config, text, agent, signal)
  139. }
  140. /**
  141. * Compact for replayed post-step pressure or one provider-confirmed context
  142. * overflow. Both triggers price the latest durable routed request envelope;
  143. * overflow bypasses the normal threshold and retained-tail policy so it can
  144. * force one useful balanced reduction.
  145. * @param agent - agent whose latest durable routed request is measured.
  146. * @param trigger - normal post-step pressure or context-overflow recovery.
  147. * @param signal - live turn cancellation signal forwarded to summarization.
  148. * @returns the latest summary compaction result, or `null` when no summary ran.
  149. */
  150. override async compactIfNeeded(
  151. agent: Agent,
  152. trigger: CompactionTrigger,
  153. signal: AbortSignal,
  154. ): Promise<CompactionResult | null> {
  155. const model = routedModel(agent.session)
  156. if (model === undefined) return null
  157. const meter = this.ctx.tokenMeter
  158. const threshold = Math.floor(meter.contextWindow * this.config.thresholdRatio)
  159. let measurement = meter.measure(agent.session)
  160. switch (trigger) {
  161. case 'context-overflow':
  162. break
  163. case 'pressure':
  164. if (measurement.totalTokens < threshold) return null
  165. break
  166. /* v8 ignore next -- closed-union exhaustiveness guard */
  167. default:
  168. assertNever(trigger, 'compaction trigger')
  169. }
  170. // Pruning is optional so compact-basic remains independently composable.
  171. // Once either trigger qualifies, land the model-free pass before choosing
  172. // a summary range, then remeasure through the singleton replay fold.
  173. const prune = this.ctx.get('toolResultPrune')
  174. if (prune !== undefined) {
  175. prune.pruneSession(agent.session)
  176. measurement = meter.measure(agent.session)
  177. }
  178. if (trigger === 'context-overflow') {
  179. const range = selectCompactableRange(agent.session, measurement, 0)
  180. if (range === null) return null
  181. return this.compactRegion(range.start, range.end, agent, signal)
  182. }
  183. if (measurement.totalTokens < threshold) return null
  184. let result: CompactionResult | null = null
  185. for (let attempt = 0; attempt <= this.config.compactionRetries; attempt += 1) {
  186. const range = selectCompactableRange(agent.session, measurement, this.config.retainTokens)
  187. if (range === null) {
  188. /* v8 ignore else -- concrete replacement preserves a compactable checkpoint; subclass hooks cannot mutate it. */
  189. if (result === null) return null
  190. /* v8 ignore next -- paired with the defensive post-success branch above. */
  191. break
  192. }
  193. result = await this.compactRegion(range.start, range.end, agent, signal)
  194. measurement = meter.measure(agent.session)
  195. if (measurement.totalTokens < threshold) return result
  196. }
  197. throw new Error(
  198. `compaction still above threshold after ${this.config.compactionRetries + 1} compaction attempts `
  199. + `(${measurement.totalTokens} estimated tokens >= threshold ${threshold})`,
  200. )
  201. }
  202. /**
  203. * Compact one inclusive positional range from the agent-owned surface using
  204. * the effective token meter for all retention and shrink pricing.
  205. * @param start - inclusive first surface-node seq.
  206. * @param end - inclusive last surface-node seq.
  207. * @param agent - owner of the target session, used by the summarizer.
  208. * @param signal - optional summarization cancellation signal.
  209. * @returns the successful durable compaction result.
  210. */
  211. override async compactRegion(
  212. start: number,
  213. end: number,
  214. agent: Agent,
  215. signal?: AbortSignal,
  216. ): Promise<CompactionResult> {
  217. const session = agent.session
  218. return compactSurfaceRegion({
  219. meter: this.ctx.tokenMeter,
  220. summarize: (text, owner, abort) => this.summarize(text, owner, abort),
  221. }, session, start, end, agent, signal)
  222. }
  223. }
  224. export default BasicCompactService