返回 DeepSeek-Reasonix
compaction_run_test.go
根目录 / internal / agent / compaction_run_test.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "reflect"
7 "testing"
8 "testing/synctest"
9 "time"
10
11 "reasonix/internal/event"
12 "reasonix/internal/provider"
13 )
14
15 type slowSummaryProvider struct {
16 output bool
17 calls int
18 }
19
20 type overflowThenSlowSummary struct {
21 requests int
22 summary slowSummaryProvider
23 }
24
25 func (*overflowThenSlowSummary) Name() string { return "overflow-summary" }
26 func (p *overflowThenSlowSummary) Stream(ctx context.Context, req provider.Request) (<-chan provider.Chunk, error) {
27 p.requests++
28 if p.requests == 1 {
29 return nil, &provider.ContextLimitError{}
30 }
31 return p.summary.Stream(ctx, req)
32 }
33
34 func TestSamplingOverflowPreservesSummaryBudgetFailure(t *testing.T) {
35 synctest.Test(t, func(t *testing.T) {
36 p := &overflowThenSlowSummary{}
37 sess := foldableSessionOverForce(6)
38 before := sess.Snapshot()
39 a := New(p, nil, sess, Options{ContextWindow: 32000, TaskBudget: TaskBudget{Wall: 10 * time.Minute}}, event.Discard)
40 started := time.Now()
41 result := a.streamWithSamplingRecovery(t.Context(), 1)
42 var summary *SummaryError
43 if !errors.Is(result.err, ErrCompactionRequired) || !errors.As(result.err, &summary) || summary.Code != "summary_budget_exceeded" {
44 t.Fatalf("sampling lost the compaction cause: %v", result.err)
45 }
46 if p.requests != 2 || time.Since(started) != compactionBudget {
47 t.Fatalf("requests=%d elapsed=%s", p.requests, time.Since(started))
48 }
49 if !reflect.DeepEqual(before, sess.Snapshot()) || a.currentProjectionVersion() != 0 {
50 t.Fatal("failed overflow recovery changed history or projection")
51 }
52 })
53 }
54
55 func (*slowSummaryProvider) Name() string { return "slow-summary" }
56 func (p *slowSummaryProvider) Stream(ctx context.Context, _ provider.Request) (<-chan provider.Chunk, error) {
57 p.calls++
58 ch := make(chan provider.Chunk)
59 go func() {
60 defer close(ch)
61 tick := time.NewTicker(10 * time.Second)
62 defer tick.Stop()
63 for {
64 select {
65 case <-ctx.Done():
66 return
67 case <-tick.C:
68 chunk := provider.Chunk{Type: provider.ChunkUsage}
69 if p.output {
70 chunk = provider.Chunk{Type: provider.ChunkReasoning, Text: "working"}
71 }
72 select {
73 case ch <- chunk:
74 case <-ctx.Done():
75 return
76 }
77 }
78 }
79 }()
80 return ch, nil
81 }
82
83 func TestCompactionBudgetIncludesContinuousOutputAndSilentStreams(t *testing.T) {
84 for _, output := range []bool{false, true} {
85 for _, trigger := range []string{CompactionTriggerManual, CompactionTriggerPressure, CompactionTriggerOverflow} {
86 t.Run(trigger+map[bool]string{false: "-silent", true: "-reasoning"}[output], func(t *testing.T) {
87 synctest.Test(t, func(t *testing.T) {
88 p := &slowSummaryProvider{output: output}
89 sess := foldableSessionOverForce(6)
90 before := sess.Snapshot()
91 a := New(p, nil, sess, Options{ContextWindow: 5000, CompactRatio: .5}, event.Discard)
92 parent := context.Background()
93 start := time.Now()
94 _, err := a.contextManager().Prepare(parent, ContextPreparePolicy{Trigger: trigger})
95 if trigger == CompactionTriggerPressure {
96 if err != nil {
97 t.Fatalf("safe pressure must continue: %v", err)
98 }
99 } else {
100 var summary *SummaryError
101 if !errors.As(err, &summary) || summary.Code != "summary_budget_exceeded" {
102 t.Fatalf("error = %v", err)
103 }
104 if trigger == CompactionTriggerOverflow && !errors.Is(err, ErrCompactionRequired) {
105 t.Fatal("budget failure lost the blocked-context classification")
106 }
107 }
108 if time.Since(start) != compactionBudget || parent.Err() != nil || p.calls != 1 {
109 t.Fatalf("elapsed=%s parent=%v calls=%d", time.Since(start), parent.Err(), p.calls)
110 }
111 if !reflect.DeepEqual(before, sess.Snapshot()) || a.currentProjectionVersion() != 0 {
112 t.Fatal("timeout changed history")
113 }
114 })
115 })
116 }
117 }
118 }
119
120 func TestCompactionBudgetCancelsQueuedOperationBeforeGateRelease(t *testing.T) {
121 synctest.Test(t, func(t *testing.T) {
122 p := &slowSummaryProvider{}
123 a := New(p, nil, foldableSessionOverForce(6), Options{ContextWindow: 5000}, event.Discard)
124 a.sess.compactionRunMu.Lock()
125 finished := make(chan error, 1)
126 go func() { finished <- a.CompactNow(context.Background(), "") }()
127 synctest.Wait()
128 time.Sleep(compactionBudget)
129 synctest.Wait()
130 select {
131 case err := <-finished:
132 if !errors.Is(err, errSummaryBudget) {
133 t.Fatal(err)
134 }
135 default:
136 t.Fatal("expired waiter still needs the gate")
137 }
138 a.sess.compactionRunMu.Unlock()
139 if p.calls != 0 {
140 t.Fatal("expired waiter issued a request")
141 }
142 })
143 }
144
145 func TestNestedCompactionRunKeepsOriginalBudget(t *testing.T) {
146 synctest.Test(t, func(t *testing.T) {
147 a := New(nil, nil, NewSession("sys"), Options{}, event.Discard)
148 ctx, finish := a.beginCompactionRun(context.Background())
149 time.Sleep(4 * time.Minute)
150 nested, end := a.beginCompactionRun(ctx)
151 deadline, _ := nested.Deadline()
152 if time.Until(deadline) != time.Minute || currentCompactionRun(ctx) != currentCompactionRun(nested) {
153 t.Fatal("nested operation reset its budget")
154 }
155 _ = end(nil)
156 _ = finish(nil)
157 })
158 }
159
159 lines GO