Merge v1.1: cancelled requests recorded as 499; empty usage is []
This commit is contained in:
@@ -5,6 +5,7 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once
|
|||||||
|
|
||||||
| Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model |
|
| Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model |
|
||||||
|---|---|---|---|---|---|---|---|
|
|---|---|---|---|---|---|---|---|
|
||||||
|
| v1.1/01-review-fixes | 2026-09-25 | done | 1 | pass | none | Copied `cancel_test.go` and `usage_empty_test.go` byte-identical from `docs/plans/v1.1/_files/`; the earlier session's fixes in `internal/proxy/proxy.go`, `internal/proxy/forward.go` and `internal/admin/admin_ops.go` were already in the working tree. `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b |
|
||||||
| v1/08-smoke-readme | 2026-09-25 | done | 1 | pass | owner-directed fix to `Free` in `proxy.Chooser` | Changed `Free` from `c.lim.FreeSlots(host)` (sum over every model) to per-model free slots, `freeForModel(cfg.Hosts[host], model, c.lim.InFlight(host, model))`, floored at 0 and 0 when the host does not list the model (new helper in hosts.go); the one code change the task directs. `go test -race ./internal/proxy/` and `make gate` pass on the first run; `make smoke` → `smoke: ok (stream spread 1006 ms)`. README intro, `## Configure` (added db/lease_idle/retention, rewrote queue_max and hosts.<name>.hosts) and `## Inspect`→`## Operate` (all six endpoints, examples taken from the smoke run) updated. | llama.cpp/ornith-1.5-35b-a3b |
|
| v1/08-smoke-readme | 2026-09-25 | done | 1 | pass | owner-directed fix to `Free` in `proxy.Chooser` | Changed `Free` from `c.lim.FreeSlots(host)` (sum over every model) to per-model free slots, `freeForModel(cfg.Hosts[host], model, c.lim.InFlight(host, model))`, floored at 0 and 0 when the host does not list the model (new helper in hosts.go); the one code change the task directs. `go test -race ./internal/proxy/` and `make gate` pass on the first run; `make smoke` → `smoke: ok (stream spread 1006 ms)`. README intro, `## Configure` (added db/lease_idle/retention, rewrote queue_max and hosts.<name>.hosts) and `## Inspect`→`## Operate` (all six endpoints, examples taken from the smoke run) updated. | llama.cpp/ornith-1.5-35b-a3b |
|
||||||
| v1/07-main | 2026-09-25 | done | 1 | pass | none | Wired store, limiter and lease table into `cmd/crossbar/main.go`: `store.Open` before the health table, `limiter.Configure` per (host, model) from `cfg.Hosts`, `lease.New` with `proxy.Chooser`, `Candidates` for every route, three background goroutines (idle expiry per minute, prune per hour logging the count, host-health recording per `poll_interval`), and `st.Close` via `defer`. The 3s SIGTERM run exits 0 with `listening`/`shutting down`; the missing-config run exits 1. | llama.cpp/ornith-1.5-35b-a3b |
|
| v1/07-main | 2026-09-25 | done | 1 | pass | none | Wired store, limiter and lease table into `cmd/crossbar/main.go`: `store.Open` before the health table, `limiter.Configure` per (host, model) from `cfg.Hosts`, `lease.New` with `proxy.Chooser`, `Candidates` for every route, three background goroutines (idle expiry per minute, prune per hour logging the count, host-health recording per `poll_interval`), and `st.Close` via `defer`. The 3s SIGTERM run exits 0 with `listening`/`shutting down`; the missing-config run exits 1. | llama.cpp/ornith-1.5-35b-a3b |
|
||||||
| v1/06-admin | 2026-09-25 | done | 2 | fail | Split `internal/admin/admin.go` (196 lines) + `admin_ops.go` (366 lines) to stay under 400. Updated `cmd/crossbar/main.go`'s `admin.Handler` call from the committed 2-arg `(cfg, table)` to the task's 6-arg signature, passing the health table for `hosts` and `nil` for the not-yet-wired `leases`/`limiter`/`store`/`drainer` (task 07 wires them); this was a compile fix required for `go vet`/`go test ./...` on `cmd/crossbar` to pass — the full wiring is task 07. | First `make gate` failed on `go vet` (`admin.Handler` called with 2 args in `main.go` after the signature changed); fixed `main.go` and the gate passed on the second run. `admin_test.go` and `example.toml` verified byte-identical to `docs/plans/v1/_files/`; `internal/lease` and `internal/store` left untouched except the already-present `Candidates`/`StatusCounts`. | llama.cpp/ornith-1.5-35b-a3b |
|
| v1/06-admin | 2026-09-25 | done | 2 | fail | Split `internal/admin/admin.go` (196 lines) + `admin_ops.go` (366 lines) to stay under 400. Updated `cmd/crossbar/main.go`'s `admin.Handler` call from the committed 2-arg `(cfg, table)` to the task's 6-arg signature, passing the health table for `hosts` and `nil` for the not-yet-wired `leases`/`limiter`/`store`/`drainer` (task 07 wires them); this was a compile fix required for `go vet`/`go test ./...` on `cmd/crossbar` to pass — the full wiring is task 07. | First `make gate` failed on `go vet` (`admin.Handler` called with 2 args in `main.go` after the signature changed); fixed `main.go` and the gate passed on the second run. `admin_test.go` and `example.toml` verified byte-identical to `docs/plans/v1/_files/`; `internal/lease` and `internal/store` left untouched except the already-present `Candidates`/`StatusCounts`. | llama.cpp/ornith-1.5-35b-a3b |
|
||||||
@@ -80,3 +81,20 @@ Follow-ups for `v1.1`: fix 1 (record the row on the cancel path with status 499
|
|||||||
(`[]`), and an acceptance test for each; consider `lease_idle` expiry while a request is in flight
|
(`[]`), and an acceptance test for each; consider `lease_idle` expiry while a request is in flight
|
||||||
and `Prune` under concurrent writes, which this review did not probe.
|
and `Prune` under concurrent writes, which this review did not probe.
|
||||||
|
|
||||||
|
### v1.1 review — 2026-09-25 (reviewer: claude, as owner)
|
||||||
|
|
||||||
|
Checked: one task commit `9f5b50a` with the trailer; both given tests byte-identical; protected
|
||||||
|
files untouched; `make gate` → `gate: ok`; `make smoke` → `smoke: ok (stream spread 1006 ms)`.
|
||||||
|
Probed: `/_crossbar/usage` on an empty store answers `[]`; a stream cut by the client after
|
||||||
|
0.4 s appears in `/_crossbar/metrics` as `crossbar_requests_total{…,status="499"} 1`.
|
||||||
|
|
||||||
|
Process: four sessions for one task. Sessions 1–3 each ended their turn right after the sandbox
|
||||||
|
refused a write or read outside the repository (the I9 pattern) — after the fix was already
|
||||||
|
correct, in sessions 2 and 3. Session 3 also chased test flakes caused by its own inference
|
||||||
|
loading the host (the owner measured 12/12 passes idle). Findings: (a) model — five
|
||||||
|
refusal-endings tonight in total; `AGENTS.md` now names the rule, and the fourth session obeyed
|
||||||
|
it; (b) task — the task text did not state that `httputil.ReverseProxy` aborts the handler with
|
||||||
|
`http.ErrAbortHandler` on client disconnect, the fact the fix depends on (added mid-run); (c)
|
||||||
|
test design — timing-based tests (limiter, queue, spread, cancel) have margins tuned for an idle
|
||||||
|
host; widen or retry in a later plan.
|
||||||
|
|
||||||
|
|||||||
@@ -141,6 +141,9 @@ func (hx *handler) usageGet(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, http.StatusInternalServerError, "usage: "+err.Error())
|
writeError(w, http.StatusInternalServerError, "usage: "+err.Error())
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if rows == nil {
|
||||||
|
rows = []store.UsageRow{}
|
||||||
|
}
|
||||||
if r.Header.Get("Accept") == "text/plain" {
|
if r.Header.Get("Accept") == "text/plain" {
|
||||||
writeUsageTable(w, rows)
|
writeUsageTable(w, rows)
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -0,0 +1,31 @@
|
|||||||
|
package admin_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
// An empty usage table is an empty JSON array, not null: clients iterate it.
|
||||||
|
func TestUsageEmptyIsAnArray(t *testing.T) {
|
||||||
|
r := newRig(t)
|
||||||
|
for _, q := range []string{"/_crossbar/usage", "/_crossbar/usage?by=host", "/_crossbar/usage?by=model&since=1h"} {
|
||||||
|
rec := r.do(t, "GET", q, "")
|
||||||
|
if rec.Code != 200 {
|
||||||
|
t.Fatalf("%s: %d", q, rec.Code)
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(rec.Body.String()) != "[]" {
|
||||||
|
t.Errorf("%s: body %q, want []", q, rec.Body.String())
|
||||||
|
}
|
||||||
|
var rows []store.UsageRow
|
||||||
|
if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil || rows == nil || len(rows) != 0 {
|
||||||
|
t.Errorf("%s: decoded %v %v, want an empty non-nil slice", q, rows, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
rec := r.do(t, "GET", "/_crossbar/usage?by=route", "", "Accept", "text/plain")
|
||||||
|
if rec.Code != 200 || !strings.Contains(rec.Body.String(), "key") {
|
||||||
|
t.Errorf("text form with no rows must still print the header: %d %q", rec.Code, rec.Body.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,109 @@
|
|||||||
|
package proxy_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A client that goes away mid-stream is still a request that happened: it held a slot, it cost
|
||||||
|
// prefill, and it belongs in the accounting. The row records status 499 and a non-empty err.
|
||||||
|
func TestClientCancelMidStreamIsRecorded(t *testing.T) {
|
||||||
|
slow := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
switch r.URL.Path {
|
||||||
|
case "/health":
|
||||||
|
fmt.Fprint(w, `{"status":"ok"}`)
|
||||||
|
case "/v1/models":
|
||||||
|
fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"}]}`)
|
||||||
|
default:
|
||||||
|
w.Header().Set("Content-Type", "text/event-stream")
|
||||||
|
w.WriteHeader(200)
|
||||||
|
fmt.Fprint(w, "data: {\"choices\":[{\"delta\":{\"content\":\"first\"}}]}\n\n")
|
||||||
|
w.(http.Flusher).Flush()
|
||||||
|
select {
|
||||||
|
case <-r.Context().Done():
|
||||||
|
case <-time.After(3 * time.Second):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}))
|
||||||
|
t.Cleanup(slow.Close)
|
||||||
|
beta := newUpstream(t, "beta")
|
||||||
|
r := newRig(t, twoHosts, &upstream{name: "alpha", srv: slow}, beta)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
body := `{"model":"alpha-only","stream":true,"messages":[{"role":"user","content":"cancel me"}]}`
|
||||||
|
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(body))
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
resp, err := http.DefaultClient.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
buf := make([]byte, 64)
|
||||||
|
if _, err := resp.Body.Read(buf); err != nil {
|
||||||
|
t.Fatalf("first chunk: %v", err)
|
||||||
|
}
|
||||||
|
cancel()
|
||||||
|
resp.Body.Close()
|
||||||
|
|
||||||
|
deadline := time.Now().Add(3 * time.Second)
|
||||||
|
var counts []store.StatusCount
|
||||||
|
for time.Now().Before(deadline) {
|
||||||
|
counts, _ = r.store.StatusCounts(time.Time{})
|
||||||
|
if len(counts) > 0 {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
time.Sleep(25 * time.Millisecond)
|
||||||
|
}
|
||||||
|
if len(counts) != 1 || counts[0].Status != 499 || counts[0].Route != "r" || counts[0].Count != 1 {
|
||||||
|
t.Fatalf("status counts after a cancelled stream = %+v, want one row: route r, status 499", counts)
|
||||||
|
}
|
||||||
|
rows, _ := r.store.Usage(time.Time{}, store.ByRoute)
|
||||||
|
if len(rows) != 1 || rows[0].Requests != 1 || rows[0].Errors != 1 {
|
||||||
|
t.Errorf("usage = %+v, want 1 request counted as an error", rows)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The same when the client gives up while waiting in the queue: a 499 row, no slot leaked.
|
||||||
|
func TestClientCancelWhileQueuedIsRecorded(t *testing.T) {
|
||||||
|
alpha := newUpstream(t, "alpha")
|
||||||
|
alpha.delay = 800 * time.Millisecond
|
||||||
|
r := newRig(t, `
|
||||||
|
listen = "127.0.0.1:1"
|
||||||
|
queue_max = 2
|
||||||
|
[hosts.alpha]
|
||||||
|
base_url = %q
|
||||||
|
models = { "shared" = { parallel = 1 } }
|
||||||
|
[routes.r]
|
||||||
|
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)
|
||||||
|
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)))
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
if _, err := http.DefaultClient.Do(req); err == nil {
|
||||||
|
t.Fatal("the queued request should have been cancelled by its context")
|
||||||
|
}
|
||||||
|
deadline := time.Now().Add(3 * time.Second)
|
||||||
|
for time.Now().Before(deadline) {
|
||||||
|
counts, _ := r.store.StatusCounts(time.Time{})
|
||||||
|
for _, c := range counts {
|
||||||
|
if c.Status == 499 {
|
||||||
|
if r.lim.Queued("alpha", "shared") != 0 {
|
||||||
|
t.Errorf("queued = %d after the waiter cancelled", r.lim.Queued("alpha", "shared"))
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
time.Sleep(25 * time.Millisecond)
|
||||||
|
}
|
||||||
|
t.Fatal("no 499 row recorded for the request cancelled while queued")
|
||||||
|
}
|
||||||
+47
-17
@@ -30,26 +30,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
|||||||
rev := &forwardState{started: started}
|
rev := &forwardState{started: started}
|
||||||
rp := newReverseProxy(p.health, host, leaseState, target, rest, rev)
|
rp := newReverseProxy(p.health, host, leaseState, target, rest, rev)
|
||||||
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
|
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
|
||||||
|
|
||||||
|
// ServeHTTP unwinds with http.ErrAbortHandler when a client leaves mid-stream; recover so the
|
||||||
|
// row the request earned is still written, then re-panic so the server keeps its semantics.
|
||||||
|
defer func() {
|
||||||
|
if pv := recover(); pv != nil {
|
||||||
|
total := time.Since(started).Milliseconds()
|
||||||
|
req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total)
|
||||||
|
if perr, ok := pv.(error); ok && errors.Is(perr, http.ErrAbortHandler) {
|
||||||
|
req.Status = 499
|
||||||
|
req.Err = "client cancelled"
|
||||||
|
} else {
|
||||||
|
req.Err = "upstream error"
|
||||||
|
}
|
||||||
|
p.writeRecord(req)
|
||||||
|
panic(pv)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
rp.ServeHTTP(rec, r)
|
rp.ServeHTTP(rec, r)
|
||||||
total := time.Since(started)
|
total := time.Since(started)
|
||||||
|
|
||||||
req := store.Request{
|
req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total.Milliseconds())
|
||||||
Route: route,
|
if r.Context().Err() != nil {
|
||||||
FP: fp,
|
req.Status = 499
|
||||||
Model: model,
|
req.Err = "client cancelled"
|
||||||
Host: host,
|
|
||||||
Started: started,
|
|
||||||
QueuedMs: waited.Milliseconds(),
|
|
||||||
TTFBMs: ttfbMs(rev),
|
|
||||||
TotalMs: total.Milliseconds(),
|
|
||||||
Status: rec.status,
|
|
||||||
Streamed: rev.streamed,
|
|
||||||
}
|
|
||||||
if rev.tee != nil {
|
|
||||||
prompt, cached, completion := rev.tee.tokens()
|
|
||||||
req.PromptTokens = int64(prompt)
|
|
||||||
req.CachedTokens = int64(cached)
|
|
||||||
req.CompletionTokens = int64(completion)
|
|
||||||
}
|
}
|
||||||
p.writeRecord(req)
|
p.writeRecord(req)
|
||||||
|
|
||||||
@@ -70,6 +75,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// forwardRow builds the accounting row from the state a forward gathered: what the tee scanned and
|
||||||
|
// what the recorder captured. totalMs is measured from start to the caller's exit, so the forward
|
||||||
|
// path and the recovery path above build identical rows.
|
||||||
|
func forwardRow(route, fp, model, host string, start time.Time, waited time.Duration, rev *forwardState, status int, totalMs int64) store.Request {
|
||||||
|
req := store.Request{
|
||||||
|
Route: route,
|
||||||
|
FP: fp,
|
||||||
|
Model: model,
|
||||||
|
Host: host,
|
||||||
|
Started: start,
|
||||||
|
QueuedMs: waited.Milliseconds(),
|
||||||
|
TTFBMs: ttfbMs(rev),
|
||||||
|
TotalMs: totalMs,
|
||||||
|
Status: status,
|
||||||
|
Streamed: rev.streamed,
|
||||||
|
}
|
||||||
|
if rev.tee != nil {
|
||||||
|
prompt, cached, completion := rev.tee.tokens()
|
||||||
|
req.PromptTokens = int64(prompt)
|
||||||
|
req.CachedTokens = int64(cached)
|
||||||
|
req.CompletionTokens = int64(completion)
|
||||||
|
}
|
||||||
|
return req
|
||||||
|
}
|
||||||
|
|
||||||
// leaseState is "reused" when the lease already held the conversation, else "new".
|
// leaseState is "reused" when the lease already held the conversation, else "new".
|
||||||
func leaseState(reused bool) string {
|
func leaseState(reused bool) string {
|
||||||
if reused {
|
if reused {
|
||||||
|
|||||||
@@ -249,6 +249,16 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
||||||
|
p.writeRecord(store.Request{
|
||||||
|
Route: route,
|
||||||
|
FP: fp,
|
||||||
|
Model: model,
|
||||||
|
Host: host,
|
||||||
|
Started: started,
|
||||||
|
TotalMs: time.Since(started).Milliseconds(),
|
||||||
|
Status: 499,
|
||||||
|
Err: "client cancelled while queued",
|
||||||
|
})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
defer release()
|
defer release()
|
||||||
|
|||||||
Reference in New Issue
Block a user