返回 DeepSeek-Reasonix
compact_projection.go
根目录 / internal / agent / compact_projection.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "slices"
8 "strings"
9
10 "reasonix/internal/event"
11 "reasonix/internal/provider"
12 "reasonix/internal/tool"
13 )
14
15 const (
16 maxCompressAnchorBytes = 512
17 maxCompressFocusBytes = 2000
18 )
19
20 var errCompressStaleContext = errors.New("compress: conversation changed while compression was running; retry with the current context")
21
22 // CompressContext implements the context-bound compress tool. It resolves the
23 // anchor against the current model-visible view and installs a projection only;
24 // the canonical transcript and checkpoint lineage remain untouched.
25 func (a *Agent) CompressContext(ctx context.Context, req tool.CompressRequest) (result tool.CompressResult, resultErr error) {
26 ctx, finish := a.beginCompactionRun(ctx)
27 defer func() { resultErr = finish(resultErr) }()
28 direction := strings.TrimSpace(req.Direction)
29 anchor := strings.TrimSpace(req.Anchor)
30 focus := strings.TrimSpace(req.Focus)
31 if direction != "before" && direction != "after" {
32 return tool.CompressResult{}, fmt.Errorf("compress: direction must be before or after")
33 }
34 if anchor == "" {
35 return tool.CompressResult{}, fmt.Errorf("compress: anchor must not be empty")
36 }
37 if len(anchor) > maxCompressAnchorBytes {
38 return tool.CompressResult{}, fmt.Errorf("compress: anchor exceeds %d bytes", maxCompressAnchorBytes)
39 }
40 if len(focus) > maxCompressFocusBytes {
41 return tool.CompressResult{}, fmt.Errorf("compress: focus exceeds %d bytes", maxCompressFocusBytes)
42 }
43
44 snap := a.snapshotExplicitCompression()
45 matches := make([]int, 0, 2)
46 for i, msg := range snap.visible {
47 if !compressAnchorCandidate(msg) {
48 continue
49 }
50 if strings.Contains(UserMessageText(msg), anchor) {
51 matches = append(matches, i)
52 }
53 }
54 if len(matches) == 0 {
55 return tool.CompressResult{}, fmt.Errorf("compress: anchor did not match any current user message; retry with an exact excerpt from a visible user turn")
56 }
57 if len(matches) > 1 {
58 return tool.CompressResult{}, fmt.Errorf("compress: anchor matched %d user messages; retry with a longer unique excerpt", len(matches))
59 }
60
61 return a.compressVisibleRange(ctx, snap, CompactionTriggerTool, direction, matches[0], anchorPreview(UserMessageText(snap.visible[matches[0]])), focus)
62 }
63
64 type explicitCompressionSnapshot struct {
65 canonical []provider.Message
66 visible []provider.Message
67 transcriptVersion uint64
68 coveredHash string
69 projectionVersion uint64
70 generation uint64
71 promptCacheKey string
72 }
73
74 func (a *Agent) snapshotExplicitCompression() explicitCompressionSnapshot {
75 canonical, version := a.sess.conversation.snapshotMessagesVersion()
76 cacheKey := a.currentPromptCacheKey()
77 a.sess.compactionMu.Lock()
78 state := a.sess.compactionState
79 a.sess.compactionMu.Unlock()
80 visible := canonical
81 if projectionValid(state, canonical, cacheKey) {
82 if projected := modelVisibleFromProjection(state.Projection, canonical); len(projected) > 0 {
83 visible = projected
84 }
85 }
86 return explicitCompressionSnapshot{
87 canonical: canonical,
88 visible: compressionVisibleMessages(visible),
89 transcriptVersion: version,
90 coveredHash: coveredPrefixHash(canonical, len(canonical)),
91 projectionVersion: state.Projection.ProjectionVersion,
92 generation: state.Generation,
93 promptCacheKey: cacheKey,
94 }
95 }
96
97 func compressionVisibleMessages(msgs []provider.Message) []provider.Message {
98 out := make([]provider.Message, 0, len(msgs)+1)
99 for _, msg := range msgs {
100 if !msg.LocalOnly {
101 summary, user, split := splitLegacyCoalescedSummary(msg)
102 if split {
103 out = append(out, summary, user)
104 } else {
105 out = append(out, msg)
106 }
107 }
108 }
109 return out
110 }
111
112 // Older schema-v1 sidecars may have persisted a strict-role merge of the
113 // summary and its following user turn. Split that legacy shape for range
114 // planning; new sidecars keep the logical messages separate and coalesce only
115 // on the provider request copy.
116 func splitLegacyCoalescedSummary(msg provider.Message) (provider.Message, provider.Message, bool) {
117 if !isCompactionSummary(msg) {
118 return provider.Message{}, provider.Message{}, false
119 }
120 separator := summaryTagClose + "\n\n"
121 i := strings.Index(msg.Content, separator)
122 if i < 0 || i+len(separator) >= len(msg.Content) {
123 return provider.Message{}, provider.Message{}, false
124 }
125 summary := msg
126 summary.Origin = provider.MessageOriginHost
127 summary.Content = msg.Content[:i+len(summaryTagClose)]
128 summary.RawContent = ""
129 summary.Images = nil
130 summary.ImageInputs = nil
131 summary.ToolCalls = nil
132 summary.ResponsesItems = nil
133 summary.ServerSearch = nil
134 summary.CreatedAt = 0
135 user := msg
136 // The legacy coalesced record did not retain the following turn's
137 // provenance. Empty keeps old-session fallback available instead of
138 // asserting that an old host continuation was user-authored.
139 user.Origin = ""
140 user.Content = msg.Content[i+len(separator):]
141 user.RawContent = ""
142 return summary, user, true
143 }
144
145 func compressAnchorCandidate(msg provider.Message) bool {
146 if msg.Role != provider.RoleUser || msg.LocalOnly || isCompactionSummary(msg) {
147 return false
148 }
149 return IsUserAuthoredTurnMessage(msg)
150 }
151
152 func anchorPreview(text string) string {
153 return truncatePreview(previewProse(text))
154 }
155
156 type visibleCompressionPlan struct {
157 result tool.CompressResult
158 foldMask []bool
159 dropMask []bool
160 fold []provider.Message
161 firstFold int
162 }
163
164 const (
165 emptyCompressionRange = "selected range is empty"
166 noVisibleCompressionMessages = "selected range has no model-visible messages"
167 noSummarizableCompressionMessages = "selected range contains no summarizable messages"
168 )
169
170 func noCompressionHistory(reason string) bool {
171 return reason == emptyCompressionRange || reason == noVisibleCompressionMessages || reason == noSummarizableCompressionMessages
172 }
173
174 type preparedVisibleCompression struct {
175 fold []provider.Message
176 instructions string
177 inputMode string
178 }
179
180 func (a *Agent) compressVisibleRange(
181 ctx context.Context,
182 snap explicitCompressionSnapshot,
183 trigger string,
184 direction string,
185 anchorIndex int,
186 preview string,
187 instructions string,
188 ) (resultOut tool.CompressResult, resultErr error) {
189 ctx, finish := a.beginCompactionRun(ctx)
190 defer func() { resultErr = finish(resultErr) }()
191 if err := a.sess.compactionRunMu.acquire(ctx); err != nil {
192 return tool.CompressResult{}, err
193 }
194 defer a.sess.compactionRunMu.Unlock()
195 if err := ctx.Err(); err != nil {
196 return tool.CompressResult{}, err
197 }
198 if !a.explicitCompressionSnapshotCurrent(snap) {
199 return tool.CompressResult{}, summaryError(errCompressStaleContext)
200 }
201 plan, ok := a.planVisibleCompression(snap, direction, anchorIndex, preview)
202 if !ok {
203 return plan.result, nil
204 }
205 result := plan.result
206 inputMode := SummaryInputNonPrefix
207 if direction == "before" && foldMatchesVisiblePrefix(snap.visible, plan.fold) {
208 inputMode = SummaryInputCachePrefix
209 }
210
211 a.svc.sink.Emit(event.Event{Kind: event.CompactionStarted, Compaction: event.Compaction{Trigger: trigger}})
212 prepared, reason, err := a.prepareVisibleCompression(ctx, trigger, plan.fold, instructions, inputMode)
213 if err != nil {
214 a.emitCompactionAborted(trigger)
215 return tool.CompressResult{}, err
216 }
217 if reason != "" {
218 a.emitCompactionAborted(trigger)
219 result.Reason = reason
220 return result, nil
221 }
222
223 res, err := a.foldToSummaryMode(ctx, prepared.fold, prepared.instructions, prepared.inputMode)
224 summary := res.Text
225 tele := compactionTelemetryFromSummary(trigger, a.CacheState(), result.SourceTokens, res)
226 if err != nil {
227 tele.Error = err.Error()
228 a.emitCompactionTelemetry(tele)
229 a.emitCompactionAborted(trigger)
230 return tool.CompressResult{}, err
231 }
232 summary, err = a.interceptCompactionComplete(ctx, summary)
233 if err != nil {
234 tele.Error = err.Error()
235 a.emitCompactionTelemetry(tele)
236 a.emitCompactionAborted(trigger)
237 return tool.CompressResult{}, err
238 }
239
240 projection := buildVisibleCompressionProjection(snap.visible, plan, summary)
241 projection, pinnedCheckpoint, err := rebasePinnedContextProjection(projection, snap.canonical, len(snap.canonical))
242 if err != nil {
243 a.emitCompactionAborted(trigger)
244 return tool.CompressResult{}, err
245 }
246 projectionTokens := a.estimatedVisibleRequestTokens(projection)
247 tele.ProjectionTokens = projectionTokens
248 result.Messages = len(plan.fold)
249 result.ProjectionTokens = projectionTokens
250 result.Mode = res.Mode
251 if projectionTokens >= result.SourceTokens {
252 if pinnedCheckpoint {
253 result.Reason = "pinned-context-too-large: checkpoint prevents compaction from reducing context"
254 a.emitCompactionTelemetry(tele)
255 a.emitCompactionAborted(trigger)
256 return result, nil
257 }
258 result.Reason = "compressed context would not be smaller"
259 a.emitCompactionTelemetry(tele)
260 a.emitCompactionAborted(trigger)
261 return result, nil
262 }
263
264 inputHash := providerVisibleFingerprint(modelInputMessages(snap.visible))
265 outputHash := providerVisibleFingerprint(projection)
266 state, err := a.commitSummaryProjection(ctx, summaryProjectionCommit{
267 canonical: snap.canonical, fold: prepared.fold, projected: projection, result: res,
268 transcriptVersion: snap.transcriptVersion, projectionVersion: snap.projectionVersion, generation: snap.generation,
269 activeTurn: a.activeTurnCreatedAt.Load(), trigger: trigger, summary: summary,
270 inputHash: inputHash, outputHash: outputHash, sourceTokens: result.SourceTokens, projectionTokens: projectionTokens,
271 covered: len(snap.canonical),
272 })
273 if err != nil {
274 if errors.Is(err, errCompressStaleContext) {
275 tele.Error = err.Error()
276 a.emitCompactionTelemetry(tele)
277 }
278 a.emitCompactionAborted(trigger)
279 return tool.CompressResult{}, err
280 }
281 a.emitCompactionTelemetry(tele)
282 a.svc.sink.Emit(event.Event{Kind: event.CompactionDone, Compaction: event.Compaction{
283 Trigger: trigger, Messages: len(plan.fold), Summary: summary, Archive: state.LastReceipt.Archive,
284 }})
285 result.Status = "ok"
286 result.Reason = ""
287 return result, nil
288 }
289
290 func foldMatchesVisiblePrefix(visible, fold []provider.Message) bool {
291 head := 0
292 if len(visible) > 0 && visible[0].Role == provider.RoleSystem {
293 head = 1
294 }
295 if len(fold) == 0 || head+len(fold) > len(visible) {
296 return false
297 }
298 return providerVisibleFingerprint(modelInputMessages(fold)) ==
299 providerVisibleFingerprint(modelInputMessages(visible[head:head+len(fold)]))
300 }
301
302 func (a *Agent) explicitCompressionSnapshotCurrent(snap explicitCompressionSnapshot) bool {
303 current, version := a.sess.conversation.snapshotMessagesVersion()
304 a.sess.compactionMu.Lock()
305 projectionVersion := a.sess.compactionState.Projection.ProjectionVersion
306 generation := a.sess.compactionState.Generation
307 a.sess.compactionMu.Unlock()
308 return version == snap.transcriptVersion && len(current) == len(snap.canonical) &&
309 coveredPrefixHash(current, len(current)) == snap.coveredHash &&
310 projectionVersion == snap.projectionVersion && generation == snap.generation &&
311 a.currentPromptCacheKey() == snap.promptCacheKey
312 }
313
314 func (a *Agent) planVisibleCompression(snap explicitCompressionSnapshot, direction string, anchorIndex int, preview string) (visibleCompressionPlan, bool) {
315 sourceTokens := a.estimatedVisibleRequestTokens(snap.visible)
316 plan := visibleCompressionPlan{result: tool.CompressResult{
317 Status: "noop",
318 Direction: direction,
319 Anchor: preview,
320 SourceTokens: sourceTokens,
321 ProjectionTokens: sourceTokens,
322 }}
323 if anchorIndex < 0 || anchorIndex >= len(snap.visible) {
324 plan.result.Reason = "anchor is no longer present in the model context"
325 return plan, false
326 }
327 head := 0
328 if len(snap.visible) > 0 && snap.visible[0].Role == provider.RoleSystem {
329 head = 1
330 }
331 completedEnd := len(snap.visible)
332 if active := a.activeTurnStart(snap.visible); active >= 0 {
333 completedEnd = active
334 }
335 start, end := head, anchorIndex
336 if direction == "after" {
337 start, end = anchorIndex, completedEnd
338 }
339 if start < head {
340 start = head
341 }
342 if end > completedEnd {
343 end = completedEnd
344 }
345 if start >= end {
346 plan.result.Reason = emptyCompressionRange
347 return plan, false
348 }
349
350 plan.foldMask = make([]bool, len(snap.visible))
351 plan.dropMask = make([]bool, len(snap.visible))
352 plan.firstFold = len(snap.visible)
353 latestContext := latestSessionContextIndex(snap.visible)
354 for i, msg := range snap.visible {
355 selected := i >= start && i < end
356 mergeSummary := i < completedEnd && isCompactionSummary(msg)
357 if isSessionContextMessage(msg) {
358 // Context never enters the summarizer. Once an older snapshot falls
359 // inside the explicitly compressed range, remove it from the
360 // projection; the latest valid snapshot remains byte-identical.
361 plan.dropMask[i] = selected && i != latestContext
362 continue
363 }
364 if msg.Role == provider.RoleSystem || i < head || (!selected && !mergeSummary) {
365 continue
366 }
367 plan.foldMask[i] = true
368 plan.fold = append(plan.fold, msg)
369 if i < plan.firstFold {
370 plan.firstFold = i
371 }
372 }
373 if len(plan.fold) == 0 {
374 plan.result.Reason = noVisibleCompressionMessages
375 return plan, false
376 }
377 return plan, true
378 }
379
380 func (a *Agent) prepareVisibleCompression(ctx context.Context, trigger string, fold []provider.Message, instructions, inputMode string) (preparedVisibleCompression, string, error) {
381 if a.svc.hooks != nil {
382 if hookInstructions := a.svc.hooks.PreCompact(ctx, trigger); hookInstructions != "" {
383 if instructions != "" {
384 instructions += "\n"
385 }
386 instructions += hookInstructions
387 }
388 }
389 filteredFold, removedPinned := withoutPinnedContextRevisions(fold)
390 if len(filteredFold) == 0 {
391 return preparedVisibleCompression{}, noSummarizableCompressionMessages, nil
392 }
393 if removedPinned {
394 inputMode = SummaryInputNonPrefix
395 }
396 originalHash := providerVisibleFingerprint(modelInputMessages(filteredFold))
397 preparedFold, preparedInstructions, err := a.interceptCompactionPrepare(ctx, filteredFold, instructions)
398 if err != nil {
399 return preparedVisibleCompression{}, "", err
400 }
401 preparedFold = modelInputMessages(preparedFold)
402 if len(preparedFold) == 0 {
403 return preparedVisibleCompression{}, "compaction hook removed the selected range", nil
404 }
405 if !removedPinned && providerVisibleFingerprint(modelInputMessages(preparedFold)) != originalHash {
406 inputMode = SummaryInputExtensionRewritten
407 }
408 return preparedVisibleCompression{fold: preparedFold, instructions: preparedInstructions, inputMode: inputMode}, "", nil
409 }
410
411 func buildVisibleCompressionProjection(visible []provider.Message, plan visibleCompressionPlan, summary string) []provider.Message {
412 projection := make([]provider.Message, 0, len(visible)-len(plan.fold)+1)
413 for i, msg := range visible {
414 if i == plan.firstFold {
415 projection = append(projection, formatSummaryMessage(summary))
416 }
417 if !plan.foldMask[i] && (len(plan.dropMask) <= i || !plan.dropMask[i]) {
418 projection = append(projection, msg)
419 }
420 }
421 return projectionMessagesPreservingPinnedContext(projection)
422 }
423
424 func compactionTelemetryFromSummary(trigger, cacheState string, sourceTokens int, res foldSummary) CompactionTelemetry {
425 tele := CompactionTelemetry{
426 Trigger: trigger, CacheState: cacheState, Mode: res.Mode,
427 SourceTokens: sourceTokens,
428 ProviderRequestID: res.RequestID,
429 FoldTokens: res.FoldTokens,
430 Spans: res.Spans,
431 SummaryInputMode: res.InputMode,
432 }
433 if tele.Spans <= 0 {
434 tele.Spans = 1
435 }
436 usage := res.Usage
437 if usage == nil {
438 return tele
439 }
440 tele.InputTokens = usage.PromptTokens
441 tele.OutputTokens = usage.CompletionTokens
442 tele.CacheHitTokens = usage.CacheHitTokens
443 tele.CacheMissTokens = usage.CacheMissTokens
444 tele.CacheWriteTokens = usage.CacheWriteTokens
445 tele.RequestCount = usage.RequestCount
446 if tele.RequestCount <= 0 {
447 tele.RequestCount = 1
448 }
449 return tele
450 }
451
452 // foldSummaryWithChunkedFallback retries summary size failures through the
453 // resilient fragment/tree-reduce path used for over-length sessions.
454 func (a *Agent) foldSummaryWithChunkedFallback(ctx context.Context, trigger string, fold []provider.Message, instructions string, sourceTokens int, inputMode string) (foldSummary, CompactionTelemetry, error) {
455 res, tele, err := a.foldSummaryWithTelemetry(ctx, trigger, fold, instructions, sourceTokens, inputMode)
456 if err == nil || !chunkedFallbackApplies(err, inputMode) {
457 return res, tele, err
458 }
459 chunked, chunkedErr := a.chunkedFoldSummary(ctx, fold, instructions, nil)
460 chunked.Usage = mergeSamplingUsage(res.Usage, chunked.Usage)
461 chunked.Spans += res.Spans
462 if chunked.FoldTokens <= 0 {
463 chunked.FoldTokens = res.FoldTokens
464 }
465 if chunked.RequestID == "" {
466 chunked.RequestID = res.RequestID
467 }
468 if chunkedErr != nil {
469 tele = compactionTelemetryFromSummary(trigger, a.CacheState(), sourceTokens, chunked)
470 tele.Error = fmt.Sprintf("%v (chunked fallback: %v)", err, chunkedErr)
471 return chunked, tele, chunkedErr
472 }
473 return chunked, compactionTelemetryFromSummary(trigger, a.CacheState(), sourceTokens, chunked), nil
474 }
475
476 // chunkedFallbackApplies reports a size failure the fragment path can fix. A
477 // provider overflow qualifies only once the transcript form has failed too;
478 // before that a re-planned replay is one request instead of many.
479 func chunkedFallbackApplies(err error, inputMode string) bool {
480 if provider.AsContextLimitError(err) != nil {
481 return inputMode == SummaryInputSlim
482 }
483 return summarySizeFailure(err)
484 }
485
486 // compact writes a context projection; trigger stays "auto"/"manual" for UI cards.
487 func (a *Agent) summarizeFold(ctx context.Context, trigger string, fold []provider.Message, instructions string, sourceTokens int, inputMode string, req foldRequest) (foldSummary, CompactionTelemetry, error) {
488 if req.allowChunked {
489 return a.foldSummaryWithChunkedFallback(ctx, trigger, fold, instructions, sourceTokens, inputMode)
490 }
491 return a.foldSummaryWithTelemetry(ctx, trigger, fold, instructions, sourceTokens, inputMode)
492 }
493
494 func (a *Agent) compactToProjectionLocked(ctx context.Context, trigger, instructions string, req foldRequest) (CompactionOutcome, error) {
495 activeTurn := a.activeTurnCreatedAt.Load()
496 canonical, transcriptVersion := a.sess.conversation.snapshotMessagesVersion()
497 a.sess.compactionMu.Lock()
498 stateSnapshot := a.sess.compactionState
499 startProjectionVersion := a.sess.compactionState.Projection.ProjectionVersion
500 startGeneration := a.sess.compactionState.Generation
501 a.sess.compactionMu.Unlock()
502 msgs, onProjection := a.visibleInputForFold(stateSnapshot, canonical, transcriptVersion)
503 viewInputHash := providerVisibleFingerprint(modelInputMessages(msgs))
504 head, start, ok := a.planFoldRegion(msgs, req.force, req.mustFree)
505 if !ok {
506 return CompactionNoop, nil
507 }
508 latestContext := latestSessionContextIndex(msgs)
509 _, preliminaryFold, _ := a.partitionFoldForProjectionAt(msgs[head:start], head, latestContext)
510 if len(preliminaryFold) == 0 || (!req.force && !foldEconomics(preliminaryFold)) {
511 return CompactionNoop, nil
512 }
513 fixedPrefixTokens := a.estimatedVisibleRequestTokens(msgs[:head])
514 if a.contextWindow > 0 && fixedPrefixTokens >= a.compactTrigger() {
515 return CompactionNoop, fmt.Errorf("%w: fixed prefix (%d tokens) already exceeds trigger (%d)", errCheckpointRejected, fixedPrefixTokens, a.compactTrigger())
516 }
517
518 a.svc.sink.Emit(event.Event{Kind: event.CompactionStarted, Compaction: event.Compaction{Trigger: trigger}})
519 if a.svc.hooks != nil {
520 if hookInstr := a.svc.hooks.PreCompact(ctx, trigger); hookInstr != "" {
521 if instructions != "" {
522 instructions += "\n"
523 }
524 instructions += hookInstr
525 }
526 }
527 // Cap every automatic summary input (#9572), including pressure folds after
528 // projection invalidation. mustFree also covers the over-ceiling manual rescue
529 // merged in #9474; ordinary manual compaction keeps its requested range.
530 if req.mustFree || trigger != CompactionTriggerManual {
531 start = a.maximumSafeSummaryPrefixEnd(msgs, head, start, instructions)
532 if start <= head {
533 a.emitCompactionAborted(trigger)
534 return CompactionNoop, fmt.Errorf("%w: no balanced prefix leaves enough room for a summary response", errCheckpointRejected)
535 }
536 }
537
538 covered, bodySuffix := projectionCoverageForFold(stateSnapshot, msgs, start, onProjection)
539 regionHadPinnedRevision := containsPinnedContextRevision(msgs[head:start])
540 kept, fold, retention := a.partitionFoldForProjectionAt(msgs[head:start], head, latestContext)
541 if len(fold) == 0 {
542 a.emitCompactionAborted(trigger)
543 return CompactionNoop, nil
544 }
545 originalFoldHash := providerVisibleFingerprint(modelInputMessages(fold))
546 var err error
547 fold, instructions, err = a.interceptCompactionPrepare(ctx, fold, instructions)
548 if err != nil {
549 a.emitCompactionAborted(trigger)
550 return CompactionNoop, err
551 }
552 if len(fold) == 0 {
553 a.emitCompactionAborted(trigger)
554 return CompactionNoop, nil
555 }
556 if req.mustFree || trigger != CompactionTriggerManual {
557 if err := a.validateSafeSummaryRequest(fold, instructions, req.slim); err != nil {
558 a.emitCompactionAborted(trigger)
559 return CompactionNoop, err
560 }
561 }
562
563 sourceTokens := a.estimatedVisibleRequestTokens(msgs)
564 inputMode := summaryInputModeFor(req, regionHadPinnedRevision,
565 providerVisibleFingerprint(modelInputMessages(fold)) != originalFoldHash)
566 res, tele, err := a.summarizeFold(ctx, trigger, fold, instructions, sourceTokens, inputMode, req)
567 if err != nil {
568 a.emitCompactionTelemetry(tele)
569 a.emitCompactionAborted(trigger)
570 return CompactionNoop, err
571 }
572 summary, err := a.interceptCompactionComplete(ctx, res.Text)
573 if err != nil {
574 tele.Error = err.Error()
575 a.emitCompactionTelemetry(tele)
576 a.emitCompactionAborted(trigger)
577 return CompactionNoop, err
578 }
579
580 // The projection body freezes only prefix + digest + kept messages; the
581 // verbatim tail splices live from canonical[start:] so tail-side rewrites
582 // (rewind truncation, snips) stay visible without rebuilding the fold.
583 projMsgs := checkpointProjectionMessages(msgs, head, kept, summary)
584 if len(bodySuffix) > 0 {
585 projMsgs = append(projMsgs, projectionMessagesPreservingPinnedContext(bodySuffix)...)
586 }
587 tele.UserTurnsKept, tele.UserTurnsDropped = retention.Kept, retention.Dropped
588 projMsgs, spliced, projTokens, err := a.preparePinnedCheckpointCandidate(trigger, projMsgs, canonical, covered, sourceTokens, &tele)
589 if err != nil {
590 a.emitCompactionAborted(trigger)
591 return CompactionNoop, err
592 }
593 viewOutputHash := providerVisibleFingerprint(modelInputMessages(spliced))
594 _, err = a.commitSummaryProjection(ctx, summaryProjectionCommit{
595 canonical: canonical, fold: fold, projected: projMsgs, result: res,
596 transcriptVersion: transcriptVersion, projectionVersion: startProjectionVersion,
597 generation: startGeneration, activeTurn: activeTurn, trigger: trigger,
598 summary: summary, inputHash: viewInputHash, outputHash: viewOutputHash,
599 sourceTokens: sourceTokens, projectionTokens: projTokens, covered: covered,
600 })
601 if err != nil {
602 a.emitCompactionAborted(trigger)
603 return CompactionNoop, err
604 }
605 a.svc.sink.Emit(event.Event{Kind: event.CompactionDone, Compaction: event.Compaction{
606 Trigger: trigger, Messages: len(fold), Summary: summary,
607 }})
608 return CompactionInstalled, nil
609 }
610
611 func (a *Agent) preparePinnedCheckpointCandidate(
612 trigger string,
613 projection, canonical []provider.Message,
614 covered, sourceTokens int,
615 tele *CompactionTelemetry,
616 ) ([]provider.Message, []provider.Message, int, error) {
617 projection, pinnedCheckpoint, err := rebasePinnedContextProjection(projection, canonical, covered)
618 if err != nil {
619 return nil, nil, 0, err
620 }
621 spliced := append(append([]provider.Message(nil), projection...), canonical[covered:]...)
622 projectionTokens := a.estimatedVisibleRequestTokens(spliced)
623 tele.ProjectionTokens = projectionTokens
624 a.emitCompactionTelemetry(*tele)
625 if err := a.acceptCheckpointCandidate(trigger, sourceTokens, projectionTokens); err != nil {
626 if pinnedCheckpoint {
627 return nil, nil, 0, fmt.Errorf("pinned-context-too-large: checkpoint prevents compaction acceptance: %w", err)
628 }
629 return nil, nil, 0, err
630 }
631 return projection, spliced, projectionTokens, nil
632 }
633
634 // projectionCoverageForFold maps a working-view boundary to canonical
635 // coverage. A suffix inside an existing frozen body remains in the new body
636 // because it has no corresponding canonical tail to splice from.
637 func projectionCoverageForFold(state CompactionState, msgs []provider.Message, start int, onProjection bool) (int, []provider.Message) {
638 if !onProjection {
639 return start, nil
640 }
641 body := len(state.Projection.Messages)
642 prior := state.Projection.CoveredCount
643 if start < body {
644 return prior, msgs[start:body]
645 }
646 return prior + (start - body), nil
647 }
648
649 // visibleInputForFold prefers the prior projection + new history over full
650 // canonical. The second return reports whether the projection was used, so
651 // fold boundaries can be translated back to canonical indices.
652 func (a *Agent) visibleInputForFold(state CompactionState, canonical []provider.Message, transcriptVersion uint64) ([]provider.Message, bool) {
653 if projectionValid(state, canonical, a.currentPromptCacheKey()) {
654 if projected := modelVisibleFromProjection(state.Projection, canonical); len(projected) > 0 {
655 return projected, true
656 }
657 }
658 return canonical, false
659 }
660
661 func checkpointProjectionMessages(msgs []provider.Message, head int, kept []provider.Message, summary string) []provider.Message {
662 projMsgs := make([]provider.Message, 0, head+1+len(kept))
663 projMsgs = append(projMsgs, msgs[:head]...)
664 projMsgs = append(projMsgs, kept...)
665 projMsgs = append(projMsgs, formatSummaryMessage(summary))
666 return provider.ProjectionMessages(projMsgs)
667 }
668
669 // acceptCheckpointCandidate requires real savings and, for automatic
670 // maintenance, a result below the physical input ceiling.
671 func (a *Agent) acceptCheckpointCandidate(trigger string, sourceTokens, candidateTokens int) error {
672 if candidateTokens >= sourceTokens {
673 return fmt.Errorf("%w: candidate would not reduce tokens (%d >= %d)", errCheckpointRejected, candidateTokens, sourceTokens)
674 }
675 hard := a.hardInputCeiling()
676 if trigger != CompactionTriggerManual && hard > 0 && candidateTokens >= hard {
677 return fmt.Errorf("%w: candidate %d still at or above physical ceiling %d", errCheckpointRejected, candidateTokens, hard)
678 }
679 return nil
680 }
681
682 // planFoldRegion returns [head:start] to fold; force shrinks the recent tail.
683 // splitActive lets an overflow rescue fold the active turn's older completed
684 // rounds as well; otherwise the active turn stays verbatim.
685 func (a *Agent) planFoldRegion(msgs []provider.Message, force, splitActive bool) (head, start int, ok bool) {
686 head, start, ok = a.planCompaction(msgs, minCompactMessages, force)
687 if !ok {
688 head, start, ok = a.planCompaction(msgs, 1, force)
689 }
690 if !ok {
691 return head, start, false
692 }
693 if active := a.activeTurnStart(msgs); active >= head && active < start {
694 if splitActive {
695 start = activeTurnFoldBoundary(msgs, active, start)
696 } else {
697 start = active
698 }
699 }
700 return head, start, start > head
701 }
702
703 type userTurnRetention struct {
704 Kept int
705 Dropped int
706 }
707
708 func (a *Agent) partitionFoldForProjection(region []provider.Message) (kept, fold []provider.Message, retention userTurnRetention) {
709 return a.partitionFoldForProjectionAt(region, 0, latestSessionContextIndex(region))
710 }
711
712 func (a *Agent) partitionFoldForProjectionAt(region []provider.Message, offset, latestContext int) (kept, fold []provider.Message, retention userTurnRetention) {
713 for i, m := range region {
714 if m.LocalOnly || IsPinnedContextRevision(m) {
715 continue
716 }
717 if isSessionContextMessage(m) {
718 if offset+i == latestContext {
719 kept = append(kept, m)
720 }
721 continue
722 }
723 fold = append(fold, m)
724 if IsUserAuthoredTurnMessage(m) {
725 retention.Dropped++
726 }
727 }
728 return kept, fold, retention
729 }
730
731 func latestSessionContextIndex(messages []provider.Message) int {
732 for i := range slices.Backward(messages) {
733 if isSessionContextMessage(messages[i]) {
734 return i
735 }
736 }
737 return -1
738 }
739
740 // runCompactionSummary uses the single local summarizer path for every provider.
741 func (a *Agent) runCompactionSummary(ctx context.Context, fold []provider.Message, instructions string) (summary, mode string, usage *provider.Usage, providerReqID string, err error) {
742 summary, usage, err = a.summarizeOnce(ctx, fold, instructions)
743 if err != nil {
744 return "", CompactionModeSummarized, usage, "", err
745 }
746 return summary, CompactionModeSummarized, usage, "", nil
747 }
748
748 lines GO