Smoke run for v1; README for leases, admin and accounting
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -1,8 +1,10 @@
|
|||||||
# crossbar
|
# crossbar
|
||||||
|
|
||||||
crossbar is an affinity router in front of several `llama-server` routers. A client's identity is
|
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
|
the first path segment of its base URL — its route. Each conversation takes a sticky lease on one
|
||||||
route's list and streams the answer back unbuffered.
|
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
|
## Build
|
||||||
|
|
||||||
@@ -16,22 +18,26 @@ streaming over real HTTP.
|
|||||||
crossbar reads one TOML file. This is `example.toml`:
|
crossbar reads one TOML file. This is `example.toml`:
|
||||||
|
|
||||||
```toml
|
```toml
|
||||||
# crossbar example configuration. Replace <tailnet> and the addresses with your own.
|
# crossbar example configuration (v1). Replace <tailnet> and the addresses with your own.
|
||||||
listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production
|
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
|
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]
|
[hosts.alpha]
|
||||||
base_url = "http://127.0.0.1:18081" # e.g. http://straylight.<tailnet>:11434
|
base_url = "http://127.0.0.1:18081" # e.g. http://straylight.<tailnet>:11434
|
||||||
weight = 1.0
|
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]
|
[hosts.beta]
|
||||||
base_url = "http://127.0.0.1:18082" # e.g. http://titan.<tailnet>:8081
|
base_url = "http://127.0.0.1:18082" # e.g. http://titan.<tailnet>:8081
|
||||||
weight = 2.0
|
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]
|
[routes.opencode-a]
|
||||||
hosts = ["alpha", "beta"]
|
hosts = ["alpha", "beta"]
|
||||||
default_model = "ornith-1.5-35b-a3b"
|
default_model = "ornith-1.5-35b-a3b"
|
||||||
@@ -43,12 +49,15 @@ hosts = ["beta", "alpha"]
|
|||||||
| Key | Meaning |
|
| Key | Meaning |
|
||||||
| --- | --- |
|
| --- | --- |
|
||||||
| `listen` | Where crossbar binds. A tailnet address, never `0.0.0.0`. |
|
| `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. |
|
| `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.<name>.base_url` | The llama-server base URL this host serves. |
|
| `hosts.<name>.base_url` | The llama-server base URL this host serves. |
|
||||||
| `hosts.<name>.weight` | Relative share of new routes this host receives. |
|
| `hosts.<name>.weight` | Relative share of new requests this host receives. |
|
||||||
| `hosts.<name>.models` | The models this host serves, with per-model parallel tuning. |
|
| `hosts.<name>.models` | The models this host serves, with per-model parallel tuning. |
|
||||||
| `routes.<name>.hosts` | Preference order: the first healthy host that serves the model wins. |
|
| `routes.<name>.hosts` | Candidate hosts, tried in order until one is healthy; a conversation leases one of them. |
|
||||||
| `routes.<name>.default_model` | Model used when a request omits one; must be served by a host in the route. |
|
| `routes.<name>.default_model` | Model used when a request omits one; must be served by a host in the route. |
|
||||||
|
|
||||||
## Run
|
## Run
|
||||||
@@ -84,23 +93,73 @@ custom_providers:
|
|||||||
models: { ornith-1.5-35b-a3b: {} }
|
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
|
```sh
|
||||||
|
curl -H 'X-Crossbar-Route: opencode-a' \
|
||||||
`GET /_crossbar/hosts` reports every host's health and loaded models:
|
https://crossbar.<tailnet>:7777/v1/chat/completions
|
||||||
|
|
||||||
```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":""}}
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`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
|
```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
|
```json
|
||||||
`PLAN.md`.
|
{"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.
|
||||||
|
|||||||
@@ -1,11 +1,13 @@
|
|||||||
// fakeupstream stands in for a llama-server router in tests and the smoke run. Do not edit.
|
// 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.
|
// /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
|
// /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
|
// SSE chunks 200 ms apart when the body has "stream": true, then a final chunk carrying
|
||||||
// response carries X-Upstream: <name>.
|
// "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: <name>.
|
||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
@@ -25,11 +27,14 @@ func main() {
|
|||||||
name := flag.String("name", "fake", "name reported in X-Upstream and answers")
|
name := flag.String("name", "fake", "name reported in X-Upstream and answers")
|
||||||
models := flag.String("models", "m", "comma-separated model ids for /v1/models")
|
models := flag.String("models", "m", "comma-separated model ids for /v1/models")
|
||||||
downFile := flag.String("down-file", "", "while this file exists, /health answers 503")
|
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()
|
flag.Parse()
|
||||||
|
|
||||||
ids := strings.Split(*models, ",")
|
ids := strings.Split(*models, ",")
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) }
|
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) {
|
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
|
||||||
stamp(w)
|
stamp(w)
|
||||||
@@ -61,11 +66,12 @@ func main() {
|
|||||||
Stream bool `json:"stream"`
|
Stream bool `json:"stream"`
|
||||||
}
|
}
|
||||||
_ = json.Unmarshal(body, &req)
|
_ = json.Unmarshal(body, &req)
|
||||||
|
time.Sleep(time.Duration(*slow) * time.Millisecond)
|
||||||
if !req.Stream {
|
if !req.Stream {
|
||||||
writeJSON(w, map[string]any{
|
writeJSON(w, map[string]any{
|
||||||
"id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model,
|
"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"}},
|
"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
|
return
|
||||||
}
|
}
|
||||||
@@ -73,16 +79,24 @@ func main() {
|
|||||||
w.Header().Set("Cache-Control", "no-cache")
|
w.Header().Set("Cache-Control", "no-cache")
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
fl, _ := w.(http.Flusher)
|
fl, _ := w.(http.Flusher)
|
||||||
|
flush := func() {
|
||||||
|
if fl != nil {
|
||||||
|
fl.Flush()
|
||||||
|
}
|
||||||
|
}
|
||||||
for i := 1; i <= 5; i++ {
|
for i := 1; i <= 5; i++ {
|
||||||
chunk := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model,
|
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)}}}}
|
"choices": []map[string]any{{"index": 0, "delta": map[string]string{"content": fmt.Sprintf("%s chunk %d ", *name, i)}}}}
|
||||||
b, _ := json.Marshal(chunk)
|
b, _ := json.Marshal(chunk)
|
||||||
fmt.Fprintf(w, "data: %s\n\n", b)
|
fmt.Fprintf(w, "data: %s\n\n", b)
|
||||||
if fl != nil {
|
flush()
|
||||||
fl.Flush()
|
|
||||||
}
|
|
||||||
time.Sleep(200 * time.Millisecond)
|
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")
|
fmt.Fprint(w, "data: [DONE]\n\n")
|
||||||
})
|
})
|
||||||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||||||
@@ -90,7 +104,7 @@ func main() {
|
|||||||
http.Error(w, `{"error":"not found"}`, http.StatusNotFound)
|
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}
|
srv := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second}
|
||||||
log.Fatal(srv.ListenAndServe())
|
log.Fatal(srv.ListenAndServe())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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/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. | ? |
|
||||||
| 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/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/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. | ? |
|
| 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. | ? |
|
||||||
|
|||||||
+15
-1
@@ -75,7 +75,7 @@ func (c *hostChooser) Choose(candidates []string, model string) (string, bool) {
|
|||||||
Draining: false,
|
Draining: false,
|
||||||
Loaded: contains(s.Loaded, model),
|
Loaded: contains(s.Loaded, model),
|
||||||
CanServe: c.cfg.Serves(host, 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),
|
Queued: c.lim.Queued(host, model),
|
||||||
Weight: c.cfg.Hosts[host].Weight,
|
Weight: c.cfg.Hosts[host].Weight,
|
||||||
}
|
}
|
||||||
@@ -91,3 +91,17 @@ func contains(list []string, v string) bool {
|
|||||||
}
|
}
|
||||||
return false
|
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
|
||||||
|
}
|
||||||
|
|||||||
+59
-24
@@ -1,45 +1,80 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# Smoke run: two fake upstreams, one crossbar, real HTTP. Prints "smoke: ok" or fails.
|
# Smoke run (v1): two fake upstreams, one crossbar with a fresh SQLite file, real HTTP.
|
||||||
# Needs: bin/crossbar and bin/fakeupstream (make build), curl.
|
# 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
|
set -eu
|
||||||
cd "$(dirname "$0")/.."
|
cd "$(dirname "$0")/.."
|
||||||
tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM
|
tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM
|
||||||
pids=""
|
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/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
|
sleep 1.5
|
||||||
fail() { echo "smoke: FAIL: $*" >&2; echo "--- crossbar.log"; cat "$tmp/crossbar.log"; exit 1; }
|
fail() { echo "smoke: FAIL: $*" >&2; echo "--- crossbar.log"; cat "$tmp/crossbar.log"; exit 1; }
|
||||||
base=http://127.0.0.1:17777
|
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")
|
# 1. a conversation gets a lease and keeps it; beta wins (2 slots × weight 2 vs 1 × 1)
|
||||||
[ "$h" = "200 alpha" ] || fail "opencode-a should go to alpha, got '$h'"
|
h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions")
|
||||||
h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/hermes-x/v1/models")
|
[ "$h" = "200 beta new" ] || fail "first turn should be '200 beta new', got '$h'"
|
||||||
[ "$h" = "200 beta" ] || fail "hermes-x should go to beta, got '$h'"
|
h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions")
|
||||||
h=$(curl -s -o /dev/null -w '%{http_code}' "$base/nope/v1/models")
|
[ "$h" = "200 beta reused" ] || fail "second turn should reuse beta, got '$h'"
|
||||||
[ "$h" = "404" ] || fail "unknown route should be 404, got '$h'"
|
|
||||||
|
|
||||||
touch "$tmp/alpha.down"; sleep 2.5 # poll_interval is 1s in example.toml
|
# 2. header route
|
||||||
h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/opencode-a/v1/models")
|
h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Route: hermes-x' -d "$(conv B)" "$base/v1/chat/completions")
|
||||||
[ "$h" = "200 beta" ] || fail "with alpha down, opencode-a should fail over to beta, got '$h'"
|
case "$h" in "200 beta new") ;; *) fail "header route hermes-x should be '200 beta new', got '$h'";; esac
|
||||||
curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{"healthy":false' || fail "/_crossbar/hosts does not show alpha unhealthy: $(curl -s $base/_crossbar/hosts)"
|
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
|
# 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 '%header{X-Crossbar-Host}' "$base/opencode-a/v1/models")
|
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" = "alpha" ] || fail "alpha should be back after two good polls, got '$h'"
|
[ "$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)
|
start=$(date +%s%N)
|
||||||
first=""
|
curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \
|
||||||
curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[]}' \
|
|
||||||
"$base/opencode-a/v1/chat/completions" | while IFS= read -r line; do
|
"$base/opencode-a/v1/chat/completions" | while IFS= read -r line; do
|
||||||
[ -n "$line" ] || continue
|
[ -n "$line" ] || continue
|
||||||
now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line"
|
now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line"
|
||||||
done > "$tmp/stream.txt"
|
done > "$tmp/stream.txt"
|
||||||
firstms=$(head -1 "$tmp/stream.txt" | cut -d' ' -f1); lastms=$(tail -1 "$tmp/stream.txt" | cut -d' ' -f1)
|
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):
|
[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst: $(cat "$tmp/stream.txt")"
|
||||||
$(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 'DONE' "$tmp/stream.txt" || fail "stream did not end with [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)"
|
echo "smoke: ok (stream spread $((lastms - firstms)) ms)"
|
||||||
|
|||||||
Reference in New Issue
Block a user