返回 DeepSeek-Reasonix
context_recovery.go
根目录 / internal / agent / context_recovery.go
1 package agent
2
3 import (
4 "context"
5 "fmt"
6
7 "reasonix/internal/event"
8 "reasonix/internal/i18n"
9 "reasonix/internal/provider"
10 )
11
12 type contextRecoveryBudget struct {
13 retries int
14 failure error
15 }
16
17 func (a *Agent) recoverContextLimit(ctx context.Context, frozen samplingRequest, err error, budget *contextRecoveryBudget) (samplingRequest, bool, string) {
18 limit := provider.AsContextLimitError(err)
19 if a == nil || limit == nil || budget == nil {
20 return samplingRequest{}, false, contextRecoveryFailed
21 }
22 omitted := frozen.req.MaxTokens == 0
23 if limit.PromptTokens > 0 {
24 a.setPromptTokenCalibrationFromActive(limit.PromptTokens)
25 }
26 a.learnContextBudget(limit.WindowTokens, limit.CompletionTokens, omitted)
27 adm := a.lastAdmission()
28 adm.ObservedWindow = limit.WindowTokens
29 adm.ObservedPrompt = limit.PromptTokens
30 adm.ObservedCompletion = limit.CompletionTokens
31 a.storeAdmission(adm)
32
33 window := a.effectiveContextWindow()
34 prompt := limit.PromptTokens
35 if prompt <= 0 {
36 prompt = a.estimatedRequestTokens(frozen.req)
37 }
38 physical := window - prompt - outputBudgetReserve
39 // An overflow without token numbers cannot size a retry: the estimate that
40 // admitted the request is the number the provider just rejected.
41 if limit.PromptTokens <= 0 && limit.WindowTokens <= 0 {
42 physical = 0
43 }
44 if physical > 0 && budget.retries == 0 {
45 next := freezeProviderRequest(frozen.req)
46 next.MaxTokens = physical
47 if frozen.req.MaxTokens > 0 && frozen.req.MaxTokens < physical {
48 next.MaxTokens = frozen.req.MaxTokens
49 }
50 budget.retries++
51 // Publish the request that will actually be retried, not the stale
52 // pre-error admission. The Context Panel reads this atomic snapshot while
53 // the turn is still active and after it completes.
54 adm.WindowMode = provider.ContextWindowShared.String()
55 adm.Source = provider.ContextBudgetSourceLearned
56 adm.WindowTokens = window
57 adm.PromptTokens = prompt
58 adm.PhysicalRemaining = physical
59 if adm.RequestedOutputTokens <= 0 {
60 adm.RequestedOutputTokens = limit.CompletionTokens
61 }
62 if omitted && adm.AutoOutputTokens <= 0 {
63 adm.AutoOutputTokens = limit.CompletionTokens
64 }
65 adm.EffectiveOutputTokens = next.MaxTokens
66 adm.Clipped = adm.RequestedOutputTokens > 0 && next.MaxTokens < adm.RequestedOutputTokens
67 adm.ApplyMaxTokens = next.MaxTokens > 0
68 adm.LastRecovery = contextRecoveryLearnedRetry
69 a.storeAdmission(adm)
70 a.emitContextRecoveryNotice(contextRecoveryLearnedRetry, limit, next.MaxTokens)
71 shape := a.requestCalibrationShape(next)
72 a.sess.output.activeReqShape.Store(&shape)
73 return samplingRequest{req: next}, true, contextRecoveryLearnedRetry
74 }
75 if physical <= 0 && budget.retries == 0 {
76 work, finish := a.beginCompactionRun(ctx)
77 defer func() { budget.failure = finish(budget.failure) }()
78 ctx = work
79 startProjectionVersion := a.currentProjectionVersion()
80 if _, perr := a.contextManager().Prepare(ctx, ContextPreparePolicy{
81 Trigger: CompactionTriggerOverflow,
82 Force: true,
83 }); perr != nil {
84 budget.failure = perr
85 a.setLastRecovery(contextRecoveryFailed)
86 return samplingRequest{}, false, contextRecoveryFailed
87 }
88 if a.currentProjectionVersion() <= startProjectionVersion {
89 budget.failure = fmt.Errorf("%w: %w", ErrCompactionRequired, summaryError(errCheckpointRejected))
90 a.setLastRecovery(contextRecoveryFailed)
91 return samplingRequest{}, false, contextRecoveryFailed
92 }
93 rebuilt, rerr := a.buildSamplingRequest(ctx, CompactionTriggerPressure)
94 if rerr != nil {
95 budget.failure = rerr
96 a.setLastRecovery(contextRecoveryFailed)
97 return samplingRequest{}, false, contextRecoveryFailed
98 }
99 if aerr := a.applyAdmissionToRequest(&rebuilt.req); aerr != nil {
100 budget.failure = aerr
101 a.setLastRecovery(contextRecoveryFailed)
102 return samplingRequest{}, false, contextRecoveryFailed
103 }
104 budget.retries++
105 a.setLastRecovery(contextRecoveryCompacted)
106 a.emitContextRecoveryNotice(contextRecoveryCompacted, limit, rebuilt.req.MaxTokens)
107 shape := a.requestCalibrationShape(rebuilt.req)
108 a.sess.output.activeReqShape.Store(&shape)
109 return samplingRequest{req: freezeProviderRequest(rebuilt.req)}, true, contextRecoveryCompacted
110 }
111 a.setLastRecovery(contextRecoveryFailed)
112 return samplingRequest{}, false, contextRecoveryFailed
113 }
114
115 func (a *Agent) emitContextRecoveryNotice(kind string, limit *provider.ContextLimitError, nextOutput int) {
116 if a == nil || a.svc.sink == nil {
117 return
118 }
119 text := i18n.M.ContextRecoveryAdjustBudget
120 if kind == contextRecoveryCompacted {
121 text = i18n.M.ContextRecoveryCompacted
122 }
123 detail := fmt.Sprintf("recovery=%s next_output=%d", kind, nextOutput)
124 if limit != nil {
125 detail = fmt.Sprintf("%s window=%d prompt=%d completion=%d requested=%d",
126 detail, limit.WindowTokens, limit.PromptTokens, limit.CompletionTokens, limit.RequestedTokens)
127 }
128 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: text, Detail: detail})
129 }
130
130 lines GO