From 3518e84dd736ce61fede7445a76f6e1f0e4d59cf Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 17:36:59 -0700 Subject: [PATCH] v2.3 plan: clients that manage their own slots (control calls, route affinity, queue = false, route listeners, llama-server ctx error); given tests Co-Authored-By: Claude Opus 5.5 --- docs/plans/v2.3/01-control-plane.md | 77 +++++ docs/plans/v2.3/02-affinity-queue.md | 96 ++++++ docs/plans/v2.3/03-route-listeners.md | 91 ++++++ docs/plans/v2.3/04-ctx-error-docs.md | 56 ++++ docs/plans/v2.3/README.md | 57 ++++ docs/plans/v2.3/_files/example.toml | 43 +++ .../_files/internal/config/config_v23_test.go | 75 +++++ .../_files/internal/config/listen_test.go | 54 ++++ .../identity/route_middleware_test.go | 47 +++ .../_files/internal/limiter/track_test.go | 91 ++++++ .../_files/internal/proxy/affinity_test.go | 153 +++++++++ .../_files/internal/proxy/control_test.go | 176 +++++++++++ .../_files/internal/proxy/ctxguard_test.go | 164 ++++++++++ .../_files/internal/proxy/listener_test.go | 112 +++++++ .../v2.3/_files/internal/proxy/proxy_test.go | 292 ++++++++++++++++++ docs/plans/v2.3/_files/tools/smoke.sh | 71 +++++ 16 files changed, 1655 insertions(+) create mode 100644 docs/plans/v2.3/01-control-plane.md create mode 100644 docs/plans/v2.3/02-affinity-queue.md create mode 100644 docs/plans/v2.3/03-route-listeners.md create mode 100644 docs/plans/v2.3/04-ctx-error-docs.md create mode 100644 docs/plans/v2.3/README.md create mode 100644 docs/plans/v2.3/_files/example.toml create mode 100644 docs/plans/v2.3/_files/internal/config/config_v23_test.go create mode 100644 docs/plans/v2.3/_files/internal/config/listen_test.go create mode 100644 docs/plans/v2.3/_files/internal/identity/route_middleware_test.go create mode 100644 docs/plans/v2.3/_files/internal/limiter/track_test.go create mode 100644 docs/plans/v2.3/_files/internal/proxy/affinity_test.go create mode 100644 docs/plans/v2.3/_files/internal/proxy/control_test.go create mode 100644 docs/plans/v2.3/_files/internal/proxy/ctxguard_test.go create mode 100644 docs/plans/v2.3/_files/internal/proxy/listener_test.go create mode 100644 docs/plans/v2.3/_files/internal/proxy/proxy_test.go create mode 100644 docs/plans/v2.3/_files/tools/smoke.sh diff --git a/docs/plans/v2.3/01-control-plane.md b/docs/plans/v2.3/01-control-plane.md new file mode 100644 index 0000000..76c7670 --- /dev/null +++ b/docs/plans/v2.3/01-control-plane.md @@ -0,0 +1,77 @@ +# v2.3 task 01: control-plane requests + +**Branch:** `v2.3` (`git switch -c v2.3 master` if it does not exist, else `git switch v2.3`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Control-plane requests follow the lease but take no slot and write no row` + +## Goal + +A client that manages its own llama-server slot makes small calls beside its chat stream: it +polls `GET /slots?model=X` while it waits, reads `GET /props?model=X`, tokenizes, and sends +`POST /v1/chat/completions/control` on a second connection **while its own stream holds a slot**. +Today `/slots` and `/tokenize` are 404, a GET is leased under the route's default model instead +of `?model=`, and every call takes a limiter slot — so `/control` can queue behind its own +stream, or get 503 when the queue is full. After this task those calls follow the lease like any +request but never wait for or take a slot, skip the context guard, and write no accounting row. + +## Files + +- Copy: `internal/proxy/control_test.go`, and the **replacement** `internal/proxy/proxy_test.go` + (overwrites the v1 copy: the `/r/slots → 404` row becomes `/r/slots/0` and `/r/metrics`; the + v2.3 copy is now the protected one) +- Create: `internal/proxy/control.go` — the control-call test and the model-from-query rule live + here (`proxy.go` is at 325 lines) +- Modify: `internal/proxy/proxy.go`, `internal/proxy/forward.go`, `docs/implementer-log.md` + +## Rules + +1. **Model.** The body's top-level `"model"` wins; else the query parameter `model` + (`r.URL.Query().Get("model")`); else the route's `default_model`. This is used for the lease + key and the limiter pair, exactly where the body model is used today. +2. **Paths.** `allowedPath` also admits `rest == "/slots"` and `rest == "/tokenize"` (exact + match on the path; the query string is not part of `rest`). `/slots/0`, `/slots/0?action=…` + and anything else stay `404 {"error":"not found"}`. +3. **Control calls** are: method `GET` or `HEAD` (any allowed path), or method `POST` with + `rest` exactly `/tokenize` or `/v1/chat/completions/control`. Everything else — in + particular `POST /v1/chat/completions` — is not a control call. +4. A control call is routed and leased exactly as today (same `lease.Acquire`, same wake path + when no host is healthy), then forwarded **without** `lim.Acquire`, **without** the context + guard, and **without** an accounting row (`writeRecord` is not called for it). Its log line + is `Debug`, not `Info` (a client polls `/slots` every 5 s). It still gets the + `X-Crossbar-Host` / `X-Crossbar-Lease` headers and still marks a host down on a transport + error, like any forward. +5. Do not duplicate `forward`. Pass what it needs to know (for example a `control bool`, or a + small options struct if the parameter list gets long) and skip the row and the `Info` log + inside it. + +## Facts you need + +- `peekModel` already reads and restores the body; `GET`/`HEAD` return `""` there. The query + fallback goes after it, in one place. +- The rig in `helpers_test.go` passes the store as the recorder, and the row is written after + the answer is sent — the tests wait for rows; do not add sleeps to production code. +- `waitUntil` is defined in `proxy_test.go`; `conversation(id, turn)` builds a chat body. + +## Steps + +- [ ] **1.** Branch as above; copy the two given files. +- [ ] **2. See them fail:** `go test -count=1 ./internal/proxy/ -run 'TestGetModel|TestControl|TestChatIsNotControl'` + → 404 on `/slots` and `/tokenize`, 503 `queue full` on control calls, 4 stray rows. +- [ ] **3.** `control.go`: the control-call test and the model rule. **4.** `proxy.go`: use them; + branch in `serveLeased` (no limiter, no guard for a control call). **5.** `forward.go`: no row, + `Debug` log for a control call. +- [ ] **6.** `gofmt -w`; `go test -race -count=3 ./internal/proxy/` → `ok`. +- [ ] **7.** `make gate`; `make smoke`. **8.** Row `v2.3/01-control-plane`; commit. + +```sh +git add internal/proxy docs/implementer-log.md +git commit +``` + +## Done when + +- The new tests and every earlier proxy test pass under `-race -count=3`; gate and smoke ok; + given files byte-identical; no file over 400 lines. + +## Stop and report if + +- Passing needs a change to any given test, or `forward` cannot skip the row without copying it. diff --git a/docs/plans/v2.3/02-affinity-queue.md b/docs/plans/v2.3/02-affinity-queue.md new file mode 100644 index 0000000..15d7590 --- /dev/null +++ b/docs/plans/v2.3/02-affinity-queue.md @@ -0,0 +1,96 @@ +# v2.3 task 02: route affinity and queue = false + +**Branch:** `v2.3` (`git switch v2.3`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Routes may share one lease (affinity = "route") and skip crossbar's queue (queue = false)` + +## Goal + +Two route keys for a client that manages its own slot: + +- `affinity = "route"` — one lease for the whole route. Today each conversation (fingerprint) + gets its own lease, and calls without a fingerprint lease "the route itself"; a client whose + `/control` and `/slots` calls must reach the host its chat is on needs them all on one lease. +- `queue = false` — crossbar never holds or refuses the route's requests. The client pins its + llama-server slot (`id_slot`), so llama-server queues it and its `/slots` shows the slot busy; + a request held in crossbar's queue instead looks idle to the client, which gives up after 30 s. + The requests still count as load on the host, so routes that do queue see the host full. + +## Files + +- Copy: `internal/config/config_v23_test.go`, `internal/limiter/track_test.go`, + `internal/proxy/affinity_test.go` +- Modify: `internal/config/route.go` (**move the `Route` struct here** from `config.go`, which + is at 398 lines, and add the fields and methods here; `config.go` keeps a one-line call into + the route validation), `internal/config/config.go` (minimal), `internal/limiter/limiter.go`, + `internal/proxy/proxy.go` (or `control.go` if `proxy.go` would pass 400 lines), + `docs/implementer-log.md` + +## Interfaces + +```go +package config + +type Route struct { + Hosts []string `toml:"hosts"` + DefaultModel string `toml:"default_model"` + Peers []string `toml:"peers"` + Affinity string `toml:"affinity"` // "" or "conversation" (the default), or "route" + Queue *bool `toml:"queue"` // nil means true +} + +// PerRoute reports affinity = "route": every request on the route shares one lease. +func (r Route) PerRoute() bool + +// Queues reports whether the route's requests wait in (and can be refused by) crossbar's +// per-(host, model) queue; false only for queue = false. +func (r Route) Queues() bool + +package limiter + +// Track counts one request against (host, model) without waiting and without refusing: in flight +// may exceed parallel. The returned release is idempotent. +func (l *Limiter) Track(host, model string) (release func()) +``` + +## Rules + +1. **Validation.** `affinity` other than `""`, `"conversation"` or `"route"` is an error whose + text contains `routes..affinity` (for example + `routes.convo.affinity: must be "conversation" or "route"`). Both keys are allowed on + templates; `cfg.Route(name)` returns them for every concrete route the template serves. +2. **Lease key.** For a `PerRoute()` route the lease key's fingerprint is `""` for **every** + request (chat or control), so every request on the route uses one lease per model. The + accounting row keeps the request's real fingerprint (it is still useful in usage views). +3. **Queue.** For a route where `Queues()` is false, a non-control request takes + `lim.Track(host, model)` instead of `lim.Acquire` (and releases it when done, like the slot). + Control calls (task 01) take neither. +4. **Release rule.** A release — from `Acquire`'s slot or from `Track` — hands the slot to the + first waiter **only when in flight ≤ parallel** at that moment; otherwise it just decrements + in flight. (With only `Acquire` in use in-flight never exceeds parallel, so today's behaviour + is unchanged.) `InFlight`, `FreeSlots` and `Queued` count tracked requests like any other. +5. The context guard runs as today on both kinds of route. + +## Steps + +- [ ] **1.** `git switch v2.3`; copy the three given tests. +- [ ] **2. See them fail** (compile: `PerRoute`, `Queues`, `Track` missing). +- [ ] **3.** `route.go` (struct move, fields, methods, validation). **4.** `limiter.go` + (`Track`, the release rule). **5.** `proxy.go` (lease key, `Track`). +- [ ] **6.** `gofmt -w`; `go test -race -count=3 ./internal/limiter/ ./internal/config/ ./internal/proxy/` → `ok`. +- [ ] **7.** `make gate`; `make smoke`. **8.** Row `v2.3/02-affinity-queue`; commit. + +```sh +git add internal/config internal/limiter internal/proxy docs/implementer-log.md +git commit +``` + +## Done when + +- The new tests and every earlier test pass under `-race -count=3`; every earlier limiter test + is still green (the release rule must not change `Acquire`-only behaviour); gate and smoke ok; + given files byte-identical; no file over 400 lines. + +## Stop and report if + +- The struct move breaks a given test, or the release rule cannot be met without changing + `Acquire`'s results in an earlier test. diff --git a/docs/plans/v2.3/03-route-listeners.md b/docs/plans/v2.3/03-route-listeners.md new file mode 100644 index 0000000..408be01 --- /dev/null +++ b/docs/plans/v2.3/03-route-listeners.md @@ -0,0 +1,91 @@ +# v2.3 task 03: a route's dedicated listener + +**Branch:** `v2.3` (`git switch v2.3`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `A route may have its own listener: every request there is that route, paths unprefixed` + +## Goal + +Boxmaker's `inferproxy` connects to one host:port and rewrites nothing: its paths are +`/v1/chat/completions`, `/slots?model=…`, and it sends no extra header. Crossbar reads the route +from the first path segment or `X-Crossbar-Route`, so it cannot route those requests. After this +task a concrete route may set `listen = "host:port"`; crossbar serves that address too, and every +request arriving there is that route, with the whole path passed upstream as it is. + +## Files + +- Copy: `internal/config/listen_test.go`, `internal/proxy/listener_test.go`, + `internal/identity/route_middleware_test.go`, and the **replacements** `example.toml` (was + v2.2's) and `tools/smoke.sh` (was v2's; adds check 6) — the v2.3 copies are now protected +- Modify: `internal/config/route.go` (the `Listen` field and its validation), + `internal/proxy/proxy.go` (or a new `internal/proxy/listener.go`), + `internal/identity/middleware.go`, `cmd/crossbar/main.go`, `docs/implementer-log.md` + +## Interfaces + +```go +package config +// in Route: + Listen string `toml:"listen"` // "" = none; else host:port of the route's own listener + +package proxy +// ForRoute serves route alone: the request path is the upstream path (no route segment is +// taken from it), and everything after routing is exactly what ServeHTTP does. +func (p *Handler) ForRoute(route string) http.Handler + +package identity +// RouteMiddleware gates every request on peers (the fixed route's allow list), whatever path or +// X-Crossbar-Route header it carries. Empty peers lets everyone through, as for Middleware. +func RouteMiddleware(c *Checker, peers []string, next http.Handler) http.Handler +``` + +## Rules + +1. **Validation** (errors name the key): + - `listen` must split with `net.SplitHostPort` and its port must be a number 1–65535 + (`strconv.Atoi`) → else an error containing `routes..listen`. + - Not on a template: `routes..listen: a template route cannot have its own listener` + (the text contains the template's name, e.g. `t-*`). + - Not the top-level `listen`: an error containing `routes..listen`. + - Unique across routes: the second route (in sorted name order) gets an error containing + `routes..listen` and the other route's name. +2. **`ForRoute(route)`** for each request: + - `cfg.Route(route)` not ok → `404 {"error":"unknown route"}`. + - `X-Crossbar-Route` set and different from `route` → `400 {"error":"conflicting route"}`; + set and equal → ignored. + - `rest` is `r.URL.Path` unchanged; `allowedPath(rest)` false → `404 {"error":"not found"}`. + So `/bm/v1/models` (a prefixed path) and `/_crossbar/hosts` are 404 on the listener. + - Then the same flow as `ServeHTTP` from the model peek onwards — factor that flow into one + function both call; do not copy it. +3. **`RouteMiddleware`**: like `Middleware`, including the `X-Crossbar-Peer` context for + header mode, but with the fixed peers and **no** admin-path exemption (there is no admin on + a route listener). +4. **`main.go`**: for each route with `Listen` set, in sorted route order, one more + `http.Server{Addr: rt.Listen, Handler: h, ReadHeaderTimeout: 10 * time.Second}` where `h` is + `p.ForRoute(name)`, wrapped in `identity.RouteMiddleware(checker, rt.Peers, …)` when identity + is not `off`. No admin mux on it. Log `listening` with `addr` and `route`. All servers shut + down together on ctx done; any server's error other than `http.ErrServerClosed` ends `run` + with that error (and shuts the others down). + +## Steps + +- [ ] **1.** `git switch v2.3`; copy the five given files. +- [ ] **2. See them fail** (compile: `Listen`, `ForRoute`, `RouteMiddleware` missing). +- [ ] **3.** `route.go`. **4.** proxy (`ForRoute`, the shared flow). **5.** `middleware.go`. + **6.** `main.go`. +- [ ] **7.** `gofmt -w`; `go test -race -count=3 ./internal/config/ ./internal/proxy/ ./internal/identity/` → `ok`. +- [ ] **8.** `make gate`; `make smoke` (check 6 is the dedicated listener on 127.0.0.1:17801). +- [ ] **9.** Row `v2.3/03-route-listeners`; commit. + +```sh +git add internal/config internal/proxy internal/identity cmd/crossbar example.toml tools/smoke.sh docs/implementer-log.md +git commit +``` + +## Done when + +- All given tests pass under `-race -count=3`; gate and smoke ok; given files byte-identical; + no file over 400 lines. + +## Stop and report if + +- Smoke check 6 fails for a reason in the fake upstream or the script rather than in crossbar. diff --git a/docs/plans/v2.3/04-ctx-error-docs.md b/docs/plans/v2.3/04-ctx-error-docs.md new file mode 100644 index 0000000..1c938d3 --- /dev/null +++ b/docs/plans/v2.3/04-ctx-error-docs.md @@ -0,0 +1,56 @@ +# v2.3 task 04: llama-server's error shape for a context refusal; README + +**Branch:** `v2.3` (`git switch v2.3`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Context refusal in llama-server's exceed_context_size_error shape; README for v2.3` + +## Goal + +When no host can fit a prompt, crossbar answers `400 {"error":"prompt too large","estimate":N,"max":M}`. +A client that already handles llama-server's own overflow error (Boxmaker keys on `error.type` +and reads only the first 4 KiB) does not recognise it. After this task the body is the server's +shape, so the client handles crossbar's refusal like the server's: + +```json +{"error":{"code":400,"type":"exceed_context_size_error","message":"prompt too large","n_prompt_tokens":N,"n_ctx":M}} +``` + +`N` is the estimate and `M` the largest per-slot context on the route, as before. + +## Files + +- Copy: the **replacement** `internal/proxy/ctxguard_test.go` (was v2's; only the 400-body + assertions changed; the v2.3 copy is now protected) +- Modify: `internal/proxy/ctxguard.go` (`refuseCtx`), `README.md`, `docs/implementer-log.md` + +## Rules + +1. The body is exactly one JSON object whose only top-level key is `"error"`, so it starts with + `{"error":`; `Content-Type: application/json`; status 400. The accounting row is unchanged + (status 400, `Err` "prompt too large"). Every other crossbar error keeps its current + `{"error":""}` shape. +2. `README.md`: + - The context-guard section: the new body. + - A new section **"Clients that manage their own slots"** covering: control calls (which + requests, and that they follow the lease but take no slot, skip the guard and write no + row); `/slots` and `/tokenize` are proxied, `/slots/` actions are not; a GET's model + comes from `?model=`; the route keys `affinity`, `queue` and `listen` with the + `boxmaker-a` example from `example.toml`; that `listen` is refused on templates and must + not be the main address; that the admin API is not served on a route listener. + - The "hosts view"/config reference tables, if they list route keys, gain the three keys. + +## Steps + +- [ ] **1.** `git switch v2.3`; copy the replacement test. **2. See it fail** (old body). +- [ ] **3.** `refuseCtx`. **4.** README. +- [ ] **5.** `gofmt -w`; `make gate`; `make smoke` (check 3 still finds `"prompt too large"`). +- [ ] **6.** Row `v2.3/04-ctx-error-docs`; commit. + +```sh +git add internal/proxy README.md docs/implementer-log.md +git commit +``` + +## Done when + +- All tests pass; gate and smoke ok; given files byte-identical; README describes what v2.3 + does and nothing it does not. diff --git a/docs/plans/v2.3/README.md b/docs/plans/v2.3/README.md new file mode 100644 index 0000000..6c911f0 --- /dev/null +++ b/docs/plans/v2.3/README.md @@ -0,0 +1,57 @@ +# v2.3 implementation plan: clients that manage their own slots + +> **For the implementing model:** do not work from this file. The owner gives you one task file at +> a time. This file is the index for the owner and the reviewer. + +**Goal:** serve Boxmaker, a harness whose `inferproxy` talks plain HTTP/1.1 to one host:port and +rewrites nothing. It pins `id_slot`, polls `GET /slots?model=` while it waits, reads +`GET /props?model=` once, and sends `POST /v1/chat/completions/control` on a second connection +while its own stream is running. Checked on 2026-09-25 against crossbar at 4c64158, it failed on +six counts (thread `i7jeubrtziru38s5gn8gmha44a`): no route in its paths; `/slots` and `/tokenize` +not proxied; side calls leased separately from the stream; `/control` taking a limiter slot +behind its own stream; crossbar's queue hiding a waiting request from the server's `/slots`; and a +context refusal that is not llama-server's `exceed_context_size_error`. + +- **01-control-plane** — every route: a GET's model comes from `?model=`; `/slots` and + `/tokenize` are proxied; control calls (any GET/HEAD, `POST /tokenize`, + `POST /v1/chat/completions/control`) follow the lease but skip the limiter, the context guard + and the accounting row. Given: `proxy/control_test.go`; replaces `proxy/proxy_test.go` (v1: + the `/r/slots → 404` row becomes `/r/slots/0` and `/r/metrics`). +- **02-affinity-queue** — route keys `affinity = "route"` (one lease for the route) and + `queue = false` (count the request as load, never hold or refuse it); `limiter.Track`. + Given: `config/config_v23_test.go`, `limiter/track_test.go`, `proxy/affinity_test.go`. +- **03-route-listeners** — route key `listen`: a dedicated listener where every request is that + route with an unprefixed path; `Handler.ForRoute`, `identity.RouteMiddleware`, one server per + listener in `main`. Given: `config/listen_test.go`, `proxy/listener_test.go`, + `identity/route_middleware_test.go`; replaces `example.toml` (v2.2: adds `boxmaker-a`) and + `tools/smoke.sh` (v2: adds check 6, the dedicated listener). +- **04-ctx-error-docs** — the context refusal in llama-server's shape + `{"error":{"code":400,"type":"exceed_context_size_error","message":"prompt too large","n_prompt_tokens":N,"n_ctx":M}}`; + README. Given: replaces `proxy/ctxguard_test.go` (v2). + +**Order matters:** 02's affinity test uses `/slots` (01); 03's listener test uses route affinity +(02). Each task is green on its own given tests plus all earlier ones. + +**How this plan was made:** acceptance tests first, no reference implementation; the given tests +compiled against a panic-only skeleton of the new names (`Route.PerRoute`, `Route.Queues`, +`Route.Listen`, `Limiter.Track`, `Handler.ForRoute`, `identity.RouteMiddleware`) on master +4c64158 and failed there for the intended reasons (404 on `/slots`/`/tokenize`, 503 queue full +on control calls, 4 stray accounting rows, the old error body, requests held behind one slot, +unvalidated `listen`/`affinity`). + +**Facts about the live hosts (2026-09-25):** all three routers run llama-server b10964; `/slots` +answers 200 on all three; `POST /v1/chat/completions/control` exists (`{"success":false,"message":"no +active completion for this id"}` for an unknown id). In router mode `GET /slots?model=X` and +`/props?model=X` **autoload X** — a control call only ever reaches the leased host, which is where +the client's chat goes anyway, so this is the load the client asked for. + +## Global constraints + +- Everything in `AGENTS.md`. Branch `v2.3` from `master`. One task, one fresh OpenCode session, + one commit. Given files are copied and never edited; earlier plans' given files stay protected, + except the four this plan replaces (`proxy/proxy_test.go`, `proxy/ctxguard_test.go`, + `example.toml`, `tools/smoke.sh`), whose v2.3 copies are then the protected ones. + +## Changes during the run + +(none yet) diff --git a/docs/plans/v2.3/_files/example.toml b/docs/plans/v2.3/_files/example.toml new file mode 100644 index 0000000..ccbb4ee --- /dev/null +++ b/docs/plans/v2.3/_files/example.toml @@ -0,0 +1,43 @@ +# 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 +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 +identity = "off" # "tailscale" gates routes with `peers` by `tailscale whois`; "header" trusts X-Crossbar-Peer (TEST ONLY) + +[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 = 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 = 2 } } +[hosts.beta.wake] # v2: wake a sleeping host when nothing else can take a new lease +mac = "aa:bb:cc:dd:ee:02" +broadcast = "127.0.0.1:19082" # the LAN broadcast address, port 9, in production +# broadcasts = ["127.0.0.1:19082", "127.0.0.1:19083"] # v2.2: a roaming host on several networks; this or broadcast, not both, and at least one +wait = "20s" + +# 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" + +[routes.hermes-x] +hosts = ["beta", "alpha"] +# peers = ["talos"] # v2: with identity = "tailscale", only these tailnet nodes may use the route + +# v2.3: a client that manages its own llama-server slot (it pins id_slot, polls /slots, steers a +# running completion through /v1/chat/completions/control) and cannot put a route in the path. +# The route gets its own port; every request there is this route and the path goes upstream as is. +[routes.boxmaker-a] +hosts = ["beta", "alpha"] +default_model = "ornith-1.5-35b-a3b" +listen = "127.0.0.1:17801" # a tailnet address in production; never the main listen address +affinity = "route" # one lease for the whole route, not one per conversation +queue = false # counted as load but never held or refused: the server's own slot queue does that diff --git a/docs/plans/v2.3/_files/internal/config/config_v23_test.go b/docs/plans/v2.3/_files/internal/config/config_v23_test.go new file mode 100644 index 0000000..c336107 --- /dev/null +++ b/docs/plans/v2.3/_files/internal/config/config_v23_test.go @@ -0,0 +1,75 @@ +package config_test + +// v2.3 task 02: the affinity and queue route keys. + +import ( + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const affinityBase = ` +listen = "127.0.0.1:1" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.plain] +hosts = ["a"] +[routes.convo] +hosts = ["a"] +affinity = "conversation" +[routes.boxmaker] +hosts = ["a"] +affinity = "route" +queue = false +[routes."bm-*"] +hosts = ["a"] +affinity = "route" +queue = false +[routes.queued] +hosts = ["a"] +queue = true +` + +func TestAffinityAndQueueKeys(t *testing.T) { + c, err := config.Parse(strings.NewReader(affinityBase)) + if err != nil { + t.Fatal(err) + } + for _, tc := range []struct { + route string + perRoute, queues bool + }{ + {"plain", false, true}, // defaults: conversation affinity, queueing on + {"convo", false, true}, + {"boxmaker", true, false}, + {"bm-agent-1", true, false}, // a template's keys reach its concrete routes + {"queued", false, true}, + } { + r, _, ok := c.Route(tc.route) + if !ok { + t.Fatalf("route %q not found", tc.route) + } + if r.PerRoute() != tc.perRoute || r.Queues() != tc.queues { + t.Errorf("%s: PerRoute %v Queues %v, want %v %v", tc.route, r.PerRoute(), r.Queues(), tc.perRoute, tc.queues) + } + } +} + +func TestAffinityRejectsUnknownValues(t *testing.T) { + for _, bad := range []string{`"session"`, `"Route"`, `1`} { + text := strings.Replace(affinityBase, `affinity = "conversation"`, "affinity = "+bad, 1) + _, err := config.Parse(strings.NewReader(text)) + if err == nil || !strings.Contains(err.Error(), "routes.convo.affinity") { + t.Errorf("affinity = %s: err %v, want one naming routes.convo.affinity", bad, err) + } + } +} + +func TestQueueMustBeABool(t *testing.T) { + text := strings.Replace(affinityBase, "queue = true", `queue = "no"`, 1) + if _, err := config.Parse(strings.NewReader(text)); err == nil { + t.Error(`queue = "no" parsed; want an error`) + } +} diff --git a/docs/plans/v2.3/_files/internal/config/listen_test.go b/docs/plans/v2.3/_files/internal/config/listen_test.go new file mode 100644 index 0000000..5e1979a --- /dev/null +++ b/docs/plans/v2.3/_files/internal/config/listen_test.go @@ -0,0 +1,54 @@ +package config_test + +// v2.3 task 03: a concrete route may own a dedicated listener. Every request that arrives on it is +// that route, with the upstream path unprefixed, for clients that cannot put a route in the path +// or a header (Boxmaker's inferproxy rewrites nothing). + +import ( + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const listenBase = ` +listen = "127.0.0.1:7777" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.bm-a] +hosts = ["a"] +listen = "127.0.0.1:7801" +[routes.bm-b] +hosts = ["a"] +listen = "127.0.0.1:7802" +[routes.plain] +hosts = ["a"] +` + +func TestRouteListen(t *testing.T) { + c, err := config.Parse(strings.NewReader(listenBase)) + if err != nil { + t.Fatal(err) + } + for route, want := range map[string]string{"bm-a": "127.0.0.1:7801", "bm-b": "127.0.0.1:7802", "plain": ""} { + if got := c.Routes[route].Listen; got != want { + t.Errorf("%s listen = %q, want %q", route, got, want) + } + } +} + +func TestRouteListenRejected(t *testing.T) { + for _, tc := range []struct{ name, text, want string }{ + {"not host:port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"7801"`, 1), "routes.bm-a.listen"}, + {"bad port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:http"`, 1), "routes.bm-a.listen"}, + {"port zero", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:0"`, 1), "routes.bm-a.listen"}, + {"same as another route", strings.Replace(listenBase, `"127.0.0.1:7802"`, `"127.0.0.1:7801"`, 1), "listen"}, + {"same as the main listener", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:7777"`, 1), "routes.bm-a.listen"}, + {"on a template", listenBase + "[routes.\"t-*\"]\nhosts = [\"a\"]\nlisten = \"127.0.0.1:7803\"\n", "t-*"}, + } { + if _, err := config.Parse(strings.NewReader(tc.text)); err == nil || !strings.Contains(err.Error(), tc.want) { + t.Errorf("%s: err %v, want one containing %q", tc.name, err, tc.want) + } + } +} diff --git a/docs/plans/v2.3/_files/internal/identity/route_middleware_test.go b/docs/plans/v2.3/_files/internal/identity/route_middleware_test.go new file mode 100644 index 0000000..c860885 --- /dev/null +++ b/docs/plans/v2.3/_files/internal/identity/route_middleware_test.go @@ -0,0 +1,47 @@ +package identity_test + +// v2.3 task 03: on a route's dedicated listener the route is fixed, so the gate is that route's +// peers for every request, whatever path or X-Crossbar-Route header the caller sends. + +import ( + "net/http" + "net/http/httptest" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/identity" +) + +func TestRouteMiddleware(t *testing.T) { + inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) }) + checker := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"}) + locked := identity.RouteMiddleware(checker, []string{"talos"}, inner) + open := identity.RouteMiddleware(checker, nil, inner) + for _, tc := range []struct { + name string + h http.Handler + path, hdr string + addr string + want int + }{ + {"right peer", locked, "/v1/chat/completions", "", "100.64.0.5:5", 204}, + {"wrong peer", locked, "/v1/chat/completions", "", "100.64.0.9:5", 403}, + {"not a peer", locked, "/slots", "", "203.0.113.1:5", 403}, + {"a path that looks like an open route is still this route", locked, "/open/v1/models", "", "100.64.0.9:5", 403}, + {"a header naming another route does not change the gate", locked, "/v1/models", "open", "100.64.0.9:5", 403}, + {"admin-looking path is gated too (no admin on this listener)", locked, "/_crossbar/hosts", "", "100.64.0.9:5", 403}, + {"open route, anyone", open, "/v1/models", "", "203.0.113.1:5", 204}, + } { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, tc.path, nil) + req.RemoteAddr = tc.addr + if tc.hdr != "" { + req.Header.Set("X-Crossbar-Route", tc.hdr) + } + rec := httptest.NewRecorder() + tc.h.ServeHTTP(rec, req) + if rec.Code != tc.want { + t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String()) + } + }) + } +} diff --git a/docs/plans/v2.3/_files/internal/limiter/track_test.go b/docs/plans/v2.3/_files/internal/limiter/track_test.go new file mode 100644 index 0000000..9e8464f --- /dev/null +++ b/docs/plans/v2.3/_files/internal/limiter/track_test.go @@ -0,0 +1,91 @@ +package limiter_test + +// v2.3 task 02: Track counts a request without holding or refusing it. A route with queue = false +// leaves queueing to llama-server's own slots, but its requests are still load on the host, so the +// routes that do queue must see them. + +import ( + "context" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/limiter" +) + +func TestTrackNeverWaitsAndCounts(t *testing.T) { + l := limiter.New() + l.Configure("alpha", "m", 1, 0) // one slot, no waiting room + + start := time.Now() + rel1 := l.Track("alpha", "m") + rel2 := l.Track("alpha", "m") + rel3 := l.Track("alpha", "m") + if d := time.Since(start); d > 50*time.Millisecond { + t.Fatalf("Track waited %v", d) + } + if n := l.InFlight("alpha", "m"); n != 3 { + t.Fatalf("in flight = %d, want 3 (Track may pass parallel)", n) + } + if n := l.FreeSlots("alpha"); n != 0 { + t.Errorf("free slots = %d, want 0", n) + } + // A queueing request sees the host full: no waiting room, so it is refused. + if _, _, err := l.Acquire(context.Background(), "alpha", "m"); err == nil { + t.Error("Acquire on an over-tracked pair succeeded; want ErrQueueFull") + } + rel1() + rel1() // idempotent + rel2() + rel3() + if n := l.InFlight("alpha", "m"); n != 0 { + t.Errorf("in flight after release = %d, want 0", n) + } +} + +// A waiter gets a slot only once in flight is back under parallel: releasing a tracked request +// while the pair is still over its limit must not hand the slot on. +func TestTrackReleaseHandsOverOnlyUnderTheLimit(t *testing.T) { + l := limiter.New() + l.Configure("alpha", "m", 1, 1) + relA := l.Track("alpha", "m") + relB := l.Track("alpha", "m") // in flight 2, parallel 1 + + got := make(chan func(), 1) + go func() { + rel, _, err := l.Acquire(context.Background(), "alpha", "m") + if err != nil { + t.Error(err) + close(got) + return + } + got <- rel + }() + waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 }) + + relA() // in flight 1 == parallel: still no free slot + select { + case <-got: + t.Fatal("waiter got a slot while in flight was still at parallel") + case <-time.After(100 * time.Millisecond): + } + if n := l.InFlight("alpha", "m"); n != 1 { + t.Fatalf("in flight = %d after one release, want 1", n) + } + + relB() // now the slot is free: hand it to the waiter + select { + case rel := <-got: + if rel == nil { + t.Fatal("waiter failed") + } + if n := l.InFlight("alpha", "m"); n != 1 { + t.Errorf("in flight = %d with the waiter running, want 1", n) + } + rel() + case <-time.After(2 * time.Second): + t.Fatal("waiter never got the freed slot") + } + if n := l.InFlight("alpha", "m"); n != 0 { + t.Errorf("in flight at the end = %d, want 0", n) + } +} diff --git a/docs/plans/v2.3/_files/internal/proxy/affinity_test.go b/docs/plans/v2.3/_files/internal/proxy/affinity_test.go new file mode 100644 index 0000000..76db349 --- /dev/null +++ b/docs/plans/v2.3/_files/internal/proxy/affinity_test.go @@ -0,0 +1,153 @@ +package proxy_test + +// v2.3 task 02: affinity = "route" puts every request on the route (every conversation, every +// control call) on one lease, so one host; queue = false counts the route's requests on the host +// without ever holding or refusing them, because the client pins its own llama-server slot and +// the server's queue is the one that must show it. + +import ( + "net/http" + "sync" + "testing" + "time" +) + +const affinityHosts = ` +listen = "127.0.0.1:1" +queue_max = 0 +lease_idle = "30m" +[hosts.alpha] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 1 } } +[hosts.beta] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 1 } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "shared" +[routes.bm] +hosts = ["alpha", "beta"] +default_model = "shared" +affinity = "route" +queue = false +[routes."agent-*"] +hosts = ["alpha", "beta"] +default_model = "shared" +affinity = "route" +` + +func TestRouteAffinityPutsEverythingOnOneHost(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + + seen := map[string]int{} + note := func(what string, resp *http.Response) { + body := drain(resp) + if resp.StatusCode != 200 { + t.Fatalf("%s: %d %s", what, resp.StatusCode, body) + } + seen[resp.Header.Get("X-Crossbar-Host")]++ + } + // Different conversations (different fingerprints), then control calls without any. + for id := 1; id <= 4; id++ { + note("chat", r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(id, 1))) + } + note("slots", r.do(http.MethodGet, "/bm/slots?model=shared", "")) + note("props", r.do(http.MethodGet, "/bm/props?model=shared", "")) + note("control", r.do(http.MethodPost, "/bm/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`)) + if len(seen) != 1 { + t.Fatalf("route-affinity requests spread over %v, want one host", seen) + } + + // Templated concrete routes each get their own route lease, and each is internally sticky. + for _, route := range []string{"agent-a", "agent-b", "agent-c"} { + hosts := map[string]bool{} + for id := 1; id <= 3; id++ { + resp := r.do(http.MethodPost, "/"+route+"/v1/chat/completions", conversation(id, 1)) + drain(resp) + hosts[resp.Header.Get("X-Crossbar-Host")] = true + } + if len(hosts) != 1 { + t.Errorf("%s spread over %v, want one host", route, hosts) + } + } +} + +func TestQueueFalseNeitherHoldsNorRefuses(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + alpha.delay, beta.delay = 400*time.Millisecond, 400*time.Millisecond + r := newRig(t, affinityHosts, alpha, beta) + + // parallel = 1 and queue_max = 0: a queueing route would refuse the second and third. + var wg sync.WaitGroup + codes := make(chan int, 3) + start := time.Now() + for id := 1; id <= 3; id++ { + wg.Add(1) + go func(id int) { + defer wg.Done() + resp := r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(id, 1)) + drain(resp) + codes <- resp.StatusCode + }(id) + } + // While they run, the host carries all three and a queueing route sees it full. + var host string + waitUntil(t, func() bool { + for _, h := range []string{"alpha", "beta"} { + if r.lim.InFlight(h, "shared") == 3 { + host = h + return true + } + } + return false + }) + if n := r.lim.FreeSlots(host); n != 0 { + t.Errorf("free slots on %s = %d while bm runs three, want 0", host, n) + } + wg.Wait() + close(codes) + for c := range codes { + if c != 200 { + t.Errorf("queue = false request: %d, want 200", c) + } + } + // Concurrent, not serialised behind one slot: three 400 ms answers well under 1.2 s. + 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) + } + // Accounting is unchanged: each chat is still a row. + waitUntil(t, func() bool { return r.rows("bm") == 3 }) +} + +// The default is unchanged: two conversations on a conversation-affinity route may land on +// different hosts (they start where there is most room). +func TestConversationAffinityStillSpreads(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + alpha.delay, beta.delay = 300*time.Millisecond, 300*time.Millisecond + r := newRig(t, affinityHosts, alpha, beta) + var wg sync.WaitGroup + var mu sync.Mutex + hosts := map[string]bool{} + for id := 1; id <= 2; id++ { + wg.Add(1) + go func(id int) { + defer wg.Done() + resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(id, 1)) + drain(resp) + mu.Lock() + hosts[resp.Header.Get("X-Crossbar-Host")] = true + mu.Unlock() + }(id) + time.Sleep(50 * time.Millisecond) // let the first take its slot so the second sees one host full + } + wg.Wait() + if len(hosts) != 2 { + t.Errorf("two concurrent conversations on route r used %v, want both hosts", hosts) + } +} diff --git a/docs/plans/v2.3/_files/internal/proxy/control_test.go b/docs/plans/v2.3/_files/internal/proxy/control_test.go new file mode 100644 index 0000000..c0d150e --- /dev/null +++ b/docs/plans/v2.3/_files/internal/proxy/control_test.go @@ -0,0 +1,176 @@ +package proxy_test + +// v2.3 task 01: control-plane requests. A client that manages its own slots (Boxmaker) polls +// /slots, reads /props, tokenizes and steers a running completion through +// /v1/chat/completions/control. Those calls follow the route's lease like any other request but +// must never wait for, or take, a slot: /control is sent while the client's own stream holds one. + +import ( + "context" + "net/http" + "strings" + "testing" + "time" +) + +// controlClient gives every control call a short deadline: a call that queues behind a full host +// is the bug, and it must fail the test rather than hang it. +var controlClient = &http.Client{Timeout: 2 * time.Second} + +func (r *rig) do(method, path, body string) *http.Response { + r.t.Helper() + var rd *strings.Reader + if body != "" { + rd = strings.NewReader(body) + } + var req *http.Request + var err error + if rd != nil { + req, err = http.NewRequest(method, r.front.URL+path, rd) + req.Header.Set("Content-Type", "application/json") + } else { + req, err = http.NewRequest(method, r.front.URL+path, nil) + } + if err != nil { + r.t.Fatal(err) + } + resp, err := controlClient.Do(req) + if err != nil { + r.t.Fatalf("%s %s: %v", method, path, err) + } + return resp +} + +func (r *rig) rows(route string) int64 { + r.t.Helper() + counts, err := r.store.StatusCounts(time.Time{}) + if err != nil { + r.t.Fatal(err) + } + var n int64 + for _, c := range counts { + if c.Route == route { + n += c.Count + } + } + return n +} + +// A GET names its model in the query string: /slots?model=alpha-only must reach the host that +// has alpha-only loaded, not whichever host the route's default model would pick. +func TestGetModelComesFromTheQuery(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + + resp := r.do(http.MethodGet, "/r/slots?model=alpha-only", "") + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != "alpha" { + t.Fatalf("GET /r/slots?model=alpha-only: %d on %q, want 200 on alpha", resp.StatusCode, resp.Header.Get("X-Crossbar-Host")) + } + if got := alpha.lastReq(); got.method != "GET" || got.path != "/slots?model=alpha-only" { + t.Errorf("alpha saw %s %s, want GET /slots?model=alpha-only", got.method, got.path) + } + // The same for beta-only, so a lucky default cannot pass the test. + resp = r.do(http.MethodGet, "/r/slots?model=beta-only", "") + drain(resp) + if resp.Header.Get("X-Crossbar-Host") != "beta" { + t.Errorf("GET /r/slots?model=beta-only went to %q, want beta", resp.Header.Get("X-Crossbar-Host")) + } +} + +// /slots and /tokenize are proxied; the per-slot actions under /slots/ (save, restore, erase) +// are not. +func TestControlPathsAllowed(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + for _, tc := range []struct { + method, path, body string + want int + }{ + {http.MethodGet, "/r/slots", "", 200}, + {http.MethodGet, "/r/slots?model=shared", "", 200}, + {http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`, 200}, + {http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`, 200}, + {http.MethodGet, "/r/slots/0", "", 404}, + {http.MethodPost, "/r/slots/0?action=erase", "", 404}, + {http.MethodPost, "/r/slots/0?action=save", `{"filename":"x"}`, 404}, + } { + resp := r.do(tc.method, tc.path, tc.body) + body := drain(resp) + if resp.StatusCode != tc.want { + t.Errorf("%s %s: %d %s, want %d", tc.method, tc.path, resp.StatusCode, body, tc.want) + } + } +} + +// With every slot on both hosts taken and the queue full, control-plane calls still go straight +// through: no 503, no wait, no slot taken, no accounting row. +func TestControlRequestsNeverTakeASlot(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + + // Take every "shared" slot (parallel 2 on each host) and the one queue place per host. + var releases []func() + for _, host := range []string{"alpha", "beta"} { + for i := 0; i < 2; i++ { + rel, _, err := r.lim.Acquire(context.Background(), host, "shared") + if err != nil { + t.Fatal(err) + } + releases = append(releases, rel) + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go func() { _, _, _ = r.lim.Acquire(ctx, host, "shared") }() + waitUntil(t, func() bool { return r.lim.Queued(host, "shared") == 1 }) + } + defer func() { + for _, rel := range releases { + rel() + } + }() + + for _, tc := range []struct{ method, path, body string }{ + {http.MethodGet, "/r/slots?model=shared", ""}, + {http.MethodGet, "/r/props?model=shared", ""}, + {http.MethodHead, "/r/props?model=shared", ""}, + {http.MethodGet, "/r/v1/models", ""}, + {http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`}, + {http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`}, + } { + resp := r.do(tc.method, tc.path, tc.body) + body := drain(resp) + if resp.StatusCode != 200 { + t.Errorf("%s %s with the host full: %d %s, want 200", tc.method, tc.path, resp.StatusCode, body) + } + } + for _, host := range []string{"alpha", "beta"} { + if n := r.lim.InFlight(host, "shared"); n != 2 { + t.Errorf("%s in flight = %d after control calls, want 2 (control takes no slot)", host, n) + } + } + time.Sleep(100 * time.Millisecond) // a row is written after the answer; give a stray one time to land + if n := r.rows("r"); n != 0 { + t.Errorf("control calls wrote %d accounting rows, want 0", n) + } + + // A chat completion on the same full route still queues or is refused as before: the bypass + // is for control calls only. + resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1)) + drain(resp) + if resp.StatusCode != http.StatusServiceUnavailable { + t.Errorf("chat on a full route: %d, want 503 (queue full)", resp.StatusCode) + } +} + +// A chat completion is not a control call just because its path starts the same way. +func TestChatIsNotControl(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1)) + drain(resp) + if resp.StatusCode != 200 { + t.Fatalf("chat: %d", resp.StatusCode) + } + waitUntil(t, func() bool { return r.rows("r") == 1 }) // the row lands just after the answer +} diff --git a/docs/plans/v2.3/_files/internal/proxy/ctxguard_test.go b/docs/plans/v2.3/_files/internal/proxy/ctxguard_test.go new file mode 100644 index 0000000..92e3318 --- /dev/null +++ b/docs/plans/v2.3/_files/internal/proxy/ctxguard_test.go @@ -0,0 +1,164 @@ +package proxy_test + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/proxy" +) + +// ctxUpstream is a fake router that reports a context size in /props and echoes completions. +func ctxUpstream(t *testing.T, name string, nCtx, slots int) *upstream { + u := &upstream{name: name} + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"shared"}]}`) }) + mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprintf(w, `{"default_generation_settings":{"n_ctx":%d},"total_slots":%d}`, nCtx, slots) + }) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + u.hits.Add(1) + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"choices":[{"message":{"role":"assistant","content":"ok"}}],"usage":{"prompt_tokens":1,"completion_tokens":1}}`) + }) + u.srv = httptest.NewServer(mux) + t.Cleanup(u.srv.Close) + return u +} + +const ctxHosts = ` +listen = "127.0.0.1:1" +queue_max = 2 +[hosts.small] +base_url = %q +weight = 10.0 +models = { "shared" = { parallel = 2 } } +[hosts.big] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 1 } } +[routes.r] +hosts = ["small", "big"] +default_model = "shared" +` + +// bodyOfTokens builds a chat body whose byte size implies roughly n tokens under the guard's +// estimate (bytes/4 × 1.2): n tokens ≈ 3.33 n bytes ≈ 2n/3 five-byte words. +func bodyOfTokens(n int) string { + text := strings.Repeat("word ", n*2/3) + return fmt.Sprintf(`{"model":"shared","stream":false,"messages":[{"role":"user","content":"%s"}]}`, text) +} + +func TestOversizedPromptMovesToAHostWhereItFits(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) // 4096 per slot + big := ctxUpstream(t, "big", 131072, 1) // 131072 per slot + r := newRig(t, ctxHosts, small, big) + // A small prompt starts on `small` (weight 10). + resp := r.post("/r/v1/chat/completions", bodyOfTokens(100)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" { + t.Fatalf("small prompt went to %q, want small", resp.Header.Get(proxy.HostHeader)) + } + // A new conversation with ~10k tokens does not fit small's 4096-token slot: it must be + // placed on big, with the reason visible in a header. + resp = r.post("/r/v1/chat/completions", bodyOfTokens(10000)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" { + t.Fatalf("oversized prompt: %d from %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader)) + } + if got := resp.Header.Get(proxy.CtxHeader); !strings.HasPrefix(got, "moved") { + t.Errorf("%s = %q, want moved:… ", proxy.CtxHeader, got) + } +} + +func TestOversizedPromptWithNoFitIs400(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) + tiny := ctxUpstream(t, "big", 4096, 2) // also too small + r := newRig(t, ctxHosts, small, tiny) + resp := r.post("/r/v1/chat/completions", bodyOfTokens(10000)) + body := drain(resp) + if resp.StatusCode != http.StatusBadRequest { + t.Fatalf("status %d body %s, want 400", resp.StatusCode, body) + } + // v2.3: llama-server's own shape for this error, so a client handles crossbar's refusal the + // way it handles the server's (Boxmaker keys on error.type; the error JSON must come first). + if !strings.HasPrefix(body, `{"error":`) { + t.Errorf("body must start with the error object: %s", body) + } + var e struct { + Error struct { + Code int `json:"code"` + Type string `json:"type"` + Message string `json:"message"` + NPromptTokens float64 `json:"n_prompt_tokens"` + NCtx float64 `json:"n_ctx"` + } `json:"error"` + } + if err := json.Unmarshal([]byte(body), &e); err != nil || e.Error.Code != 400 || e.Error.Type != "exceed_context_size_error" || e.Error.Message != "prompt too large" { + t.Fatalf("body = %s, want {\"error\":{\"code\":400,\"type\":\"exceed_context_size_error\",\"message\":\"prompt too large\",…}}", body) + } + if est := e.Error.NPromptTokens; est < 8000 || est > 13000 { + t.Errorf("n_prompt_tokens = %v, want roughly 10000 tokens", est) + } + if max := e.Error.NCtx; max != 4096 { + t.Errorf("n_ctx = %v, want the largest per-slot context among the route's hosts (4096)", max) + } + if ct := resp.Header.Get("Content-Type"); !strings.HasPrefix(ct, "application/json") { + t.Errorf("Content-Type = %q, want application/json", ct) + } + if small.hits.Load()+tiny.hits.Load() != 0 { + t.Errorf("a refused prompt must not reach any upstream") + } +} + +func TestUnknownContextNeverBlocks(t *testing.T) { + // /props missing on both hosts: NCtx 0 means "unknown", and the guard must stay out of the way. + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.post("/r/v1/chat/completions", bodyOfTokens(50000)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.CtxHeader) != "" { + t.Errorf("unknown context: %d %q, want 200 and no ctx header", resp.StatusCode, resp.Header.Get(proxy.CtxHeader)) + } +} + +// grow appends later turns to a conversation body without touching its system prompt or first +// user message, so the fingerprint — and therefore the lease — stays the same. +func grow(body string, words int) string { + turn := `,{"role":"assistant","content":"ok"},{"role":"user","content":"` + strings.Repeat("x ", words) + `"}` + return strings.Replace(body, `]}`, turn+`]}`, 1) +} + +func TestStickyLeaseSurvivesGrowthUntilItDoesNotFit(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) + big := ctxUpstream(t, "big", 131072, 1) + r := newRig(t, ctxHosts, small, big) + body := bodyOfTokens(100) + resp := r.post("/r/v1/chat/completions", body) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" { + t.Fatal("setup: first turn must be on small") + } + // Same conversation, a later turn well under 4096 tokens: stays. + resp = r.post("/r/v1/chat/completions", grow(body, 500)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Errorf("turn 2: %q %q, want small reused", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } + // A turn that outgrows the slot moves the lease — once — and the move is visible in the header. + huge := grow(body, 30000) + resp = r.post("/r/v1/chat/completions", huge) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" || !strings.HasPrefix(resp.Header.Get(proxy.CtxHeader), "moved") { + t.Fatalf("outgrown turn: %d %q ctx=%q, want 200 from big with a moved header", resp.StatusCode, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.CtxHeader)) + } + resp = r.post("/r/v1/chat/completions", huge) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "big" || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Errorf("after the move the lease is on big: %q %q", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } +} diff --git a/docs/plans/v2.3/_files/internal/proxy/listener_test.go b/docs/plans/v2.3/_files/internal/proxy/listener_test.go new file mode 100644 index 0000000..650cfc6 --- /dev/null +++ b/docs/plans/v2.3/_files/internal/proxy/listener_test.go @@ -0,0 +1,112 @@ +package proxy_test + +// v2.3 task 03: Handler.ForRoute serves one route with unprefixed paths, for a route's dedicated +// listener. + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/proxy" +) + +// dedicated serves r's route on its own test server, sharing r's health, leases, limiter and +// store, as main does for a route with listen set. +func dedicated(t *testing.T, r *rig, route string) *httptest.Server { + p := proxy.New(r.cfg, r.health, r.leases, r.lim, r.store, nil) + srv := httptest.NewServer(p.ForRoute(route)) + t.Cleanup(srv.Close) + return srv +} + +func call(t *testing.T, method, url, body string, hdr ...string) (*http.Response, string) { + t.Helper() + var req *http.Request + if body != "" { + req, _ = http.NewRequest(method, url, strings.NewReader(body)) + req.Header.Set("Content-Type", "application/json") + } else { + req, _ = http.NewRequest(method, url, nil) + } + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } + resp, err := controlClient.Do(req) + if err != nil { + t.Fatalf("%s %s: %v", method, url, err) + } + return resp, drain(resp) +} + +func TestForRouteServesUnprefixedPaths(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "bm") + + resp, body := call(t, http.MethodPost, srv.URL+"/v1/chat/completions", conversation(1, 1)) + if resp.StatusCode != 200 { + t.Fatalf("chat on the dedicated listener: %d %s", resp.StatusCode, body) + } + host := resp.Header.Get("X-Crossbar-Host") + up := map[string]*upstream{"alpha": alpha, "beta": beta}[host] + if up == nil || up.lastReq().path != "/v1/chat/completions" { + t.Fatalf("upstream %q saw %+v, want /v1/chat/completions unchanged", host, up.lastReq()) + } + resp, _ = call(t, http.MethodGet, srv.URL+"/slots?model=shared", "") + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host || up.lastReq().path != "/slots?model=shared" { + t.Errorf("/slots: %d on %q (last %+v), want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), up.lastReq(), host) + } + resp, _ = call(t, http.MethodPost, srv.URL+"/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`) + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host { + t.Errorf("/control: %d on %q, want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), host) + } + // The chat is accounted to the route the listener serves. + waitUntil(t, func() bool { return r.rows("bm") == 1 }) + // The same route through the main listener shares the lease: same host. + resp = r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(2, 1)) + drain(resp) + if resp.Header.Get("X-Crossbar-Host") != host { + t.Errorf("main listener /bm went to %q, dedicated to %q; one route, one lease", resp.Header.Get("X-Crossbar-Host"), host) + } +} + +func TestForRouteRefusals(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "bm") + for _, tc := range []struct { + name, method, path string + hdr []string + want int + msg string + }{ + {"a prefixed path is not stripped", http.MethodGet, "/bm/v1/models", nil, 404, "not found"}, + {"no admin here", http.MethodGet, "/_crossbar/hosts", nil, 404, "not found"}, + {"root", http.MethodGet, "/", nil, 404, "not found"}, + {"header naming another route", http.MethodGet, "/v1/models", []string{"X-Crossbar-Route", "r"}, 400, "conflicting route"}, + } { + resp, body := call(t, tc.method, srv.URL+tc.path, "", tc.hdr...) + var e map[string]string + if resp.StatusCode != tc.want || json.Unmarshal([]byte(body), &e) != nil || e["error"] != tc.msg { + t.Errorf("%s: %d %s, want %d %q", tc.name, resp.StatusCode, body, tc.want, tc.msg) + } + } + // A header naming this same route is harmless. + resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "", "X-Crossbar-Route", "bm") + if resp.StatusCode != 200 { + t.Errorf("header naming the listener's own route: %d %s, want 200", resp.StatusCode, body) + } +} + +func TestForRouteUnknownRoute(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "nope") + resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "") + if resp.StatusCode != 404 || !strings.Contains(body, "unknown route") { + t.Errorf("ForRoute(unknown): %d %s, want 404 unknown route", resp.StatusCode, body) + } +} diff --git a/docs/plans/v2.3/_files/internal/proxy/proxy_test.go b/docs/plans/v2.3/_files/internal/proxy/proxy_test.go new file mode 100644 index 0000000..d5f23bb --- /dev/null +++ b/docs/plans/v2.3/_files/internal/proxy/proxy_test.go @@ -0,0 +1,292 @@ +package proxy_test + +// v1 acceptance tests for the proxy: leases, queueing, accounting, header route override. +// They drive the whole handler over real HTTP against fake upstreams; only what a client or an +// operator can observe is asserted (status codes, headers, the accounting rows, the health table). +// The rig, the fake upstream and the request helpers live in helpers_test.go. + +import ( + "encoding/json" + "net/http" + "strings" + "sync" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/proxy" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +func TestConversationIsStickyAndLeaseHeaderTellsWhy(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + first := r.post("/r/v1/chat/completions", conversation(1, 1)) + drain(first) + host := first.Header.Get(proxy.HostHeader) + if first.StatusCode != 200 || host != "beta" { // beta: same free slots, double weight + t.Fatalf("first turn: %d from %q, want 200 from beta", first.StatusCode, host) + } + if got := first.Header.Get(proxy.LeaseHeader); got != "new" { + t.Errorf("%s = %q on the first turn, want new", proxy.LeaseHeader, got) + } + // Take alpha's slots away as a "better host" signal: it must not matter, the lease holds. + for turn := 2; turn <= 6; turn++ { + resp := r.post("/r/v1/chat/completions", conversation(1, turn)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != host || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Fatalf("turn %d: host %q lease %q, want %q reused", turn, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader), host) + } + } + if alpha.hits.Load() != 0 || beta.hits.Load() != 6 { + t.Errorf("hits alpha=%d beta=%d, want 0 and 6", alpha.hits.Load(), beta.hits.Load()) + } +} + +// spreadHosts: beta is preferred (weight 10) until both of its "shared" slots are busy; then +// alpha (2 free × 1) beats beta (0 free × 10), and a new conversation must start on alpha. +const spreadHosts = ` +listen = "127.0.0.1:1" +queue_max = 4 +lease_idle = "30m" +[hosts.alpha] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 2 } } +[hosts.beta] +base_url = %q +weight = 10.0 +models = { "shared" = { parallel = 2 } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "shared" +` + +func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + beta.delay = 400 * time.Millisecond + r := newRig(t, spreadHosts, alpha, beta) + // Two slow conversations occupy beta's two "shared" slots… + var wg sync.WaitGroup + for i := 1; i <= 2; i++ { + wg.Add(1) + go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) + // 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 }) + } + // …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)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "alpha" { + t.Errorf("third conversation went to %q, want alpha (free slots beat weight)", resp.Header.Get(proxy.HostHeader)) + } + wg.Wait() + if beta.hits.Load() != 2 || alpha.hits.Load() != 1 { + t.Errorf("hits beta=%d alpha=%d, want 2 and 1", beta.hits.Load(), alpha.hits.Load()) + } +} + +// waitUntil polls cond every 5 ms 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(5 * time.Millisecond) + } + t.Fatal("condition not reached within two seconds") +} + +func TestQueueFullIs503(t *testing.T) { + alpha := newUpstream(t, "alpha") + alpha.delay = 400 * time.Millisecond + r := newRig(t, ` +listen = "127.0.0.1:1" +queue_max = 1 +[hosts.alpha] +base_url = %q +models = { "shared" = { parallel = 1 } } +[routes.r] +hosts = ["alpha"] +default_model = "shared" +`, alpha) + codes := make(chan int, 3) + fire := func(i int) { + go func() { + resp := r.post("/r/v1/chat/completions", conversation(i, 1)) + drain(resp) + codes <- resp.StatusCode + }() + } + // Arrival order is enforced by watching the limiter, not by sleeping: 1 runs, 2 queues, + // 3 finds the queue full. + fire(1) + waitUntil(t, func() bool { return r.lim.InFlight("alpha", "shared") == 1 }) + fire(2) + waitUntil(t, func() bool { return r.lim.Queued("alpha", "shared") == 1 }) + fire(3) + got := map[int]int{} + for i := 0; i < 3; i++ { + got[<-codes]++ + } + if got[200] != 2 || got[503] != 1 { + t.Fatalf("status counts = %v, want two 200 and one 503", got) + } + // Rows are written after each response completes; allow the store a moment to catch up. + var rows []store.UsageRow + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + rows, _ = r.store.Usage(time.Time{}, store.ByRoute) + if len(rows) == 1 && rows[0].Requests == 3 { + break + } + time.Sleep(20 * time.Millisecond) + } + if len(rows) != 1 || rows[0].Requests != 3 || rows[0].Errors != 1 { + t.Fatalf("usage = %+v, want 3 requests, 1 error (the 503 is recorded too)", rows) + } + if rows[0].QueuedMs <= 0 { + t.Errorf("the queued request must record its wait: %+v", rows[0]) + } +} + +func TestUnhealthyHostReleasesAndMoves(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // lands on beta + beta.srv.Close() + resp := r.post("/r/v1/chat/completions", conversation(1, 2)) + drain(resp) + if resp.StatusCode != http.StatusBadGateway { + t.Fatalf("first request after beta died: %d, want 502", resp.StatusCode) + } + if s, _ := r.health.Get("beta"); s.Healthy { + t.Fatalf("beta must be marked down after the 502") + } + resp = r.post("/r/v1/chat/completions", conversation(1, 3)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" || resp.Header.Get(proxy.LeaseHeader) != "new" { + t.Errorf("after the move: %d from %q lease %q, want 200 alpha new", resp.StatusCode, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } + ev, _ := r.store.Events(time.Time{}, 10) + var reasons []string + for _, e := range ev { + reasons = append(reasons, e.Reason) + } + if len(reasons) != 2 || reasons[0] != store.ReasonNew || reasons[1] != store.ReasonUnhealthy { + t.Errorf("lease events = %v, want [new unhealthy]", reasons) + } +} + +func TestAccountingRowsFromUsageAndTimings(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // non-streamed + drain(r.post("/r/v1/chat/completions", strings.Replace(conversation(1, 2), `"stream":false`, `"stream":true`, 1))) // streamed + deadline := time.Now().Add(2 * time.Second) + var rows []store.UsageRow + for time.Now().Before(deadline) { + rows, _ = r.store.Usage(time.Time{}, store.ByHost) + if len(rows) == 1 && rows[0].Requests == 2 { + break + } + time.Sleep(20 * time.Millisecond) + } + if len(rows) != 1 || rows[0].Requests != 2 { + t.Fatalf("usage by host = %+v, want one host with 2 requests (rows may be written after the response completes, within 2 s)", rows) + } + u := rows[0] + if u.PromptTokens != 300 || u.CachedTokens != 240 || u.CompletionTokens != 30 { + t.Errorf("tokens = prompt %d cached %d completion %d, want 300/240/30 (100+200, 90+150, 10+20)", u.PromptTokens, u.CachedTokens, u.CompletionTokens) + } + if u.BusyMs <= 0 || u.Errors != 0 { + t.Errorf("busy %d errors %d", u.BusyMs, u.Errors) + } + if got := u.CacheHitRatio(); got < 0.79 || got > 0.81 { + t.Errorf("cache hit ratio = %v, want 0.8", got) + } +} + +func TestStreamIsUnalteredWhileTeed(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.post("/r/v1/chat/completions", strings.Replace(conversation(9, 1), `"stream":false`, `"stream":true`, 1)) + body := drain(resp) + want := 0 + for _, line := range strings.Split(body, "\n") { + if strings.HasPrefix(line, "data: ") { + want++ + } + } + if want != 5 || !strings.HasSuffix(strings.TrimSpace(body), "data: [DONE]") { + t.Errorf("client must receive every SSE line untouched (3 deltas, usage, DONE); got %d data lines:\n%s", want, body) + } +} + +func TestHeaderRouteOverride(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + // The header names the route; the path has none. + resp := r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other") + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" { + t.Errorf("header route 'other' (alpha only): %d from %q", resp.StatusCode, resp.Header.Get(proxy.HostHeader)) + } + if alpha.lastReq().path != "/v1/chat/completions" { + t.Errorf("upstream path = %q", alpha.lastReq().path) + } + // A path route and a header route that disagree: the header is the operator's intent → 400. + resp = r.post("/r/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other") + if drain(resp); resp.StatusCode != 400 { + t.Errorf("conflicting route in path and header: %d, want 400", resp.StatusCode) + } + resp = r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "nope") + if drain(resp); resp.StatusCode != 404 { + t.Errorf("unknown header route: %d, want 404", resp.StatusCode) + } +} + +func TestV0BehaviourStillHolds(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + for _, tc := range []struct { + method, path string + want int + msg string + }{ + {http.MethodGet, "/", 400, "missing route"}, + {http.MethodGet, "/nope/v1/models", 404, "unknown route"}, + // v2.3: /slots itself is proxied (a control-plane path); its per-slot actions are not. + {http.MethodGet, "/r/slots/0", 404, "not found"}, + {http.MethodGet, "/r/metrics", 404, "not found"}, + {http.MethodGet, "/r/_crossbar/hosts", 404, "not found"}, + } { + req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + body := drain(resp) + var e map[string]string + if resp.StatusCode != tc.want || json.Unmarshal([]byte(body), &e) != nil || e["error"] != tc.msg { + t.Errorf("%s: %d %s, want %d %q", tc.path, resp.StatusCode, body, tc.want, tc.msg) + } + } + big := strings.Repeat("x", proxy.MaxBody+1) + resp := r.post("/r/v1/chat/completions", big) + if drain(resp); resp.StatusCode != 413 { + t.Errorf("oversize body: %d, want 413", resp.StatusCode) + } + // GET pass-through with query string, Host and X-Forwarded-For as in v0. + resp, err := http.Get(r.front.URL + "/r/v1/models?x=1") + if err != nil { + t.Fatal(err) + } + drain(resp) + host := resp.Header.Get(proxy.HostHeader) + u := map[string]*upstream{"alpha": alpha, "beta": beta}[host] + if u == nil || u.lastReq().path != "/v1/models?x=1" || u.lastReq().host != strings.TrimPrefix(u.srv.URL, "http://") || u.lastReq().xff == "" { + t.Errorf("GET pass-through: host %q last %+v", host, u.lastReq()) + } +} diff --git a/docs/plans/v2.3/_files/tools/smoke.sh b/docs/plans/v2.3/_files/tools/smoke.sh new file mode 100644 index 0000000..240a2ab --- /dev/null +++ b/docs/plans/v2.3/_files/tools/smoke.sh @@ -0,0 +1,71 @@ +#!/bin/sh +# Smoke run (v2.3): everything v1 checked, plus the context guard, wake-on-LAN, identity gating and +# a route's dedicated listener. +# 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="" +sed -e "s#^db .*#db = \"$tmp/crossbar.db\"#" -e 's#^identity .*#identity = "header"#' -e 's#^\# peers = \["talos"\]#peers = ["talos"]#' example.toml > "$tmp/crossbar.toml" +# alpha: small context (4096 per slot = 8192/2); beta: large, sleeps until woken +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 -n-ctx 8192 -slots 2 >"$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" -n-ctx 131072 -slots 2 -wol-listen 127.0.0.1:19082 -wol-mac aa:bb:cc:dd:ee:02 >"$tmp/beta.log" 2>&1 & pids="$pids $!" +touch "$tmp/beta.down" # beta starts "asleep" +bin/crossbar -config "$tmp/crossbar.toml" >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" +sleep 2.5 # two polls: alpha healthy, beta down +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"; } +big() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"user","content":"%s"}]}' "$(head -c 40000 /dev/zero | tr '\0' 'x')"; } +hdrs() { curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host} %header{X-Crossbar-Lease}' "$@"; } + +# 1. v1 behaviour: with beta asleep, opencode-a goes to alpha +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 alpha new" ] || fail "with beta asleep conversation A should be '200 alpha new', got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{[^}]*"n_ctx":8192' || fail "hosts view does not show alpha n_ctx 8192: $(curl -s $base/_crossbar/hosts)" + +# 2. context guard: a ~12k-token prompt does not fit alpha's 4096-token slot; beta is asleep and +# wakeable, so crossbar must send the magic packet, wait for beta, and place the prompt there. +start=$(date +%s) +h=$(hdrs -m 40 -X POST -H 'Content-Type: application/json' -d "$(big)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta new" ] || fail "oversized prompt should wake beta and land there, got '$h' after $(( $(date +%s) - start ))s" +grep -q "magic packet received" "$tmp/beta.log" || fail "beta never saw a magic packet" +curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":true' || fail "beta not healthy after wake" + +# 3. with beta awake, a prompt that fits nowhere is a 400 (both slots too small? no — beta fits): +# check the guard's refusal with a prompt beyond beta's 65536-per-slot too +# (the body goes through a file: a 300 KB string cannot be a single argv element on Linux) +{ printf '{"model":"ornith-1.5-35b-a3b","messages":[{"role":"user","content":"'; head -c 300000 /dev/zero | tr '\0' 'x'; printf '"}]}'; } > "$tmp/toolarge-req.json" +h=$(curl -s -o "$tmp/toolarge.json" -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d @"$tmp/toolarge-req.json" "$base/opencode-a/v1/chat/completions") +[ "$h" = "400" ] && grep -q '"prompt too large"' "$tmp/toolarge.json" || fail "300 KB prompt should be 400 prompt too large, got $h $(cat "$tmp/toolarge.json")" + +# 4. identity: hermes-x is locked to peer talos (header mode) +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x without a peer header should be 403, got $h" +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: titan' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x as titan should be 403, got $h" +h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +case "$h" in 200*) ;; *) fail "hermes-x as talos should be 200, got '$h'";; esac +h=$(curl -s -o /dev/null -w '%{http_code}' "$base/_crossbar/hosts"); [ "$h" = "200" ] || fail "admin must not be gated, got $h" + +# 5. v1 regression: streaming still incremental, usage and metrics present +start=$(date +%s%N) +curl -sN -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \ + "$base/hermes-x/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" +sleep 1 +curl -s "$base/_crossbar/usage?by=host" | grep -q '"cached_tokens":[1-9]' || fail "usage has no cached tokens" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_requests_total{route="hermes-x",host="beta",status="403"}' && fail "403s are refused before a lease and must not be counted as requests" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing beta health" +# 6. v2.3: boxmaker-a has its own listener. Paths are unprefixed, the chat and a control call land +# on the same host (affinity = "route"), and the admin API is not served there. +lb=http://127.0.0.1:17801 +h1=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$lb/v1/chat/completions") +h2=$(hdrs "$lb/props?model=ornith-1.5-35b-a3b") +case "$h1" in 200*) ;; *) fail "chat on the dedicated listener should be 200, got '$h1'";; esac +[ "$(echo "$h1" | cut -d' ' -f2)" = "$(echo "$h2" | cut -d' ' -f2)" ] || fail "chat went to '$h1', /props to '$h2': one route, one host" +h=$(curl -s -o /dev/null -w '%{http_code}' "$lb/_crossbar/hosts"); [ "$h" = "404" ] || fail "admin must not be served on a dedicated listener, got $h" +h=$(curl -s -o /dev/null -w '%{http_code}' "$lb/boxmaker-a/v1/models"); [ "$h" = "404" ] || fail "a prefixed path on the dedicated listener should be 404, got $h" + +echo "smoke: ok (stream spread $((lastms - firstms)) ms)"