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)
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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.<name>.listen: the value must split into host and a numeric port 1-65535, it must not be on
|
||||
// a template route, it must not be the top-level listen, and it must be unique across routes (the
|
||||
// second route in sorted name order reports the clash and the other route's name).
|
||||
func checkListen(name, listen, mainListen string, seenListen map[string]string) *Error {
|
||||
if listen == "" {
|
||||
return nil
|
||||
}
|
||||
_, port, err := net.SplitHostPort(listen)
|
||||
if err != nil {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "must be host:port"}
|
||||
}
|
||||
n, err := strconv.Atoi(port)
|
||||
if err != nil || n < 1 || n > 65535 {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "port must be 1-65535"}
|
||||
}
|
||||
if templateName.MatchString(name) {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "a template route cannot have its own listener"}
|
||||
}
|
||||
if listen == mainListen {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: "cannot be the main listen address"}
|
||||
}
|
||||
if other, dup := seenListen[listen]; dup {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.listen", name), Msg: fmt.Sprintf("already used by route %s", other)}
|
||||
}
|
||||
seenListen[listen] = name
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -45,6 +45,21 @@ func Middleware(c *Checker, peersFor func(route string) ([]string, bool), next h
|
||||
})
|
||||
}
|
||||
|
||||
// RouteMiddleware gates every request on a fixed set of peers, whatever path or X-Crossbar-Route
|
||||
// header it carries: on a route's dedicated listener the route is fixed, so the gate is that route's
|
||||
// peers for every request. There is no admin-path exemption — there is no admin on a route listener.
|
||||
// Empty peers lets everyone through, as for Middleware.
|
||||
func RouteMiddleware(c *Checker, peers []string, next http.Handler) http.Handler {
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
ctx := WithHeaderPeer(r.Context(), r.Header.Get(peerHeader))
|
||||
if err := c.Allow(ctx, peers, r.RemoteAddr); err != nil {
|
||||
writeForbidden(w)
|
||||
return
|
||||
}
|
||||
next.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
|
||||
// writeForbidden answers the JSON 403 the tests and callers expect.
|
||||
func writeForbidden(w http.ResponseWriter) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
package identity_test
|
||||
|
||||
// v2.3 task 03: on a route's dedicated listener the route is fixed, so the gate is that route's
|
||||
// peers for every request, whatever path or X-Crossbar-Route header the caller sends.
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/identity"
|
||||
)
|
||||
|
||||
func TestRouteMiddleware(t *testing.T) {
|
||||
inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) })
|
||||
checker := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"})
|
||||
locked := identity.RouteMiddleware(checker, []string{"talos"}, inner)
|
||||
open := identity.RouteMiddleware(checker, nil, inner)
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
h http.Handler
|
||||
path, hdr string
|
||||
addr string
|
||||
want int
|
||||
}{
|
||||
{"right peer", locked, "/v1/chat/completions", "", "100.64.0.5:5", 204},
|
||||
{"wrong peer", locked, "/v1/chat/completions", "", "100.64.0.9:5", 403},
|
||||
{"not a peer", locked, "/slots", "", "203.0.113.1:5", 403},
|
||||
{"a path that looks like an open route is still this route", locked, "/open/v1/models", "", "100.64.0.9:5", 403},
|
||||
{"a header naming another route does not change the gate", locked, "/v1/models", "open", "100.64.0.9:5", 403},
|
||||
{"admin-looking path is gated too (no admin on this listener)", locked, "/_crossbar/hosts", "", "100.64.0.9:5", 403},
|
||||
{"open route, anyone", open, "/v1/models", "", "203.0.113.1:5", 204},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
req := httptest.NewRequest(http.MethodGet, tc.path, nil)
|
||||
req.RemoteAddr = tc.addr
|
||||
if tc.hdr != "" {
|
||||
req.Header.Set("X-Crossbar-Route", tc.hdr)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
tc.h.ServeHTTP(rec, req)
|
||||
if rec.Code != tc.want {
|
||||
t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
+33
-3
@@ -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)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user