返回源码地图

packages/compaction/compaction-basic/src/region.ts

main snapshot · da00f7f5358f · 正文引用章节 09 / 09;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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