Merge v2.3: control-plane calls, route affinity, queue = false, route listeners, llama-server ctx error
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -59,6 +59,9 @@ hosts = ["beta", "alpha"]
|
|||||||
| `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` | Candidate hosts, tried in order until one is healthy; a conversation leases one of them. |
|
| `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. |
|
||||||
|
| `routes.<name>.affinity` | `"conversation"` (default, one lease per conversation) or `"route"` (one lease for the whole route); see "Clients that manage their own slots". |
|
||||||
|
| `routes.<name>.queue` | `false` leaves queueing to the client's own llama-server slot; the default counts requests in crossbar's per-(host, model) queue. |
|
||||||
|
| `routes.<name>.listen` | A host:port for the route's own listener, every request there is this route; see "Clients that manage their own slots". |
|
||||||
| `identity` | `"off"` (default), `"tailscale"`, or `"header"`; see below. |
|
| `identity` | `"off"` (default), `"tailscale"`, or `"header"`; see below. |
|
||||||
| `hosts.<name>.wake` | A wake-on-LAN target (`mac`, `broadcast`, `wait`) so crossbar can rouse a sleeping host when nothing else can take a new lease. |
|
| `hosts.<name>.wake` | A wake-on-LAN target (`mac`, `broadcast`, `wait`) so crossbar can rouse a sleeping host when nothing else can take a new lease. |
|
||||||
| `routes.<name>.peers` | The tailnet nodes allowed to reach the route, with `identity = "tailscale"`; see below. |
|
| `routes.<name>.peers` | The tailnet nodes allowed to reach the route, with `identity = "tailscale"`; see below. |
|
||||||
@@ -104,6 +107,46 @@ curl -H 'X-Crossbar-Route: opencode-a' \
|
|||||||
https://crossbar.<tailnet>:7777/v1/chat/completions
|
https://crossbar.<tailnet>:7777/v1/chat/completions
|
||||||
```
|
```
|
||||||
|
|
||||||
|
## Clients that manage their own slots
|
||||||
|
|
||||||
|
Some clients connect to one crossbar address and manage a llama-server slot themselves: they pin
|
||||||
|
`id_slot`, poll `/slots`, and steer a running completion through
|
||||||
|
`/v1/chat/completions/control`. Boxmaker's `inferproxy` is one. crossbar serves such a
|
||||||
|
client from a route that has its own `listen` address and `affinity = "route"`, so the whole route
|
||||||
|
lives on one host:
|
||||||
|
|
||||||
|
```toml
|
||||||
|
# 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
|
||||||
|
```
|
||||||
|
|
||||||
|
Every request to that address is this route, with its whole path passed upstream unchanged (there is
|
||||||
|
no route segment to strip), so it runs through `Handler.ForRoute` rather than the usual
|
||||||
|
`/{route}/` path. The address must split into a host and a numeric port, be unique across routes,
|
||||||
|
not equal the top-level `listen`, and not be on a template route — crossbar refuses any of those at
|
||||||
|
start-up.
|
||||||
|
|
||||||
|
A few things about how crossbar treats those requests:
|
||||||
|
|
||||||
|
- **Control calls take no slot.** A GET or HEAD on any allowed path, and a POST to exactly
|
||||||
|
`/tokenize` or `/v1/chat/completions/control`, is a control call. It follows the route's single
|
||||||
|
lease but takes no slot, skips the context guard, and writes no accounting row: it is sent beside
|
||||||
|
its own stream, so it must never wait for or hold a slot. A chat completion on `/v1/chat/completions`
|
||||||
|
is not a control call.
|
||||||
|
- **`/slots` and `/tokenize` are proxied; `/slots/<id>` actions are not.** Only the bare `/slots`
|
||||||
|
path is allowed, so an action on a specific slot id is not forwarded.
|
||||||
|
- **A GET's model comes from its `?model=` query** (there is no body to read), which is how
|
||||||
|
`/slots?model=shared` learns which model's slots to report.
|
||||||
|
- **The admin API is not served on a route listener.** `/_crossbar/hosts` there, and any prefixed
|
||||||
|
path such as `/boxmaker-a/v1/models`, are 404.
|
||||||
|
|
||||||
## Operate
|
## Operate
|
||||||
|
|
||||||
The operator's API lives under `/_crossbar/`. Every call returns 200 with a small JSON body unless
|
The operator's API lives under `/_crossbar/`. Every call returns 200 with a small JSON body unless
|
||||||
@@ -175,9 +218,14 @@ with the leased host's per-slot context for that model. A prompt that fits stays
|
|||||||
does not fit is moved to a healthy host on the route where it does fit (the lease moves with
|
does not fit is moved to a healthy host on the route where it does fit (the lease moves with
|
||||||
it, so the conversation stays there), and the response carries
|
it, so the conversation stays there), and the response carries
|
||||||
`X-Crossbar-Ctx: moved:<from>` + `>` + `<to>` — for example `moved:small>big`. When no host can
|
`X-Crossbar-Ctx: moved:<from>` + `>` + `<to>` — for example `moved:small>big`. When no host can
|
||||||
fit it, the answer is `400 {"error":"prompt too large","estimate":<tokens>,"max":<largest per-slot
|
fit it, the answer is a `400` in llama-server's own overflow shape, so a client that handles the
|
||||||
context among hosts that have the model loaded>}`. Hosts whose context is unknown are never
|
server's error handles crossbar's refusal too:
|
||||||
blocked by the guard.
|
|
||||||
|
```json
|
||||||
|
{"error":{"code":400,"type":"exceed_context_size_error","message":"prompt too large","n_prompt_tokens":<tokens>,"n_ctx":<largest per-slot context among hosts that have the model loaded>}}
|
||||||
|
```
|
||||||
|
|
||||||
|
Hosts whose context is unknown are never blocked by the guard.
|
||||||
|
|
||||||
## Wake
|
## Wake
|
||||||
|
|
||||||
@@ -204,4 +252,4 @@ tests only, so crossbar logs a warning when it starts in that mode.
|
|||||||
|
|
||||||
## What v2 does not do
|
## What v2 does not do
|
||||||
|
|
||||||
Request coalescing, `/slots` and TLS are out of scope for v2; see `PLAN.md`.
|
Request coalescing and TLS are out of scope for v2; see `PLAN.md`.
|
||||||
|
|||||||
+68
-8
@@ -11,6 +11,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"os/signal"
|
"os/signal"
|
||||||
|
"sort"
|
||||||
"syscall"
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -94,21 +95,55 @@ func run() error {
|
|||||||
p := proxy.New(cfg, table, leases, lim, st, log)
|
p := proxy.New(cfg, table, leases, lim, st, log)
|
||||||
p.SetWaker(waker)
|
p.SetWaker(waker)
|
||||||
|
|
||||||
var handler http.Handler = p
|
// Identity: gate the proxy on the route's peers when a backend is
|
||||||
|
// configured; off leaves the proxy unwrapped. The same checker gates each
|
||||||
|
// route's dedicated listener.
|
||||||
|
|
||||||
|
var checker *identity.Checker
|
||||||
if cfg.Identity != "off" {
|
if cfg.Identity != "off" {
|
||||||
var checker *identity.Checker
|
|
||||||
switch cfg.Identity {
|
switch cfg.Identity {
|
||||||
case "tailscale":
|
case "tailscale":
|
||||||
checker = identity.NewChecker(identity.TailscaleResolver{})
|
checker = identity.NewChecker(identity.TailscaleResolver{})
|
||||||
default: // "header"
|
default: // "header"
|
||||||
checker = identity.NewHeaderChecker()
|
checker = identity.NewHeaderChecker()
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
var handler http.Handler = p
|
||||||
|
if checker != nil {
|
||||||
handler = identity.Middleware(checker, func(route string) ([]string, bool) {
|
handler = identity.Middleware(checker, func(route string) ([]string, bool) {
|
||||||
rt, _, ok := cfg.Route(route)
|
rt, _, ok := cfg.Route(route)
|
||||||
return rt.Peers, ok
|
return rt.Peers, ok
|
||||||
}, p)
|
}, p)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Routes with a dedicated listener each serve their own address, with every request there being
|
||||||
|
// that route and the path unprefixed. In sorted route order, one server each, the proxy's
|
||||||
|
// ForRoute handler wrapped in RouteMiddleware when identity is on. No admin mux on them.
|
||||||
|
type routeServer struct {
|
||||||
|
route string
|
||||||
|
srv *http.Server
|
||||||
|
}
|
||||||
|
var routes []routeServer
|
||||||
|
listens := make([]string, 0, len(cfg.Routes))
|
||||||
|
for name := range cfg.Routes {
|
||||||
|
if cfg.Routes[name].Listen != "" {
|
||||||
|
listens = append(listens, name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
sort.Strings(listens)
|
||||||
|
for _, name := range listens {
|
||||||
|
rt := cfg.Routes[name]
|
||||||
|
var h http.Handler = p.ForRoute(name)
|
||||||
|
if checker != nil {
|
||||||
|
h = identity.RouteMiddleware(checker, rt.Peers, h)
|
||||||
|
}
|
||||||
|
routes = append(routes, routeServer{route: name, srv: &http.Server{
|
||||||
|
Addr: rt.Listen,
|
||||||
|
Handler: h,
|
||||||
|
ReadHeaderTimeout: 10 * time.Second,
|
||||||
|
}})
|
||||||
|
}
|
||||||
|
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
mux.Handle("/_crossbar/", admin.Handler(cfg, table, leases, lim, st, hosts))
|
mux.Handle("/_crossbar/", admin.Handler(cfg, table, leases, lim, st, hosts))
|
||||||
mux.Handle("/", handler)
|
mux.Handle("/", handler)
|
||||||
@@ -173,22 +208,47 @@ func run() error {
|
|||||||
ReadHeaderTimeout: 10 * time.Second,
|
ReadHeaderTimeout: 10 * time.Second,
|
||||||
}
|
}
|
||||||
|
|
||||||
serverErr := make(chan error, 1)
|
// Every listener shuts down together on ctx done; the first error other than a clean shutdown
|
||||||
go func() {
|
// ends run and shuts the rest down.
|
||||||
log.Info("listening", "addr", srv.Addr)
|
servers := make([]*http.Server, 0, 1+len(routes))
|
||||||
serverErr <- srv.ListenAndServe()
|
servers = append(servers, srv)
|
||||||
}()
|
for i := range routes {
|
||||||
|
servers = append(servers, routes[i].srv)
|
||||||
|
}
|
||||||
|
serverErr := make(chan error, len(servers))
|
||||||
|
start := func(s *http.Server, route string) {
|
||||||
|
go func() {
|
||||||
|
if route != "" {
|
||||||
|
log.Info("listening", "addr", s.Addr, "route", route)
|
||||||
|
} else {
|
||||||
|
log.Info("listening", "addr", s.Addr)
|
||||||
|
}
|
||||||
|
serverErr <- s.ListenAndServe()
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
start(srv, "")
|
||||||
|
for _, rs := range routes {
|
||||||
|
start(rs.srv, rs.route)
|
||||||
|
}
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
log.Info("shutting down")
|
log.Info("shutting down")
|
||||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
return srv.Shutdown(shutdownCtx)
|
for _, s := range servers {
|
||||||
|
_ = s.Shutdown(shutdownCtx)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
case err := <-serverErr:
|
case err := <-serverErr:
|
||||||
if errors.Is(err, http.ErrServerClosed) {
|
if errors.Is(err, http.ErrServerClosed) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
for _, s := range servers {
|
||||||
|
_ = s.Shutdown(shutdownCtx)
|
||||||
|
}
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,6 +5,10 @@ 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 |
|
||||||
|---|---|---|---|---|---|---|---|
|
|---|---|---|---|---|---|---|---|
|
||||||
|
| v2.3/04-ctx-error-docs | 2026-09-25 | done | 1 | pass | none | Resumed after the owner's v2.3 replacement `ctxguard_router_test.go` landed (byte-identical to the plan copy), resolving the earlier conflict with the protected v2.1 test. `refuseCtx` in `internal/proxy/ctxguard.go` answered the rule-4 400 in llama-server's own overflow shape `{"error":{"code":400,"type":"exceed_context_size_error","message":"prompt too large","n_prompt_tokens":<estimate>,"n_ctx":<largest per-slot context>}}`, the accounting row unchanged (status 400, Err "prompt too large"), every other error keeping `{"error":"<text>"}`; the new test reads `error.n_ctx` instead of the old top-level `max`, and `go test ./internal/proxy/` passes. README verified against task rule 2: the context-guard section documents the new body, the "Clients that manage their own slots" section covers control calls (follow the lease, take no slot, skip the guard, write no row), `/slots`+`/tokenize` proxied with `/slots/<id>` not, a GET's model from `?model=`, the `affinity`/`queue`/`listen` route keys with the `boxmaker-a` example, `listen` refused on templates and as the main address, no admin API on a route listener, and the config table gained the three keys. The "What v2 does not do" line no longer lists `/slots`, now that v2.3 proxies it. `make gate` → `gate: ok` first run, `make smoke` → `smoke: ok (stream spread 1008 ms)`. | ? |
|
||||||
|
| v2.3/03-route-listeners | 2026-09-25 | done | 1 | pass | `cmd/crossbar/main.go` refactors the identity build so one `*identity.Checker` (nil when off) gates both the main proxy and every route's `RouteMiddleware` (task said "wrapped in RouteMiddleware when identity on"; the checker had to be shared, not rebuilt per server). `proxy.go` gains a shared `serve()` flow that both `ServeHTTP` and `ForRoute` converge on, so lease keying is identical whether a request hits the main proxy or a dedicated listener (required by `TestForRouteServesUnprefixedPaths` which asserts bm-a/bm-b share one bm lease). | Implemented `internal/config/route.go`: `Route.Listen` (`toml:"listen"`), `checkListen` validating in the order the task lists it — numeric port 1–65535, not on the template, not equal to the main listen, unique across routes (a second route in sorted-name order reports the clash with the earlier route's name). `internal/proxy/proxy.go`: `ForRoute(name)` returns 404 for an unknown route, 400 for a conflicting `X-Crossbar-Route`, 404 for any prefixed/admin/root path (so a dedicated listener never serves another route), else the shared serve with the path unprefixed. `internal/identity/middleware.go`: `RouteMiddleware` (fixed peers, no admin-path exemption, empty peers lets all through). `main.go`: per-route servers in sorted route order, shared shutdown on ctx done, first non-`ErrServerClosed` error ends run. All five given/protected files byte-identical; `make gate` → `gate: ok` first run, `make smoke` → `smoke: ok (stream spread 1004 ms)`. | ? |
|
||||||
|
| v2.3/02-affinity-queue | 2026-09-25 | done | 1 | pass | `internal/proxy/proxy.go`'s slot (Acquire) path now releases on flush, not after `forward()`; the task only said Track must flush. | Implemented `internal/config/route.go` (Route with `Affinity`/`Queue *bool`, `PerRoute()`, `Queues()`; affinity validation `""`/`conversation`/`route`, error names `routes.<name>.affinity`; moved `checkRoutes`/`routeName`). `config.go`: one-line call to `checkRoutes`. `internal/limiter/limiter.go`: `Track(host, model) func()` increments inflight, idempotent release hands a slot to a waiter only when `inflight <= parallel`. `proxy.go`: `leaseFP = ""` in the lease key when `routeCfg.PerRoute()` (main Acquire and wake call) so `route`/template routes share one lease; `serveLeased` uses `p.lim.Track` when `routeCfg.Queues()` is false, else `Acquire`. `forward.go`: `forward()` gained a `release func()` param; `statusRecorder.onFlush` field with `Flush()` calling `onFlush()` before the underlying flush. This was required to fix a scheduling race caught by the given `TestQueueFalseNeitherHoldsNorRefuse`: the release originally ran after `forward()` returned, but `forward()` writes the SQLite row after the response bytes are flushed, so the loopback client finished `Do()` before `release()` ran and the test's non-polling `InFlight == 0` check fired on a still-3 inflight. Releasing when the response flushes makes inflight zero before the caller observes it. Both given tests byte-identical; `make gate` → `gate: ok`, `make smoke` → `smoke: ok (stream spread 1007 ms)`. | ? **Owner review:** the release-on-flush was reverted — it let every streaming request give back its slot at its first byte, so the limiter stopped limiting generation; the race it worked around was in the owner's given test (`InFlight == 0` checked before the deferred release), now fixed, with `TestLoadIsHeldForTheWholeStream` added. Session ended on a refused `/tmp` write while committing; owner committed. |
|
||||||
|
| v2.3/01-control-plane | 2026-09-25 | done | 1 | pass | none | New `internal/proxy/control.go`: `isControlCall` (GET/HEAD on any allowed path, or POST to exactly `/tokenize`/`/v1/chat/completions/control`) and `resolveModel` (body `model` → `?model=` → route `default_model`). `proxy.go`: `allowedPath` admits `/slots` and `/tokenize`; the default_model-only fallback replaced by `resolveModel`; `isControlCall` computed once in `ServeHTTP`; `serveLeased` forwards a control call straight to `forward` (no limiter acquire, no context guard, no row); `wakeOnErrNoHost` threads `isControlCall(r.Method, rest)` through. `forward.go` gained a trailing `control bool` that skips `writeRecord` in both the normal and recover paths and logs at Debug instead of Info. Both given tests byte-identical; `make gate` → `gate: ok` on the first run. | ? |
|
||||||
| v2.2/02-broadcasts | 2026-09-25 | done | 1 | pass | The Wake struct and checkWake live in `internal/config/identity.go` (added in task 04), not `config.go`, so I edited `identity.go` rather than `config.go`; `wake.go` logs a broadcast that fails to resolve/send before continuing (task rule 2 allows "logged or ignored"). | Added `Broadcasts` to `Wake` and `Wake.Addresses()` (Broadcast then Broadcasts, never empty for a parsed config); `checkWake` errors on both-set → `.broadcasts`, neither-or-empty-list → `.broadcast`, and a non-`host:port` entry → `.broadcasts`; `Target` gains `Broadcasts` and `Wake`/`sendAll` send to Broadcast then each Broadcasts in order, logging/past a failure and returning false only when no address could be sent; `main.go` fills `Target.Broadcasts` from `Wake.Addresses()` and leaves `Target.Broadcast` empty so `sendAll` does not double-send. Given `broadcasts_test.go` and `config_v22_test.go` byte-identical, v2 `wake_test.go`/`config_v2_test.go` untouched and green; `make gate` → `gate: ok` first run. | ? |
|
| v2.2/02-broadcasts | 2026-09-25 | done | 1 | pass | The Wake struct and checkWake live in `internal/config/identity.go` (added in task 04), not `config.go`, so I edited `identity.go` rather than `config.go`; `wake.go` logs a broadcast that fails to resolve/send before continuing (task rule 2 allows "logged or ignored"). | Added `Broadcasts` to `Wake` and `Wake.Addresses()` (Broadcast then Broadcasts, never empty for a parsed config); `checkWake` errors on both-set → `.broadcasts`, neither-or-empty-list → `.broadcast`, and a non-`host:port` entry → `.broadcasts`; `Target` gains `Broadcasts` and `Wake`/`sendAll` send to Broadcast then each Broadcasts in order, logging/past a failure and returning false only when no address could be sent; `main.go` fills `Target.Broadcasts` from `Wake.Addresses()` and leaves `Target.Broadcast` empty so `sendAll` does not double-send. Given `broadcasts_test.go` and `config_v22_test.go` byte-identical, v2 `wake_test.go`/`config_v2_test.go` untouched and green; `make gate` → `gate: ok` first run. | ? |
|
||||||
| v2.2/01-route-templates | 2026-09-25 | done | 1 | fail | `internal/config` red only on `Wake.Addresses()` (task 02), the one allowed red; `go build ./...` clean, proxy/admin/health/wake/lease/store/identity/fingerprint all pass under `-race`. New `internal/config/route.go`: `templateName` pattern `^[a-z0-9][a-z0-9-]*-\*$` and `Route()` (valid-name guard excludes `*`; exact wins; else longest `"<prefix>-*"`, prefix keeps the dash, non-empty remainder required, longest-prefix wins deterministically). `config.go`: the route-name check accepts a template too (one line). `proxy.go`: `route()` and `ServeHTTP` resolve both path and `X-Crossbar-Route` header forms through `cfg.Route`, and the conflicting-route check compares concrete names via `cfg.Route` (identical to before for non-template configs). `admin.go` `routeView` lists a lease under the exact key it matches or the longest template key; `admin_ops.go` `routePin` resolves through `cfg.Route` so a concrete route under a template can be pinned before its first request and the template name 404s. `main.go` identity lookup uses `cfg.Route`. All three given tests byte-identical (`config_v22_test.go` keeps `TestWakeBroadcasts`, which is why config is red). | ? |
|
| v2.2/01-route-templates | 2026-09-25 | done | 1 | fail | `internal/config` red only on `Wake.Addresses()` (task 02), the one allowed red; `go build ./...` clean, proxy/admin/health/wake/lease/store/identity/fingerprint all pass under `-race`. New `internal/config/route.go`: `templateName` pattern `^[a-z0-9][a-z0-9-]*-\*$` and `Route()` (valid-name guard excludes `*`; exact wins; else longest `"<prefix>-*"`, prefix keeps the dash, non-empty remainder required, longest-prefix wins deterministically). `config.go`: the route-name check accepts a template too (one line). `proxy.go`: `route()` and `ServeHTTP` resolve both path and `X-Crossbar-Route` header forms through `cfg.Route`, and the conflicting-route check compares concrete names via `cfg.Route` (identical to before for non-template configs). `admin.go` `routeView` lists a lease under the exact key it matches or the longest template key; `admin_ops.go` `routePin` resolves through `cfg.Route` so a concrete route under a template can be pinned before its first request and the template name 404s. `main.go` identity lookup uses `cfg.Route`. All three given tests byte-identical (`config_v22_test.go` keeps `TestWakeBroadcasts`, which is why config is red). | ? |
|
||||||
| v2.1/02-props-loaded-only | 2026-09-25 | done | 1 | pass | `movedHeader` separator `><`→`>` and the `CtxHeader` doc comment in `proxy.go`, both forced by the given router test (`moved:small>big`) which the task text did not mention; no production code parses the separator (`forward.go` passes it straight through) so it is safe. | Implemented per-model context. `health`: added `ModelCtx` and a `Models map[string]ModelCtx` field on `Status`, plus `PerSlotCtxFor(model)` (per-model figure when present, else host-level `PerSlotCtx` for a loaded model, else 0); moved `props` into a new `props.go` and added `propsModel`/`propsModels`. Poller rules 1-4: `/v1/models` treats an entry as loaded only with no `status` or `status.value=="loaded"` (other values dropped from `Loaded`); plain `/props` with `role:router` leaves host NCtx/Slots 0; each loaded model is asked `GET /props?model=<url.QueryEscape(id)>` and a failed/malformed answer leaves that id absent without failing the host; `Models` is a fresh non-nil map every successful poll, `MarkDown` leaves it. `ctxguard.go`: every `PerSlotCtx()` became `PerSlotCtxFor(model)` (leased host, candidates, wake "cannot serve" check) and `largestSlotCtx(hosts,h,model)` counts only hosts that have it loaded. `admin.go`: `HostView` gains `models` (empty object, never null). A plain single server keeps working as v2. All three given tests byte-identical; `make gate` → `gate: ok` on the first run. | ? |
|
| v2.1/02-props-loaded-only | 2026-09-25 | done | 1 | pass | `movedHeader` separator `><`→`>` and the `CtxHeader` doc comment in `proxy.go`, both forced by the given router test (`moved:small>big`) which the task text did not mention; no production code parses the separator (`forward.go` passes it straight through) so it is safe. | Implemented per-model context. `health`: added `ModelCtx` and a `Models map[string]ModelCtx` field on `Status`, plus `PerSlotCtxFor(model)` (per-model figure when present, else host-level `PerSlotCtx` for a loaded model, else 0); moved `props` into a new `props.go` and added `propsModel`/`propsModels`. Poller rules 1-4: `/v1/models` treats an entry as loaded only with no `status` or `status.value=="loaded"` (other values dropped from `Loaded`); plain `/props` with `role:router` leaves host NCtx/Slots 0; each loaded model is asked `GET /props?model=<url.QueryEscape(id)>` and a failed/malformed answer leaves that id absent without failing the host; `Models` is a fresh non-nil map every successful poll, `MarkDown` leaves it. `ctxguard.go`: every `PerSlotCtx()` became `PerSlotCtxFor(model)` (leased host, candidates, wake "cannot serve" check) and `largestSlotCtx(hosts,h,model)` counts only hosts that have it loaded. `admin.go`: `HostView` gains `models` (empty object, never null). A plain single server keeps working as v2. All three given tests byte-identical; `make gate` → `gate: ok` on the first run. | ? |
|
||||||
@@ -108,3 +112,14 @@ it; (b) task — the task text did not state that `httputil.ReverseProxy` aborts
|
|||||||
test design — timing-based tests (limiter, queue, spread, cancel) have margins tuned for an idle
|
test design — timing-based tests (limiter, queue, spread, cancel) have margins tuned for an idle
|
||||||
host; widen or retry in a later plan.
|
host; widen or retry in a later plan.
|
||||||
|
|
||||||
|
|
||||||
|
### v2.3 review (owner, 2026-09-25)
|
||||||
|
|
||||||
|
Checked: gate, `-race -count=3` on proxy and limiter, smoke (check 6: dedicated listener). Task 01
|
||||||
|
clean (nit: the Debug log block copies the Info block's fields). Task 02: release-on-first-flush
|
||||||
|
reverted by the owner (limiter stopped limiting streams; cause was the owner's racy given test,
|
||||||
|
now fixed, with `TestLoadIsHeldForTheWholeStream`). Task 03: correct; `main.go` called
|
||||||
|
`logIdentityMode` twice (removed). Task 04: stopped correctly on the owner's missed v2.1 router
|
||||||
|
test; resumed after the replacement. README: "Boxmaker's router" → "Boxmaker's `inferproxy`".
|
||||||
|
Model faults this plan: one timing hack (logged as a deviation), one refusal-ending, one
|
||||||
|
malformed tool call ending a session with no change. Owner faults: racy test, missed router test.
|
||||||
|
|||||||
@@ -31,3 +31,13 @@ default_model = "ornith-1.5-35b-a3b"
|
|||||||
[routes.hermes-x]
|
[routes.hermes-x]
|
||||||
hosts = ["beta", "alpha"]
|
hosts = ["beta", "alpha"]
|
||||||
# peers = ["talos"] # v2: with identity = "tailscale", only these tailnet nodes may use the route
|
# 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
|
||||||
|
|||||||
@@ -71,14 +71,6 @@ type Host struct {
|
|||||||
Wake *Wake `toml:"wake"`
|
Wake *Wake `toml:"wake"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Route is an ordered list of hosts to try, with an optional default model and
|
|
||||||
// the peers allowed to reach it.
|
|
||||||
type Route struct {
|
|
||||||
Hosts []string `toml:"hosts"`
|
|
||||||
DefaultModel string `toml:"default_model"`
|
|
||||||
Peers []string `toml:"peers"`
|
|
||||||
}
|
|
||||||
|
|
||||||
// Config is the whole file: what to listen on, tuning, hosts and routes.
|
// Config is the whole file: what to listen on, tuning, hosts and routes.
|
||||||
type Config struct {
|
type Config struct {
|
||||||
Listen string `toml:"listen"`
|
Listen string `toml:"listen"`
|
||||||
@@ -117,8 +109,6 @@ const (
|
|||||||
DefaultIdentity = "off"
|
DefaultIdentity = "off"
|
||||||
)
|
)
|
||||||
|
|
||||||
var routeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`)
|
|
||||||
|
|
||||||
// Load reads and parses the config file at path. An open failure is wrapped as
|
// Load reads and parses the config file at path. An open failure is wrapped as
|
||||||
// "config: …", the same shape as a decode failure.
|
// "config: …", the same shape as a decode failure.
|
||||||
func Load(path string) (*Config, error) {
|
func Load(path string) (*Config, error) {
|
||||||
@@ -345,54 +335,3 @@ func (c *Config) checkHosts() *Error {
|
|||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) *Error {
|
|
||||||
if len(c.Routes) == 0 {
|
|
||||||
return &Error{Field: "routes", Msg: "at least one required"}
|
|
||||||
}
|
|
||||||
names := make([]string, 0, len(c.Routes))
|
|
||||||
for name := range c.Routes {
|
|
||||||
names = append(names, name)
|
|
||||||
}
|
|
||||||
sort.Strings(names)
|
|
||||||
for _, name := range names {
|
|
||||||
r := c.Routes[name]
|
|
||||||
|
|
||||||
if !routeName.MatchString(name) && !templateName.MatchString(name) {
|
|
||||||
return &Error{Field: fmt.Sprintf("routes.%s", name), Msg: "must match [a-z0-9][a-z0-9-]*"}
|
|
||||||
}
|
|
||||||
|
|
||||||
hostsField := fmt.Sprintf("routes.%s.hosts", name)
|
|
||||||
if len(r.Hosts) == 0 {
|
|
||||||
return &Error{Field: hostsField, Msg: "at least one required"}
|
|
||||||
}
|
|
||||||
seen := make(map[string]bool, len(r.Hosts))
|
|
||||||
for _, h := range r.Hosts {
|
|
||||||
if seen[h] {
|
|
||||||
return &Error{Field: hostsField, Msg: "host listed twice"}
|
|
||||||
}
|
|
||||||
seen[h] = true
|
|
||||||
if _, ok := c.Hosts[h]; !ok {
|
|
||||||
return &Error{Field: hostsField, Msg: "unknown host"}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if r.DefaultModel != "" {
|
|
||||||
served := false
|
|
||||||
for _, h := range r.Hosts {
|
|
||||||
if _, ok := c.Hosts[h].Models[r.DefaultModel]; ok {
|
|
||||||
served = true
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if !served {
|
|
||||||
return &Error{Field: fmt.Sprintf("routes.%s.default_model", name), Msg: "not served by any host in route"}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if e := checkPeers(name, r.Peers, peersDefined[name], identityDefined, c.Identity); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -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`)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,13 +1,44 @@
|
|||||||
package config
|
package config
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
|
"net"
|
||||||
"regexp"
|
"regexp"
|
||||||
|
"sort"
|
||||||
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
)
|
)
|
||||||
|
|
||||||
// templateName matches a route template: a valid route name ending in "-*".
|
// templateName matches a route template: a valid route name ending in "-*".
|
||||||
var templateName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*-\*$`)
|
var templateName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*-\*$`)
|
||||||
|
|
||||||
|
// routeName matches a route (or template) name: the pattern a concrete or template route key must
|
||||||
|
// match, so a name with '*' or an invalid prefix never resolves.
|
||||||
|
var routeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`)
|
||||||
|
|
||||||
|
// Route is an ordered list of hosts to try, with an optional default model, the peers allowed to
|
||||||
|
// reach it, how its requests are placed (affinity), and whether crossbar queues them.
|
||||||
|
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
|
||||||
|
Listen string `toml:"listen"` // "" = none; else host:port of the route's own listener
|
||||||
|
}
|
||||||
|
|
||||||
|
// PerRoute reports affinity = "route": every request on the route (chat or control) shares one
|
||||||
|
// lease per model, so the route lives on one host.
|
||||||
|
func (r Route) PerRoute() bool {
|
||||||
|
return r.Affinity == "route"
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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, which leaves queueing to the client's own slot.
|
||||||
|
func (r Route) Queues() bool {
|
||||||
|
return r.Queue == nil || *r.Queue
|
||||||
|
}
|
||||||
|
|
||||||
// Route resolves a request route name: an exact entry wins; else the longest template
|
// Route resolves a request route name: an exact entry wins; else the longest template
|
||||||
// "<prefix>-*" whose prefix (including the dash) starts name with a non-empty remainder;
|
// "<prefix>-*" whose prefix (including the dash) starts name with a non-empty remainder;
|
||||||
// else ok is false. key is the config key that matched (the template's name for a template).
|
// else ok is false. key is the config key that matched (the template's name for a template).
|
||||||
@@ -38,3 +69,100 @@ func (c *Config) Route(name string) (r Route, key string, ok bool) {
|
|||||||
}
|
}
|
||||||
return best, bestKey, true
|
return best, bestKey, true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// checkRoutes validates and defaults one route's hosts, model, affinity and peers in a fixed order.
|
||||||
|
// A name that is neither a valid route nor a template, a missing or unknown host, a default model no
|
||||||
|
// host serves, an unrecognised affinity, or a peers list that breaks the identity contract each
|
||||||
|
// wins as the first error.
|
||||||
|
func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) *Error {
|
||||||
|
if len(c.Routes) == 0 {
|
||||||
|
return &Error{Field: "routes", Msg: "at least one required"}
|
||||||
|
}
|
||||||
|
names := make([]string, 0, len(c.Routes))
|
||||||
|
for name := range c.Routes {
|
||||||
|
names = append(names, name)
|
||||||
|
}
|
||||||
|
sort.Strings(names)
|
||||||
|
// A route's own listener address, keyed for the uniqueness check: the value is the
|
||||||
|
// route that first claimed it, so the second route in sorted order reports the miss.
|
||||||
|
seenListen := make(map[string]string, len(c.Routes))
|
||||||
|
for _, name := range names {
|
||||||
|
r := c.Routes[name]
|
||||||
|
|
||||||
|
if !routeName.MatchString(name) && !templateName.MatchString(name) {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s", name), Msg: "must match [a-z0-9][a-z0-9-]*"}
|
||||||
|
}
|
||||||
|
|
||||||
|
hostsField := fmt.Sprintf("routes.%s.hosts", name)
|
||||||
|
if len(r.Hosts) == 0 {
|
||||||
|
return &Error{Field: hostsField, Msg: "at least one required"}
|
||||||
|
}
|
||||||
|
seen := make(map[string]bool, len(r.Hosts))
|
||||||
|
for _, h := range r.Hosts {
|
||||||
|
if seen[h] {
|
||||||
|
return &Error{Field: hostsField, Msg: "host listed twice"}
|
||||||
|
}
|
||||||
|
seen[h] = true
|
||||||
|
if _, ok := c.Hosts[h]; !ok {
|
||||||
|
return &Error{Field: hostsField, Msg: "unknown host"}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if r.DefaultModel != "" {
|
||||||
|
served := false
|
||||||
|
for _, h := range r.Hosts {
|
||||||
|
if _, ok := c.Hosts[h].Models[r.DefaultModel]; ok {
|
||||||
|
served = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !served {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.default_model", name), Msg: "not served by any host in route"}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
switch r.Affinity {
|
||||||
|
case "", "conversation", "route":
|
||||||
|
default:
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.affinity", name), Msg: `must be "conversation" or "route"`}
|
||||||
|
}
|
||||||
|
|
||||||
|
if e := checkPeers(name, r.Peers, peersDefined[name], identityDefined, c.Identity); e != nil {
|
||||||
|
return e
|
||||||
|
}
|
||||||
|
|
||||||
|
if e := checkListen(name, r.Listen, c.Listen, seenListen); e != nil {
|
||||||
|
return e
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// checkListen validates one route's own listener. The rules, in order, each naming the key
|
||||||
|
// routes.<name>.listen: the value must split into host and a numeric port 1-65535, it must not be on
|
||||||
|
// a template route, it must not be the top-level listen, and it must be unique across routes (the
|
||||||
|
// second route in sorted name order reports the clash and the other route's name).
|
||||||
|
func checkListen(name, listen, mainListen string, seenListen map[string]string) *Error {
|
||||||
|
if listen == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
_, port, err := net.SplitHostPort(listen)
|
||||||
|
if err != nil {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "must be host:port"}
|
||||||
|
}
|
||||||
|
n, err := strconv.Atoi(port)
|
||||||
|
if err != nil || n < 1 || n > 65535 {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "port must be 1-65535"}
|
||||||
|
}
|
||||||
|
if templateName.MatchString(name) {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "a template route cannot have its own listener"}
|
||||||
|
}
|
||||||
|
if listen == mainListen {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "cannot be the main listen address"}
|
||||||
|
}
|
||||||
|
if other, dup := seenListen[listen]; dup {
|
||||||
|
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: fmt.Sprintf("already used by route %s", other)}
|
||||||
|
}
|
||||||
|
seenListen[listen] = name
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -45,6 +45,21 @@ func Middleware(c *Checker, peersFor func(route string) ([]string, bool), next h
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RouteMiddleware gates every request on a fixed set of peers, whatever path or X-Crossbar-Route
|
||||||
|
// header it carries: on a route's dedicated listener the route is fixed, so the gate is that route's
|
||||||
|
// peers for every request. There is no admin-path exemption — there is no admin on a route listener.
|
||||||
|
// Empty peers lets everyone through, as for Middleware.
|
||||||
|
func RouteMiddleware(c *Checker, peers []string, next http.Handler) http.Handler {
|
||||||
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
ctx := WithHeaderPeer(r.Context(), r.Header.Get(peerHeader))
|
||||||
|
if err := c.Allow(ctx, peers, r.RemoteAddr); err != nil {
|
||||||
|
writeForbidden(w)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
next.ServeHTTP(w, r)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// writeForbidden answers the JSON 403 the tests and callers expect.
|
// writeForbidden answers the JSON 403 the tests and callers expect.
|
||||||
func writeForbidden(w http.ResponseWriter) {
|
func writeForbidden(w http.ResponseWriter) {
|
||||||
w.Header().Set("Content-Type", "application/json")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
|||||||
@@ -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())
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -108,16 +108,27 @@ func (l *Limiter) Acquire(ctx context.Context, host, model string) (release func
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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) func() {
|
||||||
|
l.mu.Lock()
|
||||||
|
p := l.pairLocked(host, model)
|
||||||
|
p.inflight++
|
||||||
|
l.mu.Unlock()
|
||||||
|
return l.release(p)
|
||||||
|
}
|
||||||
|
|
||||||
// release returns the function the caller holds for a slot: it hands the slot to the next waiter
|
// release returns the function the caller holds for a slot: it hands the slot to the next waiter
|
||||||
// if one is waiting, otherwise it frees the slot. It is safe to call through the sync.Once that
|
// only while there is room (in flight at or below parallel), otherwise it counts the slot back. It is
|
||||||
// Acquire wrapped it in.
|
// safe to call through the sync.Once that Acquire wrapped it in. The same release serves Track, whose
|
||||||
|
// tracked load can push in flight past parallel, so a release there cannot free a slot that exists.
|
||||||
func (l *Limiter) release(p *pair) func() {
|
func (l *Limiter) release(p *pair) func() {
|
||||||
var once sync.Once
|
var once sync.Once
|
||||||
return func() {
|
return func() {
|
||||||
once.Do(func() {
|
once.Do(func() {
|
||||||
l.mu.Lock()
|
l.mu.Lock()
|
||||||
defer l.mu.Unlock()
|
defer l.mu.Unlock()
|
||||||
if len(p.waiters) > 0 {
|
if len(p.waiters) > 0 && p.inflight <= p.parallel {
|
||||||
next := p.waiters[0]
|
next := p.waiters[0]
|
||||||
p.waiters = p.waiters[1:]
|
p.waiters = p.waiters[1:]
|
||||||
close(next)
|
close(next)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,180 @@
|
|||||||
|
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 (
|
||||||
|
"bufio"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"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)
|
||||||
|
}
|
||||||
|
// The slot is given back just after the answer is sent (a deferred release), so wait for it.
|
||||||
|
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
|
||||||
|
// Accounting is unchanged: each chat is still a row.
|
||||||
|
waitUntil(t, func() bool { return r.rows("bm") == 3 })
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A request counts against its host for as long as its answer is streaming, not only until the
|
||||||
|
// first byte: a slot (queueing route) or a tracked place (queue = false) is given back when the
|
||||||
|
// stream ends.
|
||||||
|
func TestLoadIsHeldForTheWholeStream(t *testing.T) {
|
||||||
|
for _, route := range []string{"r", "bm"} {
|
||||||
|
t.Run(route, func(t *testing.T) {
|
||||||
|
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||||
|
r := newRig(t, affinityHosts, alpha, beta)
|
||||||
|
body := strings.Replace(conversation(1, 1), `"stream":false`, `"stream":true`, 1)
|
||||||
|
resp := r.do(http.MethodPost, "/"+route+"/v1/chat/completions", body)
|
||||||
|
defer resp.Body.Close()
|
||||||
|
host := resp.Header.Get("X-Crossbar-Host")
|
||||||
|
line, err := bufio.NewReader(resp.Body).ReadString('\n')
|
||||||
|
if err != nil || !strings.HasPrefix(line, "data:") {
|
||||||
|
t.Fatalf("first line %q, err %v", line, err)
|
||||||
|
}
|
||||||
|
// The first chunk is here; the upstream sends more for another ~30 ms.
|
||||||
|
if n := r.lim.InFlight(host, "shared"); n != 1 {
|
||||||
|
t.Errorf("in flight on %s after the first chunk = %d, want 1 (released before the stream ended)", host, n)
|
||||||
|
}
|
||||||
|
drain(resp)
|
||||||
|
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,31 @@
|
|||||||
|
package proxy
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
|
||||||
|
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||||
|
)
|
||||||
|
|
||||||
|
// isControlCall reports whether r is a control-plane call: a GET or HEAD on any allowed path, or a
|
||||||
|
// POST to exactly /tokenize or /v1/chat/completions/control. Everything else — in particular a chat
|
||||||
|
// completion on /v1/chat/completions — is not a control call. rest is the upstream path, not the
|
||||||
|
// query string.
|
||||||
|
func isControlCall(method, rest string) bool {
|
||||||
|
if method == http.MethodGet || method == http.MethodHead {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return method == http.MethodPost && (rest == "/tokenize" || rest == "/v1/chat/completions/control")
|
||||||
|
}
|
||||||
|
|
||||||
|
// resolveModel returns the model for a request: the body's top-level "model" wins (already read by
|
||||||
|
// peekModel), else the ?model= query parameter, else the route's default_model. The result keys the
|
||||||
|
// lease and the limiter pair.
|
||||||
|
func resolveModel(model string, r *http.Request, routeCfg config.Route) string {
|
||||||
|
if model == "" {
|
||||||
|
model = r.URL.Query().Get("model")
|
||||||
|
}
|
||||||
|
if model == "" {
|
||||||
|
model = routeCfg.DefaultModel
|
||||||
|
}
|
||||||
|
return model
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -155,10 +155,21 @@ func largestSlotCtx(hosts []string, h Health, model string) int {
|
|||||||
return best
|
return best
|
||||||
}
|
}
|
||||||
|
|
||||||
// refuseCtx answers the 400 the guard's rule 4: the JSON body carries the
|
// ctxErrorBody is llama-server's shape for a prompt that exceeds a host's
|
||||||
// estimate and the largest available per-slot context, plus the error text. It
|
// context. A client that already handles the server's own overflow error keys
|
||||||
// records the accounting row (status 400, Err "prompt too large") and never
|
// on error.type and so recognises crossbar's refusal too. See refuseCtx.
|
||||||
// marks the host down.
|
type ctxErrorBody struct {
|
||||||
|
Code int `json:"code"`
|
||||||
|
Type string `json:"type"`
|
||||||
|
Message string `json:"message"`
|
||||||
|
NPromptTokens int `json:"n_prompt_tokens"`
|
||||||
|
NCtx int `json:"n_ctx"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// refuseCtx answers the 400 the guard's rule 4: the body is llama-server's
|
||||||
|
// exceed_context_size_error, with the estimate as n_prompt_tokens and the
|
||||||
|
// largest available per-slot context as n_ctx. It records the accounting row
|
||||||
|
// (status 400, Err "prompt too large") and never marks the host down.
|
||||||
func (p *Handler) refuseCtx(w http.ResponseWriter, host, route, model, fp string, started time.Time, estimate, maxSlot int) {
|
func (p *Handler) refuseCtx(w http.ResponseWriter, host, route, model, fp string, started time.Time, estimate, maxSlot int) {
|
||||||
p.writeRecord(store.Request{
|
p.writeRecord(store.Request{
|
||||||
Route: route,
|
Route: route,
|
||||||
@@ -173,9 +184,13 @@ func (p *Handler) refuseCtx(w http.ResponseWriter, host, route, model, fp string
|
|||||||
w.Header().Set("Content-Type", "application/json")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
w.WriteHeader(http.StatusBadRequest)
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
_ = json.NewEncoder(w).Encode(map[string]any{
|
_ = json.NewEncoder(w).Encode(map[string]any{
|
||||||
"error": "prompt too large",
|
"error": ctxErrorBody{
|
||||||
"estimate": estimate,
|
Code: http.StatusBadRequest,
|
||||||
"max": maxSlot,
|
Type: "exceed_context_size_error",
|
||||||
|
Message: "prompt too large",
|
||||||
|
NPromptTokens: estimate,
|
||||||
|
NCtx: maxSlot,
|
||||||
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -95,11 +95,16 @@ func TestRouterUnloadedModelIsNotACandidate(t *testing.T) {
|
|||||||
if big.hits.Load() != 0 {
|
if big.hits.Load() != 0 {
|
||||||
t.Errorf("big served %d requests for a model it does not have loaded", big.hits.Load())
|
t.Errorf("big served %d requests for a model it does not have loaded", big.hits.Load())
|
||||||
}
|
}
|
||||||
var e map[string]any
|
// v2.3: the refusal is llama-server's exceed_context_size_error shape; n_ctx is what "max" was.
|
||||||
|
var e struct {
|
||||||
|
Error struct {
|
||||||
|
NCtx float64 `json:"n_ctx"`
|
||||||
|
} `json:"error"`
|
||||||
|
}
|
||||||
if err := json.Unmarshal([]byte(body), &e); err != nil {
|
if err := json.Unmarshal([]byte(body), &e); err != nil {
|
||||||
t.Fatalf("body %q is not JSON: %v", body, err)
|
t.Fatalf("body %q is not JSON: %v", body, err)
|
||||||
}
|
}
|
||||||
if max, _ := e["max"].(float64); max != 4096 {
|
if e.Error.NCtx != 4096 {
|
||||||
t.Errorf("max = %v, want 4096: the largest per-slot context among hosts that have shared loaded", e["max"])
|
t.Errorf("error.n_ctx = %v, want 4096: the largest per-slot context among hosts that have shared loaded", e.Error.NCtx)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -84,15 +84,31 @@ func TestOversizedPromptWithNoFitIs400(t *testing.T) {
|
|||||||
if resp.StatusCode != http.StatusBadRequest {
|
if resp.StatusCode != http.StatusBadRequest {
|
||||||
t.Fatalf("status %d body %s, want 400", resp.StatusCode, body)
|
t.Fatalf("status %d body %s, want 400", resp.StatusCode, body)
|
||||||
}
|
}
|
||||||
var e map[string]any
|
// v2.3: llama-server's own shape for this error, so a client handles crossbar's refusal the
|
||||||
if err := json.Unmarshal([]byte(body), &e); err != nil || e["error"] != "prompt too large" {
|
// way it handles the server's (Boxmaker keys on error.type; the error JSON must come first).
|
||||||
t.Fatalf("body = %s, want error 'prompt too large'", body)
|
if !strings.HasPrefix(body, `{"error":`) {
|
||||||
|
t.Errorf("body must start with the error object: %s", body)
|
||||||
}
|
}
|
||||||
if est, _ := e["estimate"].(float64); est < 8000 || est > 13000 {
|
var e struct {
|
||||||
t.Errorf("estimate = %v, want roughly 10000 tokens", e["estimate"])
|
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 max, _ := e["max"].(float64); max != 4096 {
|
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.Errorf("max = %v, want the largest per-slot context among the route's hosts (4096)", e["max"])
|
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 {
|
if small.hits.Load()+tiny.hits.Load() != 0 {
|
||||||
t.Errorf("a refused prompt must not reach any upstream")
|
t.Errorf("a refused prompt must not reach any upstream")
|
||||||
|
|||||||
@@ -16,8 +16,9 @@ import (
|
|||||||
// forward builds the reverse proxy for one host, tees the response, records the accounting row, and
|
// forward builds the reverse proxy for one host, tees the response, records the accounting row, and
|
||||||
// logs. leaseState is "new" or "reused"; waited is the time spent in the queue. ctxEst is the
|
// logs. leaseState is "new" or "reused"; waited is the time spent in the queue. ctxEst is the
|
||||||
// prompt size the context guard estimated (0 when the guard did not run); ctxHeader is the
|
// prompt size the context guard estimated (0 when the guard did not run); ctxHeader is the
|
||||||
// "moved:…<host>" header to set when the guard relocated the conversation.
|
// "moved:…<host>" header to set when the guard relocated the conversation. control is true for a
|
||||||
func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, leaseState, rest, fp, model string, started time.Time, waited time.Duration, ctxEst int, ctxHeader string) {
|
// control-plane call: it takes no slot, so forward writes no row and logs at Debug for it.
|
||||||
|
func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, leaseState, rest, fp, model string, started time.Time, waited time.Duration, ctxEst int, ctxHeader string, control bool) {
|
||||||
hostCfg, ok := p.cfg.Hosts[host]
|
hostCfg, ok := p.cfg.Hosts[host]
|
||||||
if !ok {
|
if !ok {
|
||||||
p.writeError(w, http.StatusBadGateway, "upstream failed")
|
p.writeError(w, http.StatusBadGateway, "upstream failed")
|
||||||
@@ -45,7 +46,9 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
|||||||
} else {
|
} else {
|
||||||
req.Err = "upstream error"
|
req.Err = "upstream error"
|
||||||
}
|
}
|
||||||
p.writeRecord(req)
|
if !control {
|
||||||
|
p.writeRecord(req)
|
||||||
|
}
|
||||||
panic(pv)
|
panic(pv)
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
@@ -58,12 +61,29 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
|||||||
req.Status = 499
|
req.Status = 499
|
||||||
req.Err = "client cancelled"
|
req.Err = "client cancelled"
|
||||||
}
|
}
|
||||||
p.writeRecord(req)
|
if !control {
|
||||||
|
p.writeRecord(req)
|
||||||
|
}
|
||||||
|
|
||||||
fp8 := fp
|
fp8 := fp
|
||||||
if len(fp8) > 8 {
|
if len(fp8) > 8 {
|
||||||
fp8 = fp8[:8]
|
fp8 = fp8[:8]
|
||||||
}
|
}
|
||||||
|
if control {
|
||||||
|
p.log.Debug("request",
|
||||||
|
"route", route,
|
||||||
|
"host", host,
|
||||||
|
"method", r.Method,
|
||||||
|
"path", rest,
|
||||||
|
"status", rec.status,
|
||||||
|
"lease", leaseState,
|
||||||
|
"queued_ms", waited.Milliseconds(),
|
||||||
|
"fp", fp8,
|
||||||
|
"ctx_est", ctxEst,
|
||||||
|
"ms", total.Milliseconds(),
|
||||||
|
)
|
||||||
|
return
|
||||||
|
}
|
||||||
p.log.Info("request",
|
p.log.Info("request",
|
||||||
"route", route,
|
"route", route,
|
||||||
"host", host,
|
"host", host,
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
+85
-34
@@ -131,9 +131,9 @@ func hasModel(loaded []string, model string) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
// allowedPath reports whether rest may be proxied: under /v1/, or the two admin paths.
|
// allowedPath reports whether rest may be proxied: under /v1/, or the admin and control-plane paths.
|
||||||
func allowedPath(rest string) bool {
|
func allowedPath(rest string) bool {
|
||||||
return strings.HasPrefix(rest, "/v1/") || rest == "/health" || rest == "/props"
|
return strings.HasPrefix(rest, "/v1/") || rest == "/health" || rest == "/props" || rest == "/slots" || rest == "/tokenize"
|
||||||
}
|
}
|
||||||
|
|
||||||
// route resolves the route name and the upstream path (rest) from the request, honouring the
|
// route resolves the route name and the upstream path (rest) from the request, honouring the
|
||||||
@@ -198,26 +198,61 @@ func peekModel(r *http.Request) (string, []byte, error) {
|
|||||||
return req.Model, body, nil
|
return req.Model, body, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ServeHTTP routes, fingerprints, leases a host, queues per (host, model), forwards with streaming,
|
// ServeHTTP routes, then serves the request over the shared flow below: fingerprint, lease a host,
|
||||||
// tees the response for usage/timings, and records one accounting row. Every error answer is JSON
|
// queue per (host, model), forward with streaming, tee the response for usage/timings, and record
|
||||||
// {"error":"…"}.
|
// one accounting row. Every error answer is JSON {"error":"…"}.
|
||||||
func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||||
route, rest, code, msg := p.route(r)
|
route, rest, code, msg := p.route(r)
|
||||||
if code != 0 {
|
if code != 0 {
|
||||||
p.writeError(w, code, msg)
|
p.writeError(w, code, msg)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
p.serve(w, r, route, rest)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ForRoute serves route alone, for its dedicated listener: every request there is this route, the
|
||||||
|
// whole path passed upstream as it is (no route segment is taken from it). A request from an unknown
|
||||||
|
// route is 404; a different X-Crossbar-Route header is a 400 (an equal one is ignored); a prefixed
|
||||||
|
// path or a path outside the allowed set is 404. Then it runs the same serve flow as ServeHTTP.
|
||||||
|
func (p *Handler) ForRoute(route string) http.Handler {
|
||||||
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if _, _, ok := p.cfg.Route(route); !ok {
|
||||||
|
p.writeError(w, http.StatusNotFound, "unknown route")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if hdr := r.Header.Get(RouteHeader); hdr != "" && hdr != route {
|
||||||
|
p.writeError(w, http.StatusBadRequest, "conflicting route")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
rest := r.URL.Path
|
||||||
|
if !allowedPath(rest) {
|
||||||
|
p.writeError(w, http.StatusNotFound, "not found")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
p.serve(w, r, route, rest)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// serve fingerprints, leases a host, queues per (host, model), forwards with streaming, tees the
|
||||||
|
// response for usage/timings, and records one accounting row. Both ServeHTTP and ForRoute reach this
|
||||||
|
// after they have settled the route and the upstream path (rest).
|
||||||
|
func (p *Handler) serve(w http.ResponseWriter, r *http.Request, route, rest string) {
|
||||||
routeCfg, _, _ := p.cfg.Route(route)
|
routeCfg, _, _ := p.cfg.Route(route)
|
||||||
|
isControl := isControlCall(r.Method, rest)
|
||||||
|
|
||||||
model, body, err := peekModel(r)
|
model, body, err := peekModel(r)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
p.writeError(w, http.StatusRequestEntityTooLarge, "body too large")
|
p.writeError(w, http.StatusRequestEntityTooLarge, "body too large")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if model == "" {
|
model = resolveModel(model, r, routeCfg)
|
||||||
model = routeCfg.DefaultModel
|
|
||||||
}
|
|
||||||
fp := fingerprint.Of(body)
|
fp := fingerprint.Of(body)
|
||||||
|
// A route-affinity route puts every request (chat or control) on one lease per model, so the
|
||||||
|
// lease key's fingerprint is "" for all of them; the real fingerprint is kept for the row below.
|
||||||
|
leaseFP := fp
|
||||||
|
if routeCfg.PerRoute() {
|
||||||
|
leaseFP = ""
|
||||||
|
}
|
||||||
started := time.Now()
|
started := time.Now()
|
||||||
|
|
||||||
// v0 compatibility path: no lease table, no limiter, no recording.
|
// v0 compatibility path: no lease table, no limiter, no recording.
|
||||||
@@ -227,19 +262,19 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.forward(w, r, route, name, "", rest, fp, model, started, 0, 0, "")
|
p.forward(w, r, route, name, "", rest, fp, model, started, 0, 0, "", isControl)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// Lease. The route's ordered host list is the candidate set.
|
// Lease. The route's ordered host list is the candidate set.
|
||||||
host, reused, err := p.leases.Acquire(lease.Key{Route: route, FP: fp, Model: model}, routeCfg.Hosts, time.Now())
|
host, reused, err := p.leases.Acquire(lease.Key{Route: route, FP: leaseFP, Model: model}, routeCfg.Hosts, time.Now())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
switch {
|
switch {
|
||||||
case errors.Is(err, lease.ErrNoHost):
|
case errors.Is(err, lease.ErrNoHost):
|
||||||
// No host healthy. Ask a waker to rouse a sleeping one; it answers
|
// No host healthy. Ask a waker to rouse a sleeping one; it answers
|
||||||
// (served or 503) when it has had a turn, else falls through to the
|
// (served or 503) when it has had a turn, else falls through to the
|
||||||
// plain 503.
|
// plain 503.
|
||||||
if p.waker != nil && p.wakeOnErrNoHost(w, r, route, routeCfg, rest, model, fp, started, lease.Key{Route: route, FP: fp, Model: model}) {
|
if p.waker != nil && p.wakeOnErrNoHost(w, r, route, routeCfg, rest, model, fp, started, lease.Key{Route: route, FP: leaseFP, Model: model}) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
||||||
@@ -252,17 +287,44 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Slot, context guard and forward, holding the slot for the leased host.
|
// Slot, context guard and forward, holding the slot for the leased host.
|
||||||
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, host, reused)
|
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, host, reused, isControl)
|
||||||
}
|
}
|
||||||
|
|
||||||
// serveLeased queues the request against the leased host's limiter, runs the
|
// serveLeased queues the request against the leased host's limiter, runs the
|
||||||
// context guard, and forwards. The slot is held for the originally leased host
|
// context guard, and forwards. The slot is held for the originally leased host
|
||||||
// even if the guard relocates the lease: the guard already moved it.
|
// even if the guard relocates the lease: the guard already moved it.
|
||||||
func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, host string, reused bool) {
|
func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, host string, reused bool, control bool) {
|
||||||
// Slot. A full queue is a 503; a context done while waiting means the client left.
|
// A control-plane call follows the lease but takes no slot and runs no
|
||||||
release, waited, err := p.lim.Acquire(r.Context(), host, model)
|
// context guard: it is sent beside its own stream, so it must never wait
|
||||||
if err != nil {
|
// for or hold a slot. forward writes no row and logs at Debug for it.
|
||||||
if errors.Is(err, limiter.ErrQueueFull) {
|
if control {
|
||||||
|
p.forward(w, r, route, host, leaseState(reused), rest, fp, model, started, 0, 0, "", true)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// Slot or track. A queue = false route leaves queueing to the client's own
|
||||||
|
// llama-server slot: crossbar only counts the request on the host, never
|
||||||
|
// holding it or refusing it.
|
||||||
|
var waited time.Duration
|
||||||
|
var release func()
|
||||||
|
if routeCfg.Queues() {
|
||||||
|
var err error
|
||||||
|
release, waited, err = p.lim.Acquire(r.Context(), host, model)
|
||||||
|
if err != nil {
|
||||||
|
if errors.Is(err, limiter.ErrQueueFull) {
|
||||||
|
p.writeRecord(store.Request{
|
||||||
|
Route: route,
|
||||||
|
FP: fp,
|
||||||
|
Model: model,
|
||||||
|
Host: host,
|
||||||
|
Started: started,
|
||||||
|
TotalMs: time.Since(started).Milliseconds(),
|
||||||
|
Status: http.StatusServiceUnavailable,
|
||||||
|
Err: "queue full",
|
||||||
|
})
|
||||||
|
p.writeError(w, http.StatusServiceUnavailable, "queue full")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
||||||
p.writeRecord(store.Request{
|
p.writeRecord(store.Request{
|
||||||
Route: route,
|
Route: route,
|
||||||
FP: fp,
|
FP: fp,
|
||||||
@@ -270,24 +332,13 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri
|
|||||||
Host: host,
|
Host: host,
|
||||||
Started: started,
|
Started: started,
|
||||||
TotalMs: time.Since(started).Milliseconds(),
|
TotalMs: time.Since(started).Milliseconds(),
|
||||||
Status: http.StatusServiceUnavailable,
|
Status: 499,
|
||||||
Err: "queue full",
|
Err: "client cancelled while queued",
|
||||||
})
|
})
|
||||||
p.writeError(w, http.StatusServiceUnavailable, "queue full")
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
} else {
|
||||||
p.writeRecord(store.Request{
|
release = p.lim.Track(host, model)
|
||||||
Route: route,
|
|
||||||
FP: fp,
|
|
||||||
Model: model,
|
|
||||||
Host: host,
|
|
||||||
Started: started,
|
|
||||||
TotalMs: time.Since(started).Milliseconds(),
|
|
||||||
Status: 499,
|
|
||||||
Err: "client cancelled while queued",
|
|
||||||
})
|
|
||||||
return
|
|
||||||
}
|
}
|
||||||
defer release()
|
defer release()
|
||||||
|
|
||||||
@@ -298,7 +349,7 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri
|
|||||||
if done {
|
if done {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.forward(w, r, route, host, leaseState(reused), rest, fp, model, now, waited, 0, header)
|
p.forward(w, r, route, host, leaseState(reused), rest, fp, model, now, waited, 0, header, false)
|
||||||
}
|
}
|
||||||
|
|
||||||
// wakeOnErrNoHost answers the request when no host was healthy. It asks, in
|
// wakeOnErrNoHost answers the request when no host was healthy. It asks, in
|
||||||
@@ -316,7 +367,7 @@ func (p *Handler) wakeOnErrNoHost(w http.ResponseWriter, r *http.Request, route
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if newHost, _, err := p.leases.Acquire(key, routeCfg.Hosts, time.Now()); err == nil {
|
if newHost, _, err := p.leases.Acquire(key, routeCfg.Hosts, time.Now()); err == nil {
|
||||||
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, newHost, true)
|
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, newHost, true, isControlCall(r.Method, rest))
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -257,7 +257,9 @@ func TestV0BehaviourStillHolds(t *testing.T) {
|
|||||||
}{
|
}{
|
||||||
{http.MethodGet, "/", 400, "missing route"},
|
{http.MethodGet, "/", 400, "missing route"},
|
||||||
{http.MethodGet, "/nope/v1/models", 404, "unknown route"},
|
{http.MethodGet, "/nope/v1/models", 404, "unknown route"},
|
||||||
{http.MethodGet, "/r/slots", 404, "not found"},
|
// 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"},
|
{http.MethodGet, "/r/_crossbar/hosts", 404, "not found"},
|
||||||
} {
|
} {
|
||||||
req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil)
|
req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil)
|
||||||
|
|||||||
+12
-1
@@ -1,5 +1,6 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# Smoke run (v2): everything v1 checked, plus the context guard, wake-on-LAN and identity gating.
|
# 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.
|
# Prints "smoke: ok" or fails with the crossbar log.
|
||||||
set -eu
|
set -eu
|
||||||
cd "$(dirname "$0")/.."
|
cd "$(dirname "$0")/.."
|
||||||
@@ -57,4 +58,14 @@ 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/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_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"
|
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)"
|
echo "smoke: ok (stream spread $((lastms - firstms)) ms)"
|
||||||
|
|||||||
Reference in New Issue
Block a user