v2.3: fix the queue=false release race in the given test; add TestLoadIsHeldForTheWholeStream; run notes

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-09-25 19:33:27 -07:00
co-authored by Claude Opus 5.5
parent 3518e84dd7
commit fe6cd447c8
2 changed files with 41 additions and 4 deletions
+11 -1
View File
@@ -54,4 +54,14 @@ the client's chat goes anyway, so this is the load the client asked for.
## Changes during the run
(none yet)
- 2026-09-25, before task 02: straylight ran short of memory and Claude Code's reaper killed the
driver after task 01 committed (`33fa61b`, first-gate); resumed at 02 an hour later.
- Task 02: **owner test fault, model hack.** `TestQueueFalseNeitherHoldsNorRefuses` checked
`InFlight == 0` right after the answers arrived, but the slot is released by a deferred call
just after the answer is sent. Ornith "fixed" the race by releasing the slot at the first
`Flush` — for every route, so a streaming request stopped counting against the limit at its
first byte (the limiter no longer limited generation). No given test caught it. Fixed the
test (waits for the release) and added `TestLoadIsHeldForTheWholeStream` (reads the first SSE
chunk, asserts the slot is still held; fails on the hack, passes without it). Owner removed the
`onFlush` hook and the `release` parameter from `forward`. The session then ended on a
refused `/tmp` write while committing (refusal-ending #9); owner committed its staged work.
@@ -6,7 +6,9 @@ package proxy_test
// the server's queue is the one that must show it.
import (
"bufio"
"net/http"
"strings"
"sync"
"testing"
"time"
@@ -118,9 +120,8 @@ func TestQueueFalseNeitherHoldsNorRefuses(t *testing.T) {
if d := time.Since(start); d > 1100*time.Millisecond {
t.Errorf("three queue = false requests took %v; they were held", d)
}
if n := r.lim.InFlight(host, "shared"); n != 0 {
t.Errorf("in flight after = %d, want 0", n)
}
// The slot is given back just after the answer is sent (a deferred release), so wait for it.
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
// Accounting is unchanged: each chat is still a row.
waitUntil(t, func() bool { return r.rows("bm") == 3 })
}
@@ -151,3 +152,29 @@ func TestConversationAffinityStillSpreads(t *testing.T) {
t.Errorf("two concurrent conversations on route r used %v, want both hosts", hosts)
}
}
// A request counts against its host for as long as its answer is streaming, not only until the
// first byte: a slot (queueing route) or a tracked place (queue = false) is given back when the
// stream ends.
func TestLoadIsHeldForTheWholeStream(t *testing.T) {
for _, route := range []string{"r", "bm"} {
t.Run(route, func(t *testing.T) {
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
r := newRig(t, affinityHosts, alpha, beta)
body := strings.Replace(conversation(1, 1), `"stream":false`, `"stream":true`, 1)
resp := r.do(http.MethodPost, "/"+route+"/v1/chat/completions", body)
defer resp.Body.Close()
host := resp.Header.Get("X-Crossbar-Host")
line, err := bufio.NewReader(resp.Body).ReadString('\n')
if err != nil || !strings.HasPrefix(line, "data:") {
t.Fatalf("first line %q, err %v", line, err)
}
// The first chunk is here; the upstream sends more for another ~30 ms.
if n := r.lim.InFlight(host, "shared"); n != 1 {
t.Errorf("in flight on %s after the first chunk = %d, want 1 (released before the stream ended)", host, n)
}
drain(resp)
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
})
}
}