Control-plane requests follow the lease but take no slot and write no row

Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
2026-09-25 17:44:00 -07:00
parent 3518e84dd7
commit 33fa61bedb
6 changed files with 251 additions and 15 deletions
+31
View File
@@ -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
}
+176
View File
@@ -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
}
+24 -4
View File
@@ -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:…<host>" 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:…<host>" 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,
+16 -10
View File
@@ -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
@@ -208,15 +208,14 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
return
}
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)
started := time.Now()
@@ -227,7 +226,7 @@ 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
}
@@ -252,13 +251,20 @@ 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) {
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. 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 {
@@ -298,7 +304,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 +322,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
}
}
+3 -1
View File
@@ -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)