返回 DeepSeek-Reasonix
context_manager.go
根目录 / internal / agent / context_manager.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "sync/atomic"
9
10 "reasonix/internal/provider"
11 )
12
13 // compactionProgress is how compaction is faring in this session: whether a
14 // fold stopped reducing, how many ran back to back, and which retries already
15 // ran in the active turn. The fields are cleared together on lineage resets.
16 type compactionProgress struct {
17 stuck bool // a fold landed above the trigger, so the same-view pressure retry is pointless
18 stuckInputHash string // provider-visible view covered by stuck; changed input may retry
19 consecutive int // back-to-back folds since one last helped
20 // failedTurn backs off changed-view retries within one active tool loop.
21 // A later user turn may retry, while hard-ceiling recovery bypasses it.
22 failedTurn atomic.Int64
23 // lastTurn stops the post-turn observer and the pre-send preflight from
24 // paying for two summaries during one active tool loop.
25 lastTurn atomic.Int64
26 }
27
28 // ContextManager is the sole owner of provider-visible context maintenance.
29 // Canonical session messages are immutable inputs; Prepare evolves only the
30 // durable projection and returns the exact visible view for one sampling round.
31 type ContextManager struct {
32 agent *Agent
33 }
34
35 // ContextPreparePolicy describes one maintenance transaction.
36 type ContextPreparePolicy struct {
37 Trigger string
38 Instructions string
39 Force bool
40 // ObservedInputTokens is used by compatibility harnesses that invoke the
41 // old post-turn shim directly. Production Prepare estimates the current view
42 // from its calibrated final request shape.
43 ObservedInputTokens int
44 // AllowChunkedFallback enables fragment/tree-reduce recovery after a single
45 // summary fails. Ordinary pressure/overflow leave this false.
46 AllowChunkedFallback bool
47 }
48
49 // PreparedContext is the frozen result of a successful Prepare transaction.
50 type PreparedContext struct {
51 Messages []provider.Message
52 InputTokens int
53 ProjectionVersion uint64
54 }
55
56 func (a *Agent) contextManager() ContextManager { return ContextManager{agent: a} }
57
58 // PrepareContext is the public automatic-maintenance entry used by smoke tools
59 // and controllers that need a one-shot Prepare without sampling.
60 func (a *Agent) PrepareContext(ctx context.Context) error {
61 _, err := a.contextManager().Prepare(ctx, ContextPreparePolicy{Trigger: CompactionTriggerPressure})
62 return err
63 }
64
65 // ObserveUsage is retained as a compatibility hook. Usage observations never
66 // mutate the provider-visible checkpoint.
67 func (m ContextManager) ObserveUsage(u *provider.Usage) {
68 _ = u
69 }
70
71 // Prepare is the sole automatic maintenance entry. Below compact_ratio it does
72 // nothing. At or above the trigger it runs one single-flight prune/summary
73 // transaction, with at most two successful summary attempts under pressure.
74 func (m ContextManager) Prepare(ctx context.Context, policy ContextPreparePolicy) (result PreparedContext, err error) {
75 // Legacy desktop callers can compact before their runtime context is installed.
76 if ctx == nil {
77 ctx = context.Background()
78 }
79 if err := ctx.Err(); err != nil {
80 return PreparedContext{}, err
81 }
82 if policy.Trigger == "" {
83 policy.Trigger = CompactionTriggerPressure
84 }
85 if m.agent == nil {
86 return PreparedContext{}, nil
87 }
88 ctx, finish := m.agent.beginCompactionRun(ctx)
89 defer func() { err = finish(err) }()
90 if err := m.agent.sess.compactionRunMu.acquire(ctx); err != nil {
91 return PreparedContext{}, err
92 }
93 defer m.agent.sess.compactionRunMu.Unlock()
94 // Cancellation may have arrived while another maintenance transaction held the lock.
95 // Reject it before any fast path or projection maintenance can run.
96 if err := ctx.Err(); err != nil {
97 return PreparedContext{}, err
98 }
99 return m.prepareOnce(ctx, policy)
100 }
101
102 func (m ContextManager) prepareOnce(ctx context.Context, policy ContextPreparePolicy) (PreparedContext, error) {
103 a := m.agent
104 if a == nil || a.sess.conversation == nil {
105 return PreparedContext{}, nil
106 }
107 visible := a.modelVisibleMessages()
108 // Threshold uses the stable pre-interceptor request shape (messages + tools
109 // + role projection). Extension interceptors run only on the real sampling
110 // request so side-effecting plugins are not double-invoked; if they expand
111 // the prompt past the hard ceiling, overflow recovery still fires.
112 est := a.estimatedVisibleRequestTokens(visible)
113 viewEst := est
114 prepared := PreparedContext{
115 Messages: append([]provider.Message(nil), visible...),
116 InputTokens: est,
117 ProjectionVersion: a.currentProjectionVersion(),
118 }
119 if len(visible) == 0 {
120 return prepared, nil
121 }
122 // Disabling automatic maintenance must not bypass the shared recovery
123 // ladder for an explicit request, including when the window is unknown.
124 if a.contextWindow <= 0 && policy.Trigger != CompactionTriggerManual {
125 return prepared, nil
126 }
127 fold := a.compactTrigger()
128 hard := a.hardInputCeiling()
129 if policy.ObservedInputTokens > 0 {
130 est = policy.ObservedInputTokens
131 prepared.InputTokens = est
132 }
133 inputHash := a.contextMaintenanceInputHash(visible)
134 // Receipts back off sub-critical retries only. At the ceiling a failed
135 // summary blocks this attempt and preserves the last committed view.
136 if blocked, _ := a.contextMaintenanceBlocked(inputHash, viewEst); blocked && policy.Trigger != CompactionTriggerManual &&
137 policy.Trigger != CompactionTriggerOverflow && est < hard {
138 return prepared, nil
139 }
140 if est < fold {
141 a.resetCompactionProgress()
142 }
143 if a.sess.compaction.stuck && a.sess.compaction.stuckInputHash != inputHash {
144 // The previous projection could not reclaim enough from its exact view,
145 // but newly appended messages create a new fold boundary and may retry.
146 a.sess.compaction.stuck = false
147 a.sess.compaction.stuckInputHash = ""
148 a.sess.compaction.consecutive = 0
149 }
150 if a.sess.compaction.stuck && policy.Trigger == CompactionTriggerPressure && est < hard {
151 return prepared, nil
152 }
153 // One user trigger. Overflow is a one-shot physical recovery path only.
154 forceFold := policy.Force || policy.Trigger == CompactionTriggerManual || policy.Trigger == CompactionTriggerOverflow || est >= hard
155 if est < fold && !forceFold {
156 return prepared, nil
157 }
158
159 // A manual compact over the hard ceiling is a rescue, not a convenience:
160 // prune first so the never-folded recent tail can shrink too.
161 if shouldPruneBeforeFold(policy.Trigger, hard > 0 && est >= hard) {
162 applied, err := a.pruneToolResultsToProjectionLocked(ctx, policy.Trigger)
163 if err != nil {
164 return PreparedContext{}, err
165 }
166 if applied {
167 prepared = m.currentPrepared()
168 est = prepared.InputTokens
169 inputHash = a.contextMaintenanceInputHash(prepared.Messages)
170 if (policy.Trigger == CompactionTriggerPressure && est < fold) ||
171 (policy.Trigger == CompactionTriggerOverflow && est < hard) {
172 return prepared, nil
173 }
174 }
175 }
176
177 return m.foldContext(ctx, prepared, policy, inputHash, est, fold, hard, forceFold)
178 }
179
180 func shouldPruneBeforeFold(trigger string, overHardCeiling bool) bool {
181 switch trigger {
182 case CompactionTriggerPressure, CompactionTriggerOverflow:
183 return true
184 case CompactionTriggerManual:
185 return overHardCeiling
186 default:
187 return false
188 }
189 }
190
191 // manualRecoverySummaries bounds the rescue loop for a manual compact that
192 // starts at or above the hard input ceiling. Each batch folds the largest
193 // admissible prefix, so a handful of batches recovers even a view several
194 // times the window while capping summarizer spend on pathological input.
195 const manualRecoverySummaries = 4
196
197 func maxSummariesFor(policy ContextPreparePolicy, overCeiling bool) int {
198 switch {
199 case policy.Trigger == CompactionTriggerManual && overCeiling:
200 return manualRecoverySummaries
201 case policy.Trigger == CompactionTriggerPressure:
202 return 2
203 default:
204 return 1
205 }
206 }
207
208 func (m ContextManager) foldContext(ctx context.Context, prepared PreparedContext, policy ContextPreparePolicy, inputHash string, est, fold, hard int, forceFold bool) (PreparedContext, error) {
209 a := m.agent
210 // Reserve the manual rescue budget when an overflow may reveal a window.
211 // With no known ceiling, the first successful fold completes the request.
212 maxSummaries := maxSummariesFor(policy, hard <= 0 || est >= hard)
213 ladder := newSummaryLadder(maxSummaries)
214 result := prepared
215 for ladder.next() {
216 mustFree := policy.Trigger == CompactionTriggerOverflow || hard > 0 && result.InputTokens >= hard
217 outcome, err := a.compactToProjectionLocked(ctx, policy.Trigger, policy.Instructions,
218 ladder.request(forceFold, mustFree, policy.AllowChunkedFallback))
219 if err != nil {
220 if ctxErr := ctx.Err(); ctxErr != nil && !errors.Is(context.Cause(ctx), errSummaryBudget) {
221 return PreparedContext{}, ctxErr
222 }
223 if ladder.absorbOverflow(err) {
224 // Summary feedback may have learned both the physical window and
225 // a denser tokenizer. Re-plan against those measurements.
226 fold, hard = a.compactTrigger(), a.hardInputCeiling()
227 result = m.currentPrepared()
228 continue
229 }
230 return m.summaryFailed(ctx, policy, inputHash, hard, err)
231 }
232 if outcome == CompactionNoop {
233 return m.summaryNoop(ctx, policy, inputHash, hard)
234 }
235
236 result = m.currentPrepared()
237 if foldLanded(policy, result.InputTokens, fold, hard) {
238 a.resetCompactionProgress()
239 return result, nil
240 }
241 forceFold = false
242 inputHash = a.contextMaintenanceInputHash(result.Messages)
243 }
244
245 reason := fmt.Sprintf("summary result remains above fold trigger after %d attempts (%d >= %d)", maxSummaries, result.InputTokens, fold)
246 blockedInputHash := a.contextMaintenanceInputHash(result.Messages)
247 a.recordContextMaintenanceBlocked(blockedInputHash, policy.Trigger, "summary", reason)
248 a.sess.compaction.stuck = true
249 a.sess.compaction.stuckInputHash = blockedInputHash
250 a.sess.compaction.consecutive += maxSummaries
251 if policy.Trigger == CompactionTriggerOverflow || hard > 0 && result.InputTokens >= hard {
252 return PreparedContext{}, fmt.Errorf("%w: %w", ErrCompactionRequired, summaryError(fmt.Errorf("%w: %s", errCheckpointRejected, reason)))
253 }
254 slog.Info("agent: context maintenance paused below hard ceiling", "reason", reason)
255 return result, nil
256 }
257
258 func foldLanded(policy ContextPreparePolicy, tokens, fold, hard int) bool {
259 switch policy.Trigger {
260 case CompactionTriggerManual, CompactionTriggerOverflow:
261 return hard <= 0 || tokens < hard || tokens < fold
262 default:
263 return tokens < fold
264 }
265 }
266
267 func (m ContextManager) summaryFailed(ctx context.Context, policy ContextPreparePolicy, inputHash string, hard int, err error) (PreparedContext, error) {
268 a := m.agent
269 var persistence *compactionPersistenceError
270 if errors.As(err, &persistence) {
271 return PreparedContext{}, err
272 }
273 if ctxErr := ctx.Err(); ctxErr != nil && !errors.Is(context.Cause(ctx), errSummaryBudget) {
274 return PreparedContext{}, ctxErr
275 }
276 err = summaryError(compactionError(ctx, err))
277 if errors.Is(err, errCompressStaleContext) && policy.Trigger != CompactionTriggerManual {
278 reason := "context changed during summary; automatic retry blocked for this generation"
279 a.recordContextMaintenanceBlocked(inputHash, policy.Trigger, "summary", reason)
280 return m.rescueOrFail(ctx, policy, hard, err)
281 }
282 status := "failed"
283 if errors.Is(err, errSummaryOutputTruncated) || errors.Is(err, errCheckpointRejected) {
284 status = "blocked"
285 }
286 a.recordContextMaintenanceOutcome(inputHash, policy.Trigger, "summary", status, fmt.Sprintf("context summary failed: %v", err))
287 return m.rescueOrFail(ctx, policy, hard, err)
288 }
289
290 func (m ContextManager) summaryNoop(ctx context.Context, policy ContextPreparePolicy, inputHash string, hard int) (PreparedContext, error) {
291 if err := ctx.Err(); err != nil {
292 return PreparedContext{}, err
293 }
294 reason := "context is above the maintenance threshold but no foldable region remains"
295 latest := m.currentPrepared()
296 switch {
297 case policy.Trigger == CompactionTriggerOverflow || hard > 0 && latest.InputTokens >= hard:
298 m.agent.recordContextMaintenanceBlocked(inputHash, policy.Trigger, "summary", reason)
299 return PreparedContext{}, fmt.Errorf("%w: %w", ErrCompactionRequired, summaryError(fmt.Errorf("%w: %s", errCheckpointRejected, reason)))
300 case policy.Force:
301 // A requested compaction with no eligible history is a successful no-op.
302 // It must not poison the retry ledger or masquerade as a hard-limit failure.
303 return latest, nil
304 default:
305 return latest, nil
306 }
307 }
308
309 // rescueOrFail keeps the last committed view below the ceiling. At the hard
310 // boundary it stops the attempt without installing any lossy fallback.
311 func (m ContextManager) rescueOrFail(ctx context.Context, policy ContextPreparePolicy, hard int, cause error) (PreparedContext, error) {
312 if err := ctx.Err(); err != nil && !errors.Is(context.Cause(ctx), errSummaryBudget) {
313 return PreparedContext{}, err
314 }
315 latest := m.currentPrepared()
316 if policy.Trigger != CompactionTriggerOverflow && (hard <= 0 || latest.InputTokens < hard) {
317 if policy.Trigger == CompactionTriggerManual {
318 return PreparedContext{}, cause
319 }
320 return latest, nil
321 }
322 return PreparedContext{}, fmt.Errorf("%w: %w", ErrCompactionRequired, cause)
323 }
324
325 func (a *Agent) resetCompactionProgress() {
326 a.sess.compaction.stuck = false
327 a.sess.compaction.stuckInputHash = ""
328 a.sess.compaction.consecutive = 0
329 a.sess.compaction.failedTurn.Store(0)
330 }
331
332 func (m ContextManager) currentPrepared() PreparedContext {
333 if m.agent == nil {
334 return PreparedContext{}
335 }
336 visible := m.agent.modelVisibleMessages()
337 return PreparedContext{
338 Messages: append([]provider.Message(nil), visible...),
339 InputTokens: m.agent.estimatedVisibleRequestTokens(visible),
340 ProjectionVersion: m.agent.currentProjectionVersion(),
341 }
342 }
343
344 // estimatedVisibleRequestTokens sizes the pre-interceptor sampling shape:
345 // ModelMessages + role projection + tool schemas. Extension interceptors are
346 // intentionally omitted here (see prepareOnce) to avoid double side effects.
347 func (a *Agent) estimatedVisibleRequestTokens(visible []provider.Message) int {
348 if a == nil {
349 return 0
350 }
351 msgs := a.normalizeModelRequestMessages(visible)
352 tools := a.providerToolSchemas()
353 return a.estimatedRequestTokens(provider.Request{
354 Messages: msgs,
355 Tools: tools,
356 MaxTokens: a.maxOutputTokens,
357 Temperature: provider.OptionalTemperature(a.temperature),
358 })
359 }
360
360 lines GO