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