From 33fa61bedb7ac6bc265df2c4d070eb43dc98cdde Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 17:44:00 -0700 Subject: [PATCH] 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) --- docs/implementer-log.md | 1 + internal/proxy/control.go | 31 ++++++ internal/proxy/control_test.go | 176 +++++++++++++++++++++++++++++++++ internal/proxy/forward.go | 28 +++++- internal/proxy/proxy.go | 26 +++-- internal/proxy/proxy_test.go | 4 +- 6 files changed, 251 insertions(+), 15 deletions(-) create mode 100644 internal/proxy/control.go create mode 100644 internal/proxy/control_test.go diff --git a/docs/implementer-log.md b/docs/implementer-log.md index aafd711..0854222 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/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. | ? | 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/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/proxy.go b/internal/proxy/proxy.go index 7614a13..d3e26c5 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 @@ -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 } } 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)