index.ts 8.5 KB

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