| 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 | ) |
| 13 | |
| 14 | func TestCompactionPartialCommitThenBudgetFailure(t *testing.T) { |
| 15 | synctest.Test(t, func(t *testing.T) { |
| 16 | a := New(&fakeProvider{reply: "first durable summary"}, nil, foldableSessionOverForce(8), Options{ContextWindow: 32000}, event.Discard) |
| 17 | ctx, finish := a.beginCompactionRun(t.Context()) |
| 18 | if err := a.CompactNow(ctx, ""); err != nil { |
| 19 | t.Fatal(err) |
| 20 | } |
| 21 | version, visible := a.currentProjectionVersion(), a.ModelHistorySnapshot() |
| 22 | if version == 0 { |
| 23 | t.Fatal("first summary did not commit") |
| 24 | } |
| 25 | time.Sleep(4 * time.Minute) |
| 26 | a.svc.prov = &slowSummaryProvider{} |
| 27 | _, _, err := a.runSummaryRequest(ctx, a.summaryRequest(visible, "")) |
| 28 | if err = finish(err); !errors.Is(err, errSummaryBudget) { |
| 29 | t.Fatalf("error=%v", err) |
| 30 | } |
| 31 | if a.currentProjectionVersion() != version || !reflect.DeepEqual(visible, a.ModelHistorySnapshot()) { |
| 32 | t.Fatal("later failure changed the last committed projection") |
| 33 | } |
| 34 | }) |
| 35 | } |
| 36 | |
| 37 | func TestCompactionQueuedStopNeverIssuesRequest(t *testing.T) { |
| 38 | synctest.Test(t, func(t *testing.T) { |
| 39 | p := &slowSummaryProvider{} |
| 40 | a := New(p, nil, foldableSessionOverForce(8), Options{ContextWindow: 32000}, event.Discard) |
| 41 | a.sess.compactionRunMu.Lock() |
| 42 | ctx, cancel := context.WithCancel(t.Context()) |
| 43 | done := make(chan error, 1) |
| 44 | go func() { done <- a.CompactNow(ctx, "") }() |
| 45 | synctest.Wait() |
| 46 | cancel() |
| 47 | if err := <-done; !errors.Is(err, context.Canceled) { |
| 48 | t.Fatal(err) |
| 49 | } |
| 50 | a.sess.compactionRunMu.Unlock() |
| 51 | if p.calls != 0 { |
| 52 | t.Fatal("cancelled queue entry started a request") |
| 53 | } |
| 54 | }) |
| 55 | } |
| 56 | |
| 57 | func TestBudgetRetainsOwnershipUntilIgnoredCancellationSettles(t *testing.T) { |
| 58 | synctest.Test(t, func(t *testing.T) { |
| 59 | p := &lateSummaryProvider{started: make(chan struct{}), release: make(chan struct{})} |
| 60 | a := New(p, nil, foldableSessionOverForce(8), Options{ContextWindow: 32000}, event.Discard) |
| 61 | done := make(chan error, 1) |
| 62 | go func() { done <- a.CompactNow(t.Context(), "") }() |
| 63 | <-p.started |
| 64 | time.Sleep(compactionBudget + time.Second) |
| 65 | synctest.Wait() |
| 66 | select { |
| 67 | case <-done: |
| 68 | t.Fatal("execution ownership released while worker still running") |
| 69 | default: |
| 70 | } |
| 71 | close(p.release) |
| 72 | if err := <-done; !errors.Is(err, errSummaryBudget) { |
| 73 | t.Fatal(err) |
| 74 | } |
| 75 | if a.currentProjectionVersion() != 0 { |
| 76 | t.Fatal("late result installed") |
| 77 | } |
| 78 | }) |
| 79 | } |
| 80 |