From fa1c398fc40fd2081170e11a22cbe061df497c79 Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 19:59:21 -0700 Subject: [PATCH] A route may have its own listener: every request there is that route, paths unprefixed Implemented-By: OpenCode session (model recorded in docs/implementer-log.md) --- cmd/crossbar/main.go | 77 ++++++++++++-- docs/implementer-log.md | 1 + example.toml | 10 ++ internal/config/listen_test.go | 54 ++++++++++ internal/config/route.go | 39 +++++++ internal/identity/middleware.go | 15 +++ internal/identity/route_middleware_test.go | 47 +++++++++ internal/proxy/listener_test.go | 112 +++++++++++++++++++++ internal/proxy/proxy.go | 36 ++++++- tools/smoke.sh | 13 ++- 10 files changed, 392 insertions(+), 12 deletions(-) create mode 100644 internal/config/listen_test.go create mode 100644 internal/identity/route_middleware_test.go create mode 100644 internal/proxy/listener_test.go diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index edcfc17..d60e9db 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,56 @@ 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. + logIdentityMode(log, cfg.Identity) + + 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 +209,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 f02f921..08b5552 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -5,6 +5,7 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once | Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model | |---|---|---|---|---|---|---|---| +| 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. | ? | 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/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 c6443cf..da9c88a 100644 --- a/internal/config/route.go +++ b/internal/config/route.go @@ -2,8 +2,10 @@ package config import ( "fmt" + "net" "regexp" "sort" + "strconv" "strings" ) @@ -22,6 +24,7 @@ type Route struct { 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 @@ -80,6 +83,9 @@ func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) 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] @@ -124,6 +130,39 @@ func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) 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/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 d6b7aff..1bce4d5 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -198,15 +198,45 @@ 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) 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)"