| 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 |