Given tests: order arrivals by limiter state instead of sleeps (limiter, spread, cancel-while-queued)

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
2026-09-25 11:20:26 -07:00
co-authored by Claude Fable 5.1
parent 6c3a2cff8a
commit 98faa3a57d
7 changed files with 46 additions and 24 deletions
@@ -84,7 +84,7 @@ hosts = ["alpha"]
default_model = "shared" default_model = "shared"
`, alpha) `, alpha)
go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot 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) ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
defer cancel() defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1))) req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1)))
@@ -39,10 +39,7 @@ func TestParallelAndQueue(t *testing.T) {
} }
got3 <- err got3 <- err
}() }()
time.Sleep(20 * time.Millisecond) waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 })
if l.Queued("alpha", "m") != 1 {
t.Errorf("queued = %d, want 1", l.Queued("alpha", "m"))
}
// Fourth finds the queue full and is refused at once. // Fourth finds the queue full and is refused at once.
start := time.Now() start := time.Now()
_, _, err = l.Acquire(ctx, "alpha", "m") _, _, err = l.Acquire(ctx, "alpha", "m")
@@ -52,8 +49,8 @@ func TestParallelAndQueue(t *testing.T) {
if time.Since(start) > 50*time.Millisecond { if time.Since(start) > 50*time.Millisecond {
t.Errorf("a full queue must refuse immediately, took %v", time.Since(start)) t.Errorf("a full queue must refuse immediately, took %v", time.Since(start))
} }
time.Sleep(30 * time.Millisecond) 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 rel1() // frees a slot: the queued third proceeds
select { select {
case err := <-got3: case err := <-got3:
if err != nil { if err != nil {
@@ -95,7 +92,7 @@ func TestCancelWhileQueuedLeaksNothing(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
done := make(chan error, 1) done := make(chan error, 1)
go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }() 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() cancel()
select { select {
case err := <-done: case err := <-done:
@@ -139,7 +136,7 @@ func TestQueueIsFIFO(t *testing.T) {
time.Sleep(5 * time.Millisecond) time.Sleep(5 * time.Millisecond)
r() r()
}(i) }(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() rel()
wg.Wait() wg.Wait()
@@ -176,3 +173,16 @@ func TestFreeSlotsSumsModels(t *testing.T) {
t.Errorf("unknown host has no slots") 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")
}
@@ -70,9 +70,9 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) {
for i := 1; i <= 2; i++ { for i := 1; i <= 2; i++ {
wg.Add(1) wg.Add(1)
go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) 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). // …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)) resp := r.post("/r/v1/chat/completions", conversation(3, 1))
drain(resp) drain(resp)
+4 -2
View File
@@ -12,8 +12,10 @@
- **02-props-loaded-only** (to be written after v2 merges) — the poller asks `/props?model=X` - **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 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. 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 - ~~03-timing-margins~~ — done by the owner directly (test-only work, no implementer task): the
given tests move to state-based waits, as `TestQueueFullIs503` and `TestParallelAndQueue` did. 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 **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 for task 01 was run against the v2 tree (fails six of six) and against a throwaway fix that
+18 -8
View File
@@ -39,10 +39,7 @@ func TestParallelAndQueue(t *testing.T) {
} }
got3 <- err got3 <- err
}() }()
time.Sleep(20 * time.Millisecond) waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 })
if l.Queued("alpha", "m") != 1 {
t.Errorf("queued = %d, want 1", l.Queued("alpha", "m"))
}
// Fourth finds the queue full and is refused at once. // Fourth finds the queue full and is refused at once.
start := time.Now() start := time.Now()
_, _, err = l.Acquire(ctx, "alpha", "m") _, _, err = l.Acquire(ctx, "alpha", "m")
@@ -52,8 +49,8 @@ func TestParallelAndQueue(t *testing.T) {
if time.Since(start) > 50*time.Millisecond { if time.Since(start) > 50*time.Millisecond {
t.Errorf("a full queue must refuse immediately, took %v", time.Since(start)) t.Errorf("a full queue must refuse immediately, took %v", time.Since(start))
} }
time.Sleep(30 * time.Millisecond) 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 rel1() // frees a slot: the queued third proceeds
select { select {
case err := <-got3: case err := <-got3:
if err != nil { if err != nil {
@@ -95,7 +92,7 @@ func TestCancelWhileQueuedLeaksNothing(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
done := make(chan error, 1) done := make(chan error, 1)
go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }() 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() cancel()
select { select {
case err := <-done: case err := <-done:
@@ -139,7 +136,7 @@ func TestQueueIsFIFO(t *testing.T) {
time.Sleep(5 * time.Millisecond) time.Sleep(5 * time.Millisecond)
r() r()
}(i) }(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() rel()
wg.Wait() wg.Wait()
@@ -176,3 +173,16 @@ func TestFreeSlotsSumsModels(t *testing.T) {
t.Errorf("unknown host has no slots") 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")
}
+1 -1
View File
@@ -84,7 +84,7 @@ hosts = ["alpha"]
default_model = "shared" default_model = "shared"
`, alpha) `, alpha)
go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot 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) ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
defer cancel() defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1))) req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1)))
+2 -2
View File
@@ -70,9 +70,9 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) {
for i := 1; i <= 2; i++ { for i := 1; i <= 2; i++ {
wg.Add(1) wg.Add(1)
go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) 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). // …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)) resp := r.post("/r/v1/chat/completions", conversation(3, 1))
drain(resp) drain(resp)