| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "sync" |
| 7 | "testing" |
| 8 | "time" |
| 9 | ) |
| 10 | |
| 11 | func TestRebuildSlotsPriorityOrder(t *testing.T) { |
| 12 | slots := newRebuildSlots(1) |
| 13 | if !slots.tryAcquire() { |
| 14 | t.Fatal("initial tryAcquire failed") |
| 15 | } |
| 16 | if slots.tryAcquire() { |
| 17 | t.Fatal("over-capacity tryAcquire succeeded") |
| 18 | } |
| 19 | |
| 20 | // Queue one waiter per priority while the only slot is held. |
| 21 | order := make([]rebuildPriority, 0, 3) |
| 22 | var mu sync.Mutex |
| 23 | ready := make(chan struct{}) |
| 24 | var wg sync.WaitGroup |
| 25 | start := func(prio rebuildPriority) { |
| 26 | wg.Go(func() { |
| 27 | if err := slots.acquire(context.Background(), prio); err != nil { |
| 28 | t.Errorf("acquire(%d): %v", prio, err) |
| 29 | return |
| 30 | } |
| 31 | mu.Lock() |
| 32 | order = append(order, prio) |
| 33 | mu.Unlock() |
| 34 | slots.release() |
| 35 | ready <- struct{}{} |
| 36 | }) |
| 37 | } |
| 38 | start(rebuildPriorityPrefetch) |
| 39 | time.Sleep(20 * time.Millisecond) // let the prefetch waiter queue first |
| 40 | start(rebuildPrioritySearch) |
| 41 | start(rebuildPriorityUser) |
| 42 | time.Sleep(50 * time.Millisecond) |
| 43 | |
| 44 | slots.release() |
| 45 | for range 3 { |
| 46 | select { |
| 47 | case <-ready: |
| 48 | case <-time.After(5 * time.Second): |
| 49 | t.Fatal("waiter never granted") |
| 50 | } |
| 51 | } |
| 52 | wg.Wait() |
| 53 | mu.Lock() |
| 54 | defer mu.Unlock() |
| 55 | // Highest priority wins each freed slot regardless of arrival order. |
| 56 | if len(order) != 3 || order[0] != rebuildPriorityUser || order[1] != rebuildPrioritySearch || order[2] != rebuildPriorityPrefetch { |
| 57 | t.Fatalf("grant order %v, want user,search,prefetch", order) |
| 58 | } |
| 59 | } |
| 60 | |
| 61 | func TestRebuildSlotsCancelDoesNotDropGrant(t *testing.T) { |
| 62 | slots := newRebuildSlots(1) |
| 63 | if !slots.tryAcquire() { |
| 64 | t.Fatal("initial tryAcquire failed") |
| 65 | } |
| 66 | ctx, cancel := context.WithCancel(context.Background()) |
| 67 | errCh := make(chan error, 1) |
| 68 | go func() { |
| 69 | errCh <- slots.acquire(ctx, rebuildPriorityUser) |
| 70 | }() |
| 71 | time.Sleep(20 * time.Millisecond) |
| 72 | cancel() |
| 73 | slots.release() |
| 74 | select { |
| 75 | case err := <-errCh: |
| 76 | // cancel precedes release, so either interleaving is correct; what must |
| 77 | // never happen is a lost grant. A cancelled waiter leaves the freed |
| 78 | // slot available, and a successful one owns it and must release it. |
| 79 | if err == nil { |
| 80 | slots.release() |
| 81 | break |
| 82 | } |
| 83 | if !errors.Is(err, context.Canceled) { |
| 84 | t.Fatalf("cancelled acquire returned %v, want context.Canceled", err) |
| 85 | } |
| 86 | if !slots.tryAcquire() { |
| 87 | t.Fatal("cancelled acquire consumed the freed slot") |
| 88 | } |
| 89 | slots.release() |
| 90 | case <-time.After(5 * time.Second): |
| 91 | t.Fatal("acquire never returned") |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | func TestRebuildSlotsCancelRemovesWaiter(t *testing.T) { |
| 96 | slots := newRebuildSlots(1) |
| 97 | if !slots.tryAcquire() { |
| 98 | t.Fatal("initial tryAcquire failed") |
| 99 | } |
| 100 | ctx, cancel := context.WithCancel(context.Background()) |
| 101 | errCh := make(chan error, 1) |
| 102 | go func() { |
| 103 | errCh <- slots.acquire(ctx, rebuildPrioritySearch) |
| 104 | }() |
| 105 | time.Sleep(20 * time.Millisecond) |
| 106 | cancel() |
| 107 | if err := <-errCh; err == nil { |
| 108 | t.Fatal("cancelled waiter acquired without grant") |
| 109 | } |
| 110 | slots.release() |
| 111 | // The queue must be empty: the released slot is free for the next caller. |
| 112 | if !slots.tryAcquire() { |
| 113 | t.Fatal("cancelled waiter still queued in front of new callers") |
| 114 | } |
| 115 | slots.release() |
| 116 | } |
| 117 |