diff --git a/docs/plans/v1.1/_files/internal/proxy/cancel_test.go b/docs/plans/v1.1/_files/internal/proxy/cancel_test.go index fba808b..a941eac 100644 --- a/docs/plans/v1.1/_files/internal/proxy/cancel_test.go +++ b/docs/plans/v1.1/_files/internal/proxy/cancel_test.go @@ -84,7 +84,7 @@ hosts = ["alpha"] default_model = "shared" `, alpha) go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot - time.Sleep(100 * time.Millisecond) + waitUntil(t, func() bool { return r.lim.InFlight("alpha", "shared") == 1 }) ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond) defer cancel() req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1))) diff --git a/docs/plans/v1/_files/internal/limiter/limiter_test.go b/docs/plans/v1/_files/internal/limiter/limiter_test.go index 6ab30e2..92b2ddb 100644 --- a/docs/plans/v1/_files/internal/limiter/limiter_test.go +++ b/docs/plans/v1/_files/internal/limiter/limiter_test.go @@ -39,10 +39,7 @@ func TestParallelAndQueue(t *testing.T) { } got3 <- err }() - time.Sleep(20 * time.Millisecond) - if l.Queued("alpha", "m") != 1 { - t.Errorf("queued = %d, want 1", l.Queued("alpha", "m")) - } + waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 }) // Fourth finds the queue full and is refused at once. start := time.Now() _, _, err = l.Acquire(ctx, "alpha", "m") @@ -52,8 +49,8 @@ func TestParallelAndQueue(t *testing.T) { if time.Since(start) > 50*time.Millisecond { t.Errorf("a full queue must refuse immediately, took %v", time.Since(start)) } - time.Sleep(30 * time.Millisecond) - rel1() // frees a slot: the queued third proceeds + time.Sleep(50 * time.Millisecond) // a lower bound on the third's wait, checked above as >= 40 ms + rel1() // frees a slot: the queued third proceeds select { case err := <-got3: if err != nil { @@ -95,7 +92,7 @@ func TestCancelWhileQueuedLeaksNothing(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }() - time.Sleep(20 * time.Millisecond) + waitUntil(t, func() bool { return l.Queued("h", "m") == 1 }) cancel() select { case err := <-done: @@ -139,7 +136,7 @@ func TestQueueIsFIFO(t *testing.T) { time.Sleep(5 * time.Millisecond) r() }(i) - time.Sleep(15 * time.Millisecond) // stagger arrivals so the order is defined + waitUntil(t, func() bool { return l.Queued("h", "m") == i }) // arrivals in order, by observation } rel() wg.Wait() @@ -176,3 +173,16 @@ func TestFreeSlotsSumsModels(t *testing.T) { t.Errorf("unknown host has no slots") } } + +// waitUntil polls cond every millisecond for up to two seconds and fails the test if it never holds. +func waitUntil(t *testing.T, cond func() bool) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(time.Millisecond) + } + t.Fatal("condition not reached within two seconds") +} diff --git a/docs/plans/v1/_files/internal/proxy/proxy_test.go b/docs/plans/v1/_files/internal/proxy/proxy_test.go index 39980aa..897982a 100644 --- a/docs/plans/v1/_files/internal/proxy/proxy_test.go +++ b/docs/plans/v1/_files/internal/proxy/proxy_test.go @@ -70,9 +70,9 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) { for i := 1; i <= 2; i++ { wg.Add(1) go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) - time.Sleep(50 * time.Millisecond) // arrive one after the other so both pick beta (10 > 2) + // arrive one after the other so both pick beta (10 > 2): wait until beta holds i slots + waitUntil(t, func() bool { return r.lim.InFlight("beta", "shared") == i }) } - time.Sleep(50 * time.Millisecond) // …so a third conversation starting now is sent to alpha (beta has 0 free slots, alpha 2). resp := r.post("/r/v1/chat/completions", conversation(3, 1)) drain(resp) diff --git a/docs/plans/v2.1/README.md b/docs/plans/v2.1/README.md index e50fe58..3b256f7 100644 --- a/docs/plans/v2.1/README.md +++ b/docs/plans/v2.1/README.md @@ -12,8 +12,10 @@ - **02-props-loaded-only** (to be written after v2 merges) — the poller asks `/props?model=X` only for models `/v1/models` lists as loaded, because llama-server's router autoloads a model named in that query (`models_autoload`). See the v2 README's note of 2026-09-25. -- **03-timing-margins** (to be written) — the remaining sleep-ordered assertions in v1 and v2 - given tests move to state-based waits, as `TestQueueFullIs503` and `TestParallelAndQueue` did. +- ~~03-timing-margins~~ — done by the owner directly (test-only work, no implementer task): the + remaining ordering sleeps in `limiter_test.go`, `proxy_test.go` and `cancel_test.go` now wait on + limiter state (`Queued`/`InFlight`); the one sleep left is a deliberate lower bound. Five clean + `-race` runs of both packages except the 499 defect task 01 fixes. **How this plan was made:** acceptance tests first; no reference implementation. The given test for task 01 was run against the v2 tree (fails six of six) and against a throwaway fix that diff --git a/internal/limiter/limiter_test.go b/internal/limiter/limiter_test.go index 6ab30e2..92b2ddb 100644 --- a/internal/limiter/limiter_test.go +++ b/internal/limiter/limiter_test.go @@ -39,10 +39,7 @@ func TestParallelAndQueue(t *testing.T) { } got3 <- err }() - time.Sleep(20 * time.Millisecond) - if l.Queued("alpha", "m") != 1 { - t.Errorf("queued = %d, want 1", l.Queued("alpha", "m")) - } + waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 }) // Fourth finds the queue full and is refused at once. start := time.Now() _, _, err = l.Acquire(ctx, "alpha", "m") @@ -52,8 +49,8 @@ func TestParallelAndQueue(t *testing.T) { if time.Since(start) > 50*time.Millisecond { t.Errorf("a full queue must refuse immediately, took %v", time.Since(start)) } - time.Sleep(30 * time.Millisecond) - rel1() // frees a slot: the queued third proceeds + time.Sleep(50 * time.Millisecond) // a lower bound on the third's wait, checked above as >= 40 ms + rel1() // frees a slot: the queued third proceeds select { case err := <-got3: if err != nil { @@ -95,7 +92,7 @@ func TestCancelWhileQueuedLeaksNothing(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }() - time.Sleep(20 * time.Millisecond) + waitUntil(t, func() bool { return l.Queued("h", "m") == 1 }) cancel() select { case err := <-done: @@ -139,7 +136,7 @@ func TestQueueIsFIFO(t *testing.T) { time.Sleep(5 * time.Millisecond) r() }(i) - time.Sleep(15 * time.Millisecond) // stagger arrivals so the order is defined + waitUntil(t, func() bool { return l.Queued("h", "m") == i }) // arrivals in order, by observation } rel() wg.Wait() @@ -176,3 +173,16 @@ func TestFreeSlotsSumsModels(t *testing.T) { t.Errorf("unknown host has no slots") } } + +// waitUntil polls cond every millisecond for up to two seconds and fails the test if it never holds. +func waitUntil(t *testing.T, cond func() bool) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(time.Millisecond) + } + t.Fatal("condition not reached within two seconds") +} diff --git a/internal/proxy/cancel_test.go b/internal/proxy/cancel_test.go index fba808b..a941eac 100644 --- a/internal/proxy/cancel_test.go +++ b/internal/proxy/cancel_test.go @@ -84,7 +84,7 @@ hosts = ["alpha"] default_model = "shared" `, alpha) go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot - time.Sleep(100 * time.Millisecond) + waitUntil(t, func() bool { return r.lim.InFlight("alpha", "shared") == 1 }) ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond) defer cancel() req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1))) diff --git a/internal/proxy/proxy_test.go b/internal/proxy/proxy_test.go index 39980aa..897982a 100644 --- a/internal/proxy/proxy_test.go +++ b/internal/proxy/proxy_test.go @@ -70,9 +70,9 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) { for i := 1; i <= 2; i++ { wg.Add(1) go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) - time.Sleep(50 * time.Millisecond) // arrive one after the other so both pick beta (10 > 2) + // arrive one after the other so both pick beta (10 > 2): wait until beta holds i slots + waitUntil(t, func() bool { return r.lim.InFlight("beta", "shared") == i }) } - time.Sleep(50 * time.Millisecond) // …so a third conversation starting now is sent to alpha (beta has 0 free slots, alpha 2). resp := r.post("/r/v1/chat/completions", conversation(3, 1)) drain(resp)