diff --git a/README.md b/README.md index 5e0f103..09f80fb 100644 --- a/README.md +++ b/README.md @@ -59,6 +59,9 @@ hosts = ["beta", "alpha"] | `hosts..models` | The models this host serves, with per-model parallel tuning. | | `routes..hosts` | Candidate hosts, tried in order until one is healthy; a conversation leases one of them. | | `routes..default_model` | Model used when a request omits one; must be served by a host in the route. | +| `routes..affinity` | `"conversation"` (default, one lease per conversation) or `"route"` (one lease for the whole route); see "Clients that manage their own slots". | +| `routes..queue` | `false` leaves queueing to the client's own llama-server slot; the default counts requests in crossbar's per-(host, model) queue. | +| `routes..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. | | `hosts..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..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.: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/` 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 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 it, so the conversation stays there), and the response carries `X-Crossbar-Ctx: moved:` + `>` + `` — for example `moved:small>big`. When no host can -fit it, the answer is `400 {"error":"prompt too large","estimate":,"max":}`. Hosts whose context is unknown are never -blocked by the guard. +fit it, the answer is a `400` in llama-server's own overflow shape, so a client that handles the +server's error handles crossbar's refusal too: + +```json +{"error":{"code":400,"type":"exceed_context_size_error","message":"prompt too large","n_prompt_tokens":,"n_ctx":}} +``` + +Hosts whose context is unknown are never blocked by the guard. ## Wake @@ -204,4 +252,4 @@ tests only, so crossbar logs a warning when it starts in that mode. ## 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`. diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index edcfc17..3831ba8 100644 --- a/cmd/crossbar/main.go +++ b/cmd/crossbar/main.go @@ -11,6 +11,7 @@ import ( "net/http" "os" "os/signal" + "sort" "syscall" "time" @@ -94,21 +95,55 @@ func run() error { p := proxy.New(cfg, table, leases, lim, st, log) 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" { - var checker *identity.Checker switch cfg.Identity { case "tailscale": checker = identity.NewChecker(identity.TailscaleResolver{}) default: // "header" checker = identity.NewHeaderChecker() } + } + var handler http.Handler = p + if checker != nil { handler = identity.Middleware(checker, func(route string) ([]string, bool) { rt, _, ok := cfg.Route(route) return rt.Peers, ok }, 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.Handle("/_crossbar/", admin.Handler(cfg, table, leases, lim, st, hosts)) mux.Handle("/", handler) @@ -173,22 +208,47 @@ func run() error { ReadHeaderTimeout: 10 * time.Second, } - serverErr := make(chan error, 1) - go func() { - log.Info("listening", "addr", srv.Addr) - serverErr <- srv.ListenAndServe() - }() + // Every listener shuts down together on ctx done; the first error other than a clean shutdown + // ends run and shuts the rest down. + servers := make([]*http.Server, 0, 1+len(routes)) + 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 { case <-ctx.Done(): log.Info("shutting down") shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() - return srv.Shutdown(shutdownCtx) + for _, s := range servers { + _ = s.Shutdown(shutdownCtx) + } + return nil case err := <-serverErr: if errors.Is(err, http.ErrServerClosed) { return nil } + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + for _, s := range servers { + _ = s.Shutdown(shutdownCtx) + } return err } } diff --git a/docs/implementer-log.md b/docs/implementer-log.md index aafd711..d8c6de5 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -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 | |---|---|---|---|---|---|---|---| +| 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":,"n_ctx":}}`, the accounting row unchanged (status 400, Err "prompt too large"), every other error keeping `{"error":""}`; 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/` 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..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/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 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=` 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 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. diff --git a/example.toml b/example.toml index 4c3a1ff..ccbb4ee 100644 --- a/example.toml +++ b/example.toml @@ -31,3 +31,13 @@ default_model = "ornith-1.5-35b-a3b" [routes.hermes-x] hosts = ["beta", "alpha"] # peers = ["talos"] # v2: with identity = "tailscale", only these tailnet nodes may use the route + +# v2.3: a client that manages its own llama-server slot (it pins id_slot, polls /slots, steers a +# running completion through /v1/chat/completions/control) and cannot put a route in the path. +# The route gets its own port; every request there is this route and the path goes upstream as is. +[routes.boxmaker-a] +hosts = ["beta", "alpha"] +default_model = "ornith-1.5-35b-a3b" +listen = "127.0.0.1:17801" # a tailnet address in production; never the main listen address +affinity = "route" # one lease for the whole route, not one per conversation +queue = false # counted as load but never held or refused: the server's own slot queue does that diff --git a/internal/config/config.go b/internal/config/config.go index 6559877..c59afc3 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -71,14 +71,6 @@ type Host struct { 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. type Config struct { Listen string `toml:"listen"` @@ -117,8 +109,6 @@ const ( 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 // "config: …", the same shape as a decode failure. func Load(path string) (*Config, error) { @@ -345,54 +335,3 @@ func (c *Config) checkHosts() *Error { } 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 -} diff --git a/internal/config/config_v23_test.go b/internal/config/config_v23_test.go new file mode 100644 index 0000000..c336107 --- /dev/null +++ b/internal/config/config_v23_test.go @@ -0,0 +1,75 @@ +package config_test + +// v2.3 task 02: the affinity and queue route keys. + +import ( + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const affinityBase = ` +listen = "127.0.0.1:1" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.plain] +hosts = ["a"] +[routes.convo] +hosts = ["a"] +affinity = "conversation" +[routes.boxmaker] +hosts = ["a"] +affinity = "route" +queue = false +[routes."bm-*"] +hosts = ["a"] +affinity = "route" +queue = false +[routes.queued] +hosts = ["a"] +queue = true +` + +func TestAffinityAndQueueKeys(t *testing.T) { + c, err := config.Parse(strings.NewReader(affinityBase)) + if err != nil { + t.Fatal(err) + } + for _, tc := range []struct { + route string + perRoute, queues bool + }{ + {"plain", false, true}, // defaults: conversation affinity, queueing on + {"convo", false, true}, + {"boxmaker", true, false}, + {"bm-agent-1", true, false}, // a template's keys reach its concrete routes + {"queued", false, true}, + } { + r, _, ok := c.Route(tc.route) + if !ok { + t.Fatalf("route %q not found", tc.route) + } + if r.PerRoute() != tc.perRoute || r.Queues() != tc.queues { + t.Errorf("%s: PerRoute %v Queues %v, want %v %v", tc.route, r.PerRoute(), r.Queues(), tc.perRoute, tc.queues) + } + } +} + +func TestAffinityRejectsUnknownValues(t *testing.T) { + for _, bad := range []string{`"session"`, `"Route"`, `1`} { + text := strings.Replace(affinityBase, `affinity = "conversation"`, "affinity = "+bad, 1) + _, err := config.Parse(strings.NewReader(text)) + if err == nil || !strings.Contains(err.Error(), "routes.convo.affinity") { + t.Errorf("affinity = %s: err %v, want one naming routes.convo.affinity", bad, err) + } + } +} + +func TestQueueMustBeABool(t *testing.T) { + text := strings.Replace(affinityBase, "queue = true", `queue = "no"`, 1) + if _, err := config.Parse(strings.NewReader(text)); err == nil { + t.Error(`queue = "no" parsed; want an error`) + } +} diff --git a/internal/config/listen_test.go b/internal/config/listen_test.go new file mode 100644 index 0000000..5e1979a --- /dev/null +++ b/internal/config/listen_test.go @@ -0,0 +1,54 @@ +package config_test + +// v2.3 task 03: a concrete route may own a dedicated listener. Every request that arrives on it is +// that route, with the upstream path unprefixed, for clients that cannot put a route in the path +// or a header (Boxmaker's inferproxy rewrites nothing). + +import ( + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const listenBase = ` +listen = "127.0.0.1:7777" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.bm-a] +hosts = ["a"] +listen = "127.0.0.1:7801" +[routes.bm-b] +hosts = ["a"] +listen = "127.0.0.1:7802" +[routes.plain] +hosts = ["a"] +` + +func TestRouteListen(t *testing.T) { + c, err := config.Parse(strings.NewReader(listenBase)) + if err != nil { + t.Fatal(err) + } + for route, want := range map[string]string{"bm-a": "127.0.0.1:7801", "bm-b": "127.0.0.1:7802", "plain": ""} { + if got := c.Routes[route].Listen; got != want { + t.Errorf("%s listen = %q, want %q", route, got, want) + } + } +} + +func TestRouteListenRejected(t *testing.T) { + for _, tc := range []struct{ name, text, want string }{ + {"not host:port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"7801"`, 1), "routes.bm-a.listen"}, + {"bad port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:http"`, 1), "routes.bm-a.listen"}, + {"port zero", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:0"`, 1), "routes.bm-a.listen"}, + {"same as another route", strings.Replace(listenBase, `"127.0.0.1:7802"`, `"127.0.0.1:7801"`, 1), "listen"}, + {"same as the main listener", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:7777"`, 1), "routes.bm-a.listen"}, + {"on a template", listenBase + "[routes.\"t-*\"]\nhosts = [\"a\"]\nlisten = \"127.0.0.1:7803\"\n", "t-*"}, + } { + if _, err := config.Parse(strings.NewReader(tc.text)); err == nil || !strings.Contains(err.Error(), tc.want) { + t.Errorf("%s: err %v, want one containing %q", tc.name, err, tc.want) + } + } +} diff --git a/internal/config/route.go b/internal/config/route.go index 43f3ce4..da9c88a 100644 --- a/internal/config/route.go +++ b/internal/config/route.go @@ -1,13 +1,44 @@ package config import ( + "fmt" + "net" "regexp" + "sort" + "strconv" "strings" ) // templateName matches a route template: a valid route name ending in "-*". 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 // "-*" 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). @@ -38,3 +69,100 @@ func (c *Config) Route(name string) (r Route, key string, ok bool) { } 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..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 +} diff --git a/internal/identity/middleware.go b/internal/identity/middleware.go index e91b3b8..0397f51 100644 --- a/internal/identity/middleware.go +++ b/internal/identity/middleware.go @@ -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. func writeForbidden(w http.ResponseWriter) { w.Header().Set("Content-Type", "application/json") diff --git a/internal/identity/route_middleware_test.go b/internal/identity/route_middleware_test.go new file mode 100644 index 0000000..c860885 --- /dev/null +++ b/internal/identity/route_middleware_test.go @@ -0,0 +1,47 @@ +package identity_test + +// v2.3 task 03: on a route's dedicated listener the route is fixed, so the gate is that route's +// peers for every request, whatever path or X-Crossbar-Route header the caller sends. + +import ( + "net/http" + "net/http/httptest" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/identity" +) + +func TestRouteMiddleware(t *testing.T) { + inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) }) + checker := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"}) + locked := identity.RouteMiddleware(checker, []string{"talos"}, inner) + open := identity.RouteMiddleware(checker, nil, inner) + for _, tc := range []struct { + name string + h http.Handler + path, hdr string + addr string + want int + }{ + {"right peer", locked, "/v1/chat/completions", "", "100.64.0.5:5", 204}, + {"wrong peer", locked, "/v1/chat/completions", "", "100.64.0.9:5", 403}, + {"not a peer", locked, "/slots", "", "203.0.113.1:5", 403}, + {"a path that looks like an open route is still this route", locked, "/open/v1/models", "", "100.64.0.9:5", 403}, + {"a header naming another route does not change the gate", locked, "/v1/models", "open", "100.64.0.9:5", 403}, + {"admin-looking path is gated too (no admin on this listener)", locked, "/_crossbar/hosts", "", "100.64.0.9:5", 403}, + {"open route, anyone", open, "/v1/models", "", "203.0.113.1:5", 204}, + } { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, tc.path, nil) + req.RemoteAddr = tc.addr + if tc.hdr != "" { + req.Header.Set("X-Crossbar-Route", tc.hdr) + } + rec := httptest.NewRecorder() + tc.h.ServeHTTP(rec, req) + if rec.Code != tc.want { + t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String()) + } + }) + } +} diff --git a/internal/limiter/limiter.go b/internal/limiter/limiter.go index 267e603..29ea8b4 100644 --- a/internal/limiter/limiter.go +++ b/internal/limiter/limiter.go @@ -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 -// if one is waiting, otherwise it frees the slot. It is safe to call through the sync.Once that -// Acquire wrapped it in. +// only while there is room (in flight at or below parallel), otherwise it counts the slot back. It is +// 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() { var once sync.Once return func() { once.Do(func() { l.mu.Lock() defer l.mu.Unlock() - if len(p.waiters) > 0 { + if len(p.waiters) > 0 && p.inflight <= p.parallel { next := p.waiters[0] p.waiters = p.waiters[1:] close(next) diff --git a/internal/limiter/track_test.go b/internal/limiter/track_test.go new file mode 100644 index 0000000..9e8464f --- /dev/null +++ b/internal/limiter/track_test.go @@ -0,0 +1,91 @@ +package limiter_test + +// v2.3 task 02: Track counts a request without holding or refusing it. A route with queue = false +// leaves queueing to llama-server's own slots, but its requests are still load on the host, so the +// routes that do queue must see them. + +import ( + "context" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/limiter" +) + +func TestTrackNeverWaitsAndCounts(t *testing.T) { + l := limiter.New() + l.Configure("alpha", "m", 1, 0) // one slot, no waiting room + + start := time.Now() + rel1 := l.Track("alpha", "m") + rel2 := l.Track("alpha", "m") + rel3 := l.Track("alpha", "m") + if d := time.Since(start); d > 50*time.Millisecond { + t.Fatalf("Track waited %v", d) + } + if n := l.InFlight("alpha", "m"); n != 3 { + t.Fatalf("in flight = %d, want 3 (Track may pass parallel)", n) + } + if n := l.FreeSlots("alpha"); n != 0 { + t.Errorf("free slots = %d, want 0", n) + } + // A queueing request sees the host full: no waiting room, so it is refused. + if _, _, err := l.Acquire(context.Background(), "alpha", "m"); err == nil { + t.Error("Acquire on an over-tracked pair succeeded; want ErrQueueFull") + } + rel1() + rel1() // idempotent + rel2() + rel3() + if n := l.InFlight("alpha", "m"); n != 0 { + t.Errorf("in flight after release = %d, want 0", n) + } +} + +// A waiter gets a slot only once in flight is back under parallel: releasing a tracked request +// while the pair is still over its limit must not hand the slot on. +func TestTrackReleaseHandsOverOnlyUnderTheLimit(t *testing.T) { + l := limiter.New() + l.Configure("alpha", "m", 1, 1) + relA := l.Track("alpha", "m") + relB := l.Track("alpha", "m") // in flight 2, parallel 1 + + got := make(chan func(), 1) + go func() { + rel, _, err := l.Acquire(context.Background(), "alpha", "m") + if err != nil { + t.Error(err) + close(got) + return + } + got <- rel + }() + waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 }) + + relA() // in flight 1 == parallel: still no free slot + select { + case <-got: + t.Fatal("waiter got a slot while in flight was still at parallel") + case <-time.After(100 * time.Millisecond): + } + if n := l.InFlight("alpha", "m"); n != 1 { + t.Fatalf("in flight = %d after one release, want 1", n) + } + + relB() // now the slot is free: hand it to the waiter + select { + case rel := <-got: + if rel == nil { + t.Fatal("waiter failed") + } + if n := l.InFlight("alpha", "m"); n != 1 { + t.Errorf("in flight = %d with the waiter running, want 1", n) + } + rel() + case <-time.After(2 * time.Second): + t.Fatal("waiter never got the freed slot") + } + if n := l.InFlight("alpha", "m"); n != 0 { + t.Errorf("in flight at the end = %d, want 0", n) + } +} diff --git a/internal/proxy/affinity_test.go b/internal/proxy/affinity_test.go new file mode 100644 index 0000000..476889e --- /dev/null +++ b/internal/proxy/affinity_test.go @@ -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 }) + }) + } +} diff --git a/internal/proxy/control.go b/internal/proxy/control.go new file mode 100644 index 0000000..1ba6664 --- /dev/null +++ b/internal/proxy/control.go @@ -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 +} diff --git a/internal/proxy/control_test.go b/internal/proxy/control_test.go new file mode 100644 index 0000000..c0d150e --- /dev/null +++ b/internal/proxy/control_test.go @@ -0,0 +1,176 @@ +package proxy_test + +// v2.3 task 01: control-plane requests. A client that manages its own slots (Boxmaker) polls +// /slots, reads /props, tokenizes and steers a running completion through +// /v1/chat/completions/control. Those calls follow the route's lease like any other request but +// must never wait for, or take, a slot: /control is sent while the client's own stream holds one. + +import ( + "context" + "net/http" + "strings" + "testing" + "time" +) + +// controlClient gives every control call a short deadline: a call that queues behind a full host +// is the bug, and it must fail the test rather than hang it. +var controlClient = &http.Client{Timeout: 2 * time.Second} + +func (r *rig) do(method, path, body string) *http.Response { + r.t.Helper() + var rd *strings.Reader + if body != "" { + rd = strings.NewReader(body) + } + var req *http.Request + var err error + if rd != nil { + req, err = http.NewRequest(method, r.front.URL+path, rd) + req.Header.Set("Content-Type", "application/json") + } else { + req, err = http.NewRequest(method, r.front.URL+path, nil) + } + if err != nil { + r.t.Fatal(err) + } + resp, err := controlClient.Do(req) + if err != nil { + r.t.Fatalf("%s %s: %v", method, path, err) + } + return resp +} + +func (r *rig) rows(route string) int64 { + r.t.Helper() + counts, err := r.store.StatusCounts(time.Time{}) + if err != nil { + r.t.Fatal(err) + } + var n int64 + for _, c := range counts { + if c.Route == route { + n += c.Count + } + } + return n +} + +// A GET names its model in the query string: /slots?model=alpha-only must reach the host that +// has alpha-only loaded, not whichever host the route's default model would pick. +func TestGetModelComesFromTheQuery(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + + resp := r.do(http.MethodGet, "/r/slots?model=alpha-only", "") + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != "alpha" { + t.Fatalf("GET /r/slots?model=alpha-only: %d on %q, want 200 on alpha", resp.StatusCode, resp.Header.Get("X-Crossbar-Host")) + } + if got := alpha.lastReq(); got.method != "GET" || got.path != "/slots?model=alpha-only" { + t.Errorf("alpha saw %s %s, want GET /slots?model=alpha-only", got.method, got.path) + } + // The same for beta-only, so a lucky default cannot pass the test. + resp = r.do(http.MethodGet, "/r/slots?model=beta-only", "") + drain(resp) + if resp.Header.Get("X-Crossbar-Host") != "beta" { + t.Errorf("GET /r/slots?model=beta-only went to %q, want beta", resp.Header.Get("X-Crossbar-Host")) + } +} + +// /slots and /tokenize are proxied; the per-slot actions under /slots/ (save, restore, erase) +// are not. +func TestControlPathsAllowed(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + for _, tc := range []struct { + method, path, body string + want int + }{ + {http.MethodGet, "/r/slots", "", 200}, + {http.MethodGet, "/r/slots?model=shared", "", 200}, + {http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`, 200}, + {http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`, 200}, + {http.MethodGet, "/r/slots/0", "", 404}, + {http.MethodPost, "/r/slots/0?action=erase", "", 404}, + {http.MethodPost, "/r/slots/0?action=save", `{"filename":"x"}`, 404}, + } { + resp := r.do(tc.method, tc.path, tc.body) + body := drain(resp) + if resp.StatusCode != tc.want { + t.Errorf("%s %s: %d %s, want %d", tc.method, tc.path, resp.StatusCode, body, tc.want) + } + } +} + +// With every slot on both hosts taken and the queue full, control-plane calls still go straight +// through: no 503, no wait, no slot taken, no accounting row. +func TestControlRequestsNeverTakeASlot(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + + // Take every "shared" slot (parallel 2 on each host) and the one queue place per host. + var releases []func() + for _, host := range []string{"alpha", "beta"} { + for i := 0; i < 2; i++ { + rel, _, err := r.lim.Acquire(context.Background(), host, "shared") + if err != nil { + t.Fatal(err) + } + releases = append(releases, rel) + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go func() { _, _, _ = r.lim.Acquire(ctx, host, "shared") }() + waitUntil(t, func() bool { return r.lim.Queued(host, "shared") == 1 }) + } + defer func() { + for _, rel := range releases { + rel() + } + }() + + for _, tc := range []struct{ method, path, body string }{ + {http.MethodGet, "/r/slots?model=shared", ""}, + {http.MethodGet, "/r/props?model=shared", ""}, + {http.MethodHead, "/r/props?model=shared", ""}, + {http.MethodGet, "/r/v1/models", ""}, + {http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`}, + {http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`}, + } { + resp := r.do(tc.method, tc.path, tc.body) + body := drain(resp) + if resp.StatusCode != 200 { + t.Errorf("%s %s with the host full: %d %s, want 200", tc.method, tc.path, resp.StatusCode, body) + } + } + for _, host := range []string{"alpha", "beta"} { + if n := r.lim.InFlight(host, "shared"); n != 2 { + t.Errorf("%s in flight = %d after control calls, want 2 (control takes no slot)", host, n) + } + } + time.Sleep(100 * time.Millisecond) // a row is written after the answer; give a stray one time to land + if n := r.rows("r"); n != 0 { + t.Errorf("control calls wrote %d accounting rows, want 0", n) + } + + // A chat completion on the same full route still queues or is refused as before: the bypass + // is for control calls only. + resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1)) + drain(resp) + if resp.StatusCode != http.StatusServiceUnavailable { + t.Errorf("chat on a full route: %d, want 503 (queue full)", resp.StatusCode) + } +} + +// A chat completion is not a control call just because its path starts the same way. +func TestChatIsNotControl(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1)) + drain(resp) + if resp.StatusCode != 200 { + t.Fatalf("chat: %d", resp.StatusCode) + } + waitUntil(t, func() bool { return r.rows("r") == 1 }) // the row lands just after the answer +} diff --git a/internal/proxy/ctxguard.go b/internal/proxy/ctxguard.go index 696de8b..0bad978 100644 --- a/internal/proxy/ctxguard.go +++ b/internal/proxy/ctxguard.go @@ -155,10 +155,21 @@ func largestSlotCtx(hosts []string, h Health, model string) int { return best } -// refuseCtx answers the 400 the guard's rule 4: the JSON body carries the -// estimate and the largest available per-slot context, plus the error text. It -// records the accounting row (status 400, Err "prompt too large") and never -// marks the host down. +// ctxErrorBody is llama-server's shape for a prompt that exceeds a host's +// context. A client that already handles the server's own overflow error keys +// on error.type and so recognises crossbar's refusal too. See refuseCtx. +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) { p.writeRecord(store.Request{ 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.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]any{ - "error": "prompt too large", - "estimate": estimate, - "max": maxSlot, + "error": ctxErrorBody{ + Code: http.StatusBadRequest, + Type: "exceed_context_size_error", + Message: "prompt too large", + NPromptTokens: estimate, + NCtx: maxSlot, + }, }) } diff --git a/internal/proxy/ctxguard_router_test.go b/internal/proxy/ctxguard_router_test.go index 7c4ffb1..b30dcf8 100644 --- a/internal/proxy/ctxguard_router_test.go +++ b/internal/proxy/ctxguard_router_test.go @@ -95,11 +95,16 @@ func TestRouterUnloadedModelIsNotACandidate(t *testing.T) { if big.hits.Load() != 0 { 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 { t.Fatalf("body %q is not JSON: %v", body, err) } - if max, _ := e["max"].(float64); max != 4096 { - t.Errorf("max = %v, want 4096: the largest per-slot context among hosts that have shared loaded", e["max"]) + if e.Error.NCtx != 4096 { + t.Errorf("error.n_ctx = %v, want 4096: the largest per-slot context among hosts that have shared loaded", e.Error.NCtx) } } diff --git a/internal/proxy/ctxguard_test.go b/internal/proxy/ctxguard_test.go index b56729b..92e3318 100644 --- a/internal/proxy/ctxguard_test.go +++ b/internal/proxy/ctxguard_test.go @@ -84,15 +84,31 @@ func TestOversizedPromptWithNoFitIs400(t *testing.T) { if resp.StatusCode != http.StatusBadRequest { t.Fatalf("status %d body %s, want 400", resp.StatusCode, body) } - var e map[string]any - if err := json.Unmarshal([]byte(body), &e); err != nil || e["error"] != "prompt too large" { - t.Fatalf("body = %s, want error 'prompt too large'", body) + // v2.3: llama-server's own shape for this error, so a client handles crossbar's refusal the + // way it handles the server's (Boxmaker keys on error.type; the error JSON must come first). + if !strings.HasPrefix(body, `{"error":`) { + t.Errorf("body must start with the error object: %s", body) } - if est, _ := e["estimate"].(float64); est < 8000 || est > 13000 { - t.Errorf("estimate = %v, want roughly 10000 tokens", e["estimate"]) + var e struct { + Error struct { + Code int `json:"code"` + Type string `json:"type"` + Message string `json:"message"` + NPromptTokens float64 `json:"n_prompt_tokens"` + NCtx float64 `json:"n_ctx"` + } `json:"error"` } - if max, _ := e["max"].(float64); max != 4096 { - t.Errorf("max = %v, want the largest per-slot context among the route's hosts (4096)", e["max"]) + if err := json.Unmarshal([]byte(body), &e); err != nil || e.Error.Code != 400 || e.Error.Type != "exceed_context_size_error" || e.Error.Message != "prompt too large" { + t.Fatalf("body = %s, want {\"error\":{\"code\":400,\"type\":\"exceed_context_size_error\",\"message\":\"prompt too large\",…}}", body) + } + if est := e.Error.NPromptTokens; est < 8000 || est > 13000 { + t.Errorf("n_prompt_tokens = %v, want roughly 10000 tokens", est) + } + if max := e.Error.NCtx; max != 4096 { + t.Errorf("n_ctx = %v, want the largest per-slot context among the route's hosts (4096)", max) + } + if ct := resp.Header.Get("Content-Type"); !strings.HasPrefix(ct, "application/json") { + t.Errorf("Content-Type = %q, want application/json", ct) } if small.hits.Load()+tiny.hits.Load() != 0 { t.Errorf("a refused prompt must not reach any upstream") diff --git a/internal/proxy/forward.go b/internal/proxy/forward.go index 3bbf17f..84c6906 100644 --- a/internal/proxy/forward.go +++ b/internal/proxy/forward.go @@ -16,8 +16,9 @@ import ( // 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 // prompt size the context guard estimated (0 when the guard did not run); ctxHeader is the -// "moved:…" header to set when the guard relocated the conversation. -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) { +// "moved:…" header to set when the guard relocated the conversation. control is true for a +// 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] if !ok { 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 { req.Err = "upstream error" } - p.writeRecord(req) + if !control { + p.writeRecord(req) + } panic(pv) } }() @@ -58,12 +61,29 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l req.Status = 499 req.Err = "client cancelled" } - p.writeRecord(req) + if !control { + p.writeRecord(req) + } fp8 := fp if len(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", "route", route, "host", host, diff --git a/internal/proxy/listener_test.go b/internal/proxy/listener_test.go new file mode 100644 index 0000000..650cfc6 --- /dev/null +++ b/internal/proxy/listener_test.go @@ -0,0 +1,112 @@ +package proxy_test + +// v2.3 task 03: Handler.ForRoute serves one route with unprefixed paths, for a route's dedicated +// listener. + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/proxy" +) + +// dedicated serves r's route on its own test server, sharing r's health, leases, limiter and +// store, as main does for a route with listen set. +func dedicated(t *testing.T, r *rig, route string) *httptest.Server { + p := proxy.New(r.cfg, r.health, r.leases, r.lim, r.store, nil) + srv := httptest.NewServer(p.ForRoute(route)) + t.Cleanup(srv.Close) + return srv +} + +func call(t *testing.T, method, url, body string, hdr ...string) (*http.Response, string) { + t.Helper() + var req *http.Request + if body != "" { + req, _ = http.NewRequest(method, url, strings.NewReader(body)) + req.Header.Set("Content-Type", "application/json") + } else { + req, _ = http.NewRequest(method, url, nil) + } + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } + resp, err := controlClient.Do(req) + if err != nil { + t.Fatalf("%s %s: %v", method, url, err) + } + return resp, drain(resp) +} + +func TestForRouteServesUnprefixedPaths(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "bm") + + resp, body := call(t, http.MethodPost, srv.URL+"/v1/chat/completions", conversation(1, 1)) + if resp.StatusCode != 200 { + t.Fatalf("chat on the dedicated listener: %d %s", resp.StatusCode, body) + } + host := resp.Header.Get("X-Crossbar-Host") + up := map[string]*upstream{"alpha": alpha, "beta": beta}[host] + if up == nil || up.lastReq().path != "/v1/chat/completions" { + t.Fatalf("upstream %q saw %+v, want /v1/chat/completions unchanged", host, up.lastReq()) + } + resp, _ = call(t, http.MethodGet, srv.URL+"/slots?model=shared", "") + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host || up.lastReq().path != "/slots?model=shared" { + t.Errorf("/slots: %d on %q (last %+v), want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), up.lastReq(), host) + } + resp, _ = call(t, http.MethodPost, srv.URL+"/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`) + if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host { + t.Errorf("/control: %d on %q, want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), host) + } + // The chat is accounted to the route the listener serves. + waitUntil(t, func() bool { return r.rows("bm") == 1 }) + // The same route through the main listener shares the lease: same host. + resp = r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(2, 1)) + drain(resp) + if resp.Header.Get("X-Crossbar-Host") != host { + t.Errorf("main listener /bm went to %q, dedicated to %q; one route, one lease", resp.Header.Get("X-Crossbar-Host"), host) + } +} + +func TestForRouteRefusals(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "bm") + for _, tc := range []struct { + name, method, path string + hdr []string + want int + msg string + }{ + {"a prefixed path is not stripped", http.MethodGet, "/bm/v1/models", nil, 404, "not found"}, + {"no admin here", http.MethodGet, "/_crossbar/hosts", nil, 404, "not found"}, + {"root", http.MethodGet, "/", nil, 404, "not found"}, + {"header naming another route", http.MethodGet, "/v1/models", []string{"X-Crossbar-Route", "r"}, 400, "conflicting route"}, + } { + resp, body := call(t, tc.method, srv.URL+tc.path, "", tc.hdr...) + var e map[string]string + if resp.StatusCode != tc.want || json.Unmarshal([]byte(body), &e) != nil || e["error"] != tc.msg { + t.Errorf("%s: %d %s, want %d %q", tc.name, resp.StatusCode, body, tc.want, tc.msg) + } + } + // A header naming this same route is harmless. + resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "", "X-Crossbar-Route", "bm") + if resp.StatusCode != 200 { + t.Errorf("header naming the listener's own route: %d %s, want 200", resp.StatusCode, body) + } +} + +func TestForRouteUnknownRoute(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, affinityHosts, alpha, beta) + srv := dedicated(t, r, "nope") + resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "") + if resp.StatusCode != 404 || !strings.Contains(body, "unknown route") { + t.Errorf("ForRoute(unknown): %d %s, want 404 unknown route", resp.StatusCode, body) + } +} diff --git a/internal/proxy/proxy.go b/internal/proxy/proxy.go index 7614a13..1bce4d5 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -131,9 +131,9 @@ func hasModel(loaded []string, model string) bool { 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 { - 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 @@ -198,26 +198,61 @@ func peekModel(r *http.Request) (string, []byte, error) { return req.Model, body, nil } -// ServeHTTP routes, fingerprints, leases a host, queues per (host, model), forwards with streaming, -// tees the response for usage/timings, and records one accounting row. Every error answer is JSON -// {"error":"…"}. +// ServeHTTP routes, then serves the request over the shared flow below: fingerprint, lease a host, +// queue per (host, model), forward with streaming, tee the response for usage/timings, and record +// one accounting row. Every error answer is JSON {"error":"…"}. func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { route, rest, code, msg := p.route(r) if code != 0 { p.writeError(w, code, msg) 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) + isControl := isControlCall(r.Method, rest) model, body, err := peekModel(r) if err != nil { p.writeError(w, http.StatusRequestEntityTooLarge, "body too large") return } - if model == "" { - model = routeCfg.DefaultModel - } + model = resolveModel(model, r, routeCfg) 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() // 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") 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 } // 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 { switch { case errors.Is(err, lease.ErrNoHost): // 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 // 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 } 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. - 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 // 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. -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) { - // Slot. A full queue is a 503; a context done while waiting means the client left. - release, waited, err := p.lim.Acquire(r.Context(), host, model) - if err != nil { - if errors.Is(err, limiter.ErrQueueFull) { +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) { + // A control-plane call follows the lease but takes no slot and runs no + // context guard: it is sent beside its own stream, so it must never wait + // for or hold a slot. forward writes no row and logs at Debug for it. + 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{ Route: route, FP: fp, @@ -270,24 +332,13 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri Host: host, Started: started, TotalMs: time.Since(started).Milliseconds(), - Status: http.StatusServiceUnavailable, - Err: "queue full", + Status: 499, + Err: "client cancelled while queued", }) - 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{ - Route: route, - FP: fp, - Model: model, - Host: host, - Started: started, - TotalMs: time.Since(started).Milliseconds(), - Status: 499, - Err: "client cancelled while queued", - }) - return + } else { + release = p.lim.Track(host, model) } defer release() @@ -298,7 +349,7 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri if done { 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 @@ -316,7 +367,7 @@ func (p *Handler) wakeOnErrNoHost(w http.ResponseWriter, r *http.Request, route continue } 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 } } diff --git a/internal/proxy/proxy_test.go b/internal/proxy/proxy_test.go index 897982a..d5f23bb 100644 --- a/internal/proxy/proxy_test.go +++ b/internal/proxy/proxy_test.go @@ -257,7 +257,9 @@ func TestV0BehaviourStillHolds(t *testing.T) { }{ {http.MethodGet, "/", 400, "missing 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"}, } { req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil) diff --git a/tools/smoke.sh b/tools/smoke.sh index 420d546..240a2ab 100755 --- a/tools/smoke.sh +++ b/tools/smoke.sh @@ -1,5 +1,6 @@ #!/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. set -eu 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/metrics" | grep -q 'crossbar_requests_total{route="hermes-x",host="beta",status="403"}' && fail "403s are refused before a lease and must not be counted as requests" curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing beta health" +# 6. v2.3: boxmaker-a has its own listener. Paths are unprefixed, the chat and a control call land +# on the same host (affinity = "route"), and the admin API is not served there. +lb=http://127.0.0.1:17801 +h1=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$lb/v1/chat/completions") +h2=$(hdrs "$lb/props?model=ornith-1.5-35b-a3b") +case "$h1" in 200*) ;; *) fail "chat on the dedicated listener should be 200, got '$h1'";; esac +[ "$(echo "$h1" | cut -d' ' -f2)" = "$(echo "$h2" | cut -d' ' -f2)" ] || fail "chat went to '$h1', /props to '$h2': one route, one host" +h=$(curl -s -o /dev/null -w '%{http_code}' "$lb/_crossbar/hosts"); [ "$h" = "404" ] || fail "admin must not be served on a dedicated listener, got $h" +h=$(curl -s -o /dev/null -w '%{http_code}' "$lb/boxmaker-a/v1/models"); [ "$h" = "404" ] || fail "a prefixed path on the dedicated listener should be 404, got $h" + echo "smoke: ok (stream spread $((lastms - firstms)) ms)"