From cf2aa243939e472df3936ab362d63976f727f688 Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 06:56:54 -0700 Subject: [PATCH] Smoke run for v1; README for leases, admin and accounting Implemented-By: OpenCode session (model recorded in docs/implementer-log.md) --- README.md | 103 ++++++++++++++++++++++++++++++--------- cmd/fakeupstream/main.go | 30 +++++++++--- docs/implementer-log.md | 1 + internal/proxy/hosts.go | 16 +++++- tools/smoke.sh | 83 ++++++++++++++++++++++--------- 5 files changed, 178 insertions(+), 55 deletions(-) diff --git a/README.md b/README.md index 9cf69eb..c6fc69c 100644 --- a/README.md +++ b/README.md @@ -1,8 +1,10 @@ # crossbar crossbar is an affinity router in front of several `llama-server` routers. A client's identity is -the first path segment of its base URL; v0 routes each request to the first healthy host on that -route's list and streams the answer back unbuffered. +the first path segment of its base URL — its route. Each conversation takes a sticky lease on one +host, chosen for the most free slots for its model times weight, and streams the answer back +incrementally with the usage chunk intact. Pins, drains, queueing, leases and accounting are all +new in v1. ## Build @@ -16,22 +18,26 @@ streaming over real HTTP. crossbar reads one TOML file. This is `example.toml`: ```toml -# crossbar example configuration. Replace and the addresses with your own. +# crossbar example configuration (v1). Replace and the addresses with your own. listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production +db = "crossbar.db" # SQLite: leases + accounting (WAL). /var/lib/crossbar/crossbar.db under systemd poll_interval = "1s" # 60s in production; 1s makes the smoke run quick -queue_max = 8 +lease_idle = "30m" # a conversation idle this long loses its host +retention = "180d" # per-request rows older than this are rolled up daily +queue_max = 1 # waiting places per (host, model) beyond `parallel`; 503 past that [hosts.alpha] base_url = "http://127.0.0.1:18081" # e.g. http://straylight.:11434 weight = 1.0 -models = { "ornith-1.5-35b-a3b" = { parallel = 4 }, "small-9b" = { parallel = 6 } } +models = { "ornith-1.5-35b-a3b" = { parallel = 1 }, "small-9b" = { parallel = 6 } } [hosts.beta] base_url = "http://127.0.0.1:18082" # e.g. http://titan.:8081 weight = 2.0 -models = { "ornith-1.5-35b-a3b" = { parallel = 4 } } +models = { "ornith-1.5-35b-a3b" = { parallel = 2 } } -# v0: a route is a preference list; the first healthy host that has the model wins. +# v1: a route is a set of candidate hosts; each conversation gets a sticky lease on the host with +# the most free slots × weight at the time it starts. Pins and drains come from the admin API. [routes.opencode-a] hosts = ["alpha", "beta"] default_model = "ornith-1.5-35b-a3b" @@ -43,12 +49,15 @@ hosts = ["beta", "alpha"] | Key | Meaning | | --- | --- | | `listen` | Where crossbar binds. A tailnet address, never `0.0.0.0`. | +| `db` | SQLite file holding leases and the accounting rows. | +| `lease_idle` | A conversation idle this long loses its host. | +| `retention` | Per-request rows older than this are rolled up daily. | | `poll_interval` | How often each host is health-checked. 60s in production; 1s makes the smoke run quick. | -| `queue_max` | Reserved for v1 queueing; no effect in v0. | +| `queue_max` | Waiting places per (host, model) beyond `parallel`; a full queue returns 503. | | `hosts..base_url` | The llama-server base URL this host serves. | -| `hosts..weight` | Relative share of new routes this host receives. | +| `hosts..weight` | Relative share of new requests this host receives. | | `hosts..models` | The models this host serves, with per-model parallel tuning. | -| `routes..hosts` | Preference order: the first healthy host that serves the model wins. | +| `routes..hosts` | Candidate hosts, tried in order until one is healthy; a conversation leases one of them. | | `routes..default_model` | Model used when a request omits one; must be served by a host in the route. | ## Run @@ -84,23 +93,73 @@ custom_providers: models: { ornith-1.5-35b-a3b: {} } ``` -The route name in the URL must exist in `[routes]`; unknown routes are 404. +The route name in the URL must exist in `[routes]`; unknown routes are 404. A client may instead +name the route on an `X-Crossbar-Route` header and point at the bare `/v1` base: -## Inspect - -`GET /_crossbar/hosts` reports every host's health and loaded models: - -```json -{"alpha":{"healthy":true,"loaded":["ornith-1.5-35b-a3b","small-9b"],"last_ok":"2026-09-25T09:34:18Z","last_err":""},"beta":{"healthy":true,"loaded":["ornith-1.5-35b-a3b"],"last_ok":"2026-09-25T09:34:18Z","last_err":""}} +```sh +curl -H 'X-Crossbar-Route: opencode-a' \ + https://crossbar.:7777/v1/chat/completions ``` -`GET /_crossbar/routes` reports each route's preference order and default model: +## Operate + +The operator's API lives under `/_crossbar/`. Every call returns 200 with a small JSON body unless +stated otherwise. + +`GET /_crossbar/hosts` reports every host's health, loaded models, live concurrency from the +limiter and drain state: ```json -{"hermes-x":{"hosts":["beta","alpha"],"default_model":""},"opencode-a":{"hosts":["alpha","beta"],"default_model":"ornith-1.5-35b-a3b"}} +{"alpha":{"healthy":true,"loaded":["ornith-1.5-35b-a3b","small-9b"],"last_ok":"2026-09-25T13:53:25Z","last_err":"","free_slots":7,"in_flight":0,"queued":0,"draining":false},"beta":{"healthy":true,"loaded":["ornith-1.5-35b-a3b"],"last_ok":"2026-09-25T13:53:25Z","last_err":"","free_slots":2,"in_flight":0,"queued":0,"draining":false}} ``` -## What v0 does not do +`GET /_crossbar/routes` reports each route's candidate hosts, default model, any pin and its live +leases: -Leases and stickiness, SQLite, `/slots`, queueing and wake-on-LAN are out of scope for v0; see -`PLAN.md`. +```json +{"hermes-x":{"hosts":["beta","alpha"],"default_model":"","pinned":"","leases":[]},"opencode-a":{"hosts":["alpha","beta"],"default_model":"ornith-1.5-35b-a3b","pinned":"","leases":[]}} +``` + +`POST /_crossbar/routes/{route}` pins a route to a host (`{"host":"alpha","pin":true}`) or releases +it and clears the pin (`{"release":true}`): + +```json +{"ok":true} +``` + +`POST /_crossbar/hosts/{host}` sets or clears drain (`{"drain":true}`); a draining host takes no +new conversations but keeps its existing leases: + +```json +{"ok":true} +``` + +`GET /_crossbar/usage` summarizes the accounting rows, grouped by `by=host`, `by=model` or +`by=route` (the default). Ask for JSON, or a fixed-width table with `Accept: text/plain`: + +```json +[{"key":"beta","requests":2,"errors":0,"busy_ms":4,"queued_ms":0,"prompt_tokens":200,"cached_tokens":180,"completion_tokens":20}] +``` + +``` +key requests errors busy_ms queued_ms prompt cached completion cache_hit +hermes-x 1 0 1 0 100 90 10 0.90 +opencode-a 1 0 3 0 100 90 10 0.90 +``` + +`GET /_crossbar/metrics` emits the Prometheus text exposition for request counts, token totals, +queue wait, host health and live slots: + +``` +# TYPE crossbar_requests_total counter +crossbar_requests_total{route="hermes-x",host="beta",status="200"} 1 +crossbar_requests_total{route="opencode-a",host="beta",status="200"} 1 +# TYPE crossbar_host_healthy gauge +crossbar_host_healthy{host="alpha"} 1 +crossbar_host_healthy{host="beta"} 1 +``` + +## What v1 does not do + +The context-size guard, wake-on-LAN, Tailscale identity and `/slots` are out of scope for v1; see +`PLAN.md` v2. diff --git a/cmd/fakeupstream/main.go b/cmd/fakeupstream/main.go index dbb7801..c222d38 100644 --- a/cmd/fakeupstream/main.go +++ b/cmd/fakeupstream/main.go @@ -1,11 +1,13 @@ // fakeupstream stands in for a llama-server router in tests and the smoke run. Do not edit. // -// fakeupstream -listen 127.0.0.1:18081 -name alpha -models a,b -down-file /tmp/alpha.down +// fakeupstream -listen 127.0.0.1:18081 -name alpha -models a,b -down-file /tmp/alpha.down -slow 0 // // /health answers 503 while the down file exists, 200 otherwise. /v1/models lists -models. // /props answers a small JSON object. /v1/chat/completions echoes: a streamed answer of five -// SSE chunks 200 ms apart when the body has "stream": true, one JSON answer otherwise. Every -// response carries X-Upstream: . +// SSE chunks 200 ms apart when the body has "stream": true, then a final chunk carrying +// "usage" and llama-server style "timings", then [DONE]; one JSON answer with usage and +// timings otherwise. -slow adds that many milliseconds before answering (for queue tests). +// Every response carries X-Upstream: . package main import ( @@ -25,11 +27,14 @@ func main() { name := flag.String("name", "fake", "name reported in X-Upstream and answers") models := flag.String("models", "m", "comma-separated model ids for /v1/models") downFile := flag.String("down-file", "", "while this file exists, /health answers 503") + slow := flag.Int("slow", 0, "milliseconds to wait before answering a completion") flag.Parse() ids := strings.Split(*models, ",") mux := http.NewServeMux() stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) } + usage := map[string]any{"prompt_tokens": 100, "completion_tokens": 10, "total_tokens": 110} + timings := map[string]any{"prompt_n": 100, "cache_n": 90, "predicted_n": 10, "predicted_ms": 50.0} mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { stamp(w) @@ -61,11 +66,12 @@ func main() { Stream bool `json:"stream"` } _ = json.Unmarshal(body, &req) + time.Sleep(time.Duration(*slow) * time.Millisecond) if !req.Stream { writeJSON(w, map[string]any{ "id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model, "choices": []map[string]any{{"index": 0, "message": map[string]string{"role": "assistant", "content": "hello from " + *name}, "finish_reason": "stop"}}, - "usage": map[string]int{"prompt_tokens": 3, "completion_tokens": 3, "total_tokens": 6}, + "usage": usage, "timings": timings, }) return } @@ -73,16 +79,24 @@ func main() { w.Header().Set("Cache-Control", "no-cache") w.WriteHeader(http.StatusOK) fl, _ := w.(http.Flusher) + flush := func() { + if fl != nil { + fl.Flush() + } + } for i := 1; i <= 5; i++ { chunk := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model, "choices": []map[string]any{{"index": 0, "delta": map[string]string{"content": fmt.Sprintf("%s chunk %d ", *name, i)}}}} b, _ := json.Marshal(chunk) fmt.Fprintf(w, "data: %s\n\n", b) - if fl != nil { - fl.Flush() - } + flush() time.Sleep(200 * time.Millisecond) } + final := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model, + "choices": []map[string]any{}, "usage": usage, "timings": timings} + b, _ := json.Marshal(final) + fmt.Fprintf(w, "data: %s\n\n", b) + flush() fmt.Fprint(w, "data: [DONE]\n\n") }) mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { @@ -90,7 +104,7 @@ func main() { http.Error(w, `{"error":"not found"}`, http.StatusNotFound) }) - log.Printf("fakeupstream %s listening on %s models=%v", *name, *listen, ids) + log.Printf("fakeupstream %s listening on %s models=%v slow=%dms", *name, *listen, ids, *slow) srv := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second} log.Fatal(srv.ListenAndServe()) } diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 332e1b1..3344ed5 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -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 | |---|---|---|---|---|---|---|---| +| 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..hosts) and `## Inspect`→`## Operate` (all six endpoints, examples taken from the smoke run) updated. | ? | | 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. | ? | | 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`. | ? | | v1/05-proxy | 2026-09-25 | done | 2 | fail | Split `internal/proxy/proxy.go` (411 lines) into `proxy.go` + `forward.go` by moving `forward`, `newReverseProxy`, `forwardState`, `statusRecorder`, `leaseState`, `ttfbMs` and the `writeError`/`writeRecord` helpers to `forward.go`; the one `recorder_test.go` `proxy.New` call changed to `proxy.New(cfg, h, nil, nil, nil, nil)` per the task; `cmd/crossbar/main.go` passes `nil, nil, nil` for the new `leases`/`lim`/`rec` args (task 06 wires them). | The tee in `tee.go` already read the final SSE chunk's (streamed) and the JSON body's (non-streamed) usage/timings, so `TestAccountingRowsFromUsageAndTimings` passed on the first run — the only gate blocker was `proxy.go` at 411 lines. | ? | diff --git a/internal/proxy/hosts.go b/internal/proxy/hosts.go index c3b0bc0..93bc65b 100644 --- a/internal/proxy/hosts.go +++ b/internal/proxy/hosts.go @@ -75,7 +75,7 @@ func (c *hostChooser) Choose(candidates []string, model string) (string, bool) { Draining: false, Loaded: contains(s.Loaded, model), CanServe: c.cfg.Serves(host, model), - Free: c.lim.FreeSlots(host), + Free: freeForModel(c.cfg.Hosts[host], model, c.lim.InFlight(host, model)), Queued: c.lim.Queued(host, model), Weight: c.cfg.Hosts[host].Weight, } @@ -91,3 +91,17 @@ func contains(list []string, v string) bool { } return false } + +// freeForModel returns the free slots for one model on one host: its parallel minus the in-flight +// count, floored at zero, and zero when the host does not list that model. +func freeForModel(h config.Host, model string, inflight int) int { + m, ok := h.Models[model] + if !ok { + return 0 + } + free := m.Parallel - inflight + if free < 0 { + return 0 + } + return free +} diff --git a/tools/smoke.sh b/tools/smoke.sh index 4b21792..7c2bcb4 100755 --- a/tools/smoke.sh +++ b/tools/smoke.sh @@ -1,45 +1,80 @@ #!/bin/sh -# Smoke run: two fake upstreams, one crossbar, real HTTP. Prints "smoke: ok" or fails. -# Needs: bin/crossbar and bin/fakeupstream (make build), curl. +# Smoke run (v1): two fake upstreams, one crossbar with a fresh SQLite file, real HTTP. +# Checks routing, leases (sticky + header), failover, recovery, streaming, queueing, pin, drain, +# usage and metrics. Prints "smoke: ok" or fails with the crossbar log. set -eu cd "$(dirname "$0")/.." tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM pids="" -bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" >"$tmp/alpha.log" 2>&1 & pids="$pids $!" +sed "s#^db .*#db = \"$tmp/crossbar.db\"#" example.toml > "$tmp/crossbar.toml" +bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" -slow 600 >"$tmp/alpha.log" 2>&1 & pids="$pids $!" bin/fakeupstream -listen 127.0.0.1:18082 -name beta -models ornith-1.5-35b-a3b -down-file "$tmp/beta.down" >"$tmp/beta.log" 2>&1 & pids="$pids $!" -bin/crossbar -config example.toml >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" +bin/crossbar -config "$tmp/crossbar.toml" >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" sleep 1.5 fail() { echo "smoke: FAIL: $*" >&2; echo "--- crossbar.log"; cat "$tmp/crossbar.log"; exit 1; } base=http://127.0.0.1:17777 +conv() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"system","content":"smoke"},{"role":"user","content":"conversation %s"}]}' "$1"; } +hdrs() { curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host} %header{X-Crossbar-Lease}' "$@"; } -h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/opencode-a/v1/models") -[ "$h" = "200 alpha" ] || fail "opencode-a should go to alpha, got '$h'" -h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/hermes-x/v1/models") -[ "$h" = "200 beta" ] || fail "hermes-x should go to beta, got '$h'" -h=$(curl -s -o /dev/null -w '%{http_code}' "$base/nope/v1/models") -[ "$h" = "404" ] || fail "unknown route should be 404, got '$h'" +# 1. a conversation gets a lease and keeps it; beta wins (2 slots × weight 2 vs 1 × 1) +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta new" ] || fail "first turn should be '200 beta new', got '$h'" +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta reused" ] || fail "second turn should reuse beta, got '$h'" -touch "$tmp/alpha.down"; sleep 2.5 # poll_interval is 1s in example.toml -h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/opencode-a/v1/models") -[ "$h" = "200 beta" ] || fail "with alpha down, opencode-a should fail over to beta, got '$h'" -curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{"healthy":false' || fail "/_crossbar/hosts does not show alpha unhealthy: $(curl -s $base/_crossbar/hosts)" +# 2. header route +h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Route: hermes-x' -d "$(conv B)" "$base/v1/chat/completions") +case "$h" in "200 beta new") ;; *) fail "header route hermes-x should be '200 beta new', got '$h'";; esac +h=$(curl -s -o /dev/null -w '%{http_code}' "$base/nope/v1/models"); [ "$h" = "404" ] || fail "unknown route 404, got $h" -rm "$tmp/alpha.down"; sleep 3.5 # recovery needs two good polls -h=$(curl -s -o /dev/null -w '%header{X-Crossbar-Host}' "$base/opencode-a/v1/models") -[ "$h" = "alpha" ] || fail "alpha should be back after two good polls, got '$h'" +# 3. pin opencode-a to alpha: conversation A's next turn moves (an operator pin outranks the lease) +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d '{"host":"alpha","pin":true}' "$base/_crossbar/routes/opencode-a") +[ "$h" = "200" ] || fail "pin returned $h" +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 alpha new" ] || fail "after pin, conversation A should be '200 alpha new', got '$h'" +curl -s "$base/_crossbar/routes" | grep -q '"pinned":"alpha"' || fail "routes view does not show the pin: $(curl -s $base/_crossbar/routes)" -# Streaming: five chunks 200 ms apart must arrive over >= 0.6 s, not all at once at the end. +# 4. queue: alpha has parallel 1, queue_max 1, and answers in 600 ms → of three concurrent, one is 503 +for i in 1 2 3; do (curl -s -o /dev/null -w '%{http_code}\n' -X POST -H 'Content-Type: application/json' -d "$(conv Q$i)" "$base/opencode-a/v1/chat/completions" >> "$tmp/codes") & sleep 0.1; done; wait $! 2>/dev/null || true +sleep 2.5 +sort "$tmp/codes" | uniq -c | tr -s ' ' > "$tmp/counts" +grep -q '2 200' "$tmp/counts" && grep -q '1 503' "$tmp/counts" || fail "queue test wanted two 200 and one 503, got: $(cat "$tmp/counts")" + +# 5. release the pin, drain alpha: new conversations go to beta, A stays on alpha +curl -s -o /dev/null -X POST -H 'Content-Type: application/json' -d '{"release":true}' "$base/_crossbar/routes/opencode-a" +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d '{"drain":true}' "$base/_crossbar/hosts/alpha"); [ "$h" = "200" ] || fail "drain returned $h" +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta new" ] || fail "with alpha draining a new conversation should go to beta, got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{[^}]*"draining":true' || fail "hosts view does not show alpha draining" +curl -s -o /dev/null -X POST -H 'Content-Type: application/json' -d '{"drain":false}' "$base/_crossbar/hosts/alpha" + +# 6. failover + recovery +touch "$tmp/beta.down"; sleep 2.5 +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 alpha new" ] || fail "with beta down conversation C should move to alpha, got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":false' || fail "hosts view does not show beta unhealthy" +rm "$tmp/beta.down"; sleep 3.5 +curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":true' || fail "beta did not recover after two good polls" + +# 7. streaming still arrives incrementally, and the final usage chunk is untouched start=$(date +%s%N) -first="" -curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[]}' \ +curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \ "$base/opencode-a/v1/chat/completions" | while IFS= read -r line; do [ -n "$line" ] || continue now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line" done > "$tmp/stream.txt" firstms=$(head -1 "$tmp/stream.txt" | cut -d' ' -f1); lastms=$(tail -1 "$tmp/stream.txt" | cut -d' ' -f1) -[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst (first ${firstms:-?} ms, last ${lastms:-?} ms): -$(cat "$tmp/stream.txt")" -grep -q 'DONE' "$tmp/stream.txt" || fail "stream did not end with [DONE]" +[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst: $(cat "$tmp/stream.txt")" +grep -q '"usage"' "$tmp/stream.txt" && grep -q 'DONE' "$tmp/stream.txt" || fail "stream lost the usage chunk or DONE" -grep -q 'route=opencode-a host=alpha' "$tmp/crossbar.log" || fail "no request log line" +# 8. accounting and metrics +sleep 1 +u=$(curl -s "$base/_crossbar/usage?by=host") +echo "$u" | grep -q '"key":"alpha"' && echo "$u" | grep -q '"key":"beta"' || fail "usage by host: $u" +echo "$u" | grep -q '"cached_tokens":[1-9]' || fail "usage has no cached tokens (SSE/JSON usage not captured): $u" +curl -s -H 'Accept: text/plain' "$base/_crossbar/usage?by=route" | grep -qi 'cache' || fail "text usage table missing" +m=$(curl -s "$base/_crossbar/metrics") +echo "$m" | grep -q 'crossbar_requests_total{route="opencode-a",host="alpha",status="503"} 1' || fail "metrics missing the 503: $m" +echo "$m" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing host health" +grep -q 'route=opencode-a host=' "$tmp/crossbar.log" || fail "no request log line" echo "smoke: ok (stream spread $((lastms - firstms)) ms)"