Files

7.7 KiB

v1 task 05: the proxy uses leases, the limiter, the SSE tee and the store

Branch: v1 (run git switch v1; git status --short must be empty, otherwise stop) Commit subject: Route by lease, queue per host and model, record every request

Goal

internal/proxy becomes the v1 proxy: route from path or header, fingerprint the body, acquire a lease, take a limiter slot (queue or 503), forward with streaming, tee the response to read usage/timings, mark hosts down on failure, and record one store.Request per request. PLAN.md §4, §6, §7a.

Context

The v0 proxy stays the skeleton of this one: SplitRoute, the ordered error answers, the model peek, httputil.ReverseProxy with FlushInterval: -1, the status recorder with a checked Flush, the log line. What changes is who picks the host and what happens around the forward. Two new response headers make the behaviour observable: X-Crossbar-Host (existing) and X-Crossbar-Lease: new|reused. The given test drives the whole handler over real HTTP with a real store, lease.Table, limiter and health.Table.

Files

  • Copy (replaces v0's file): internal/proxy/proxy_test.go
  • Copy: internal/proxy/helpers_test.go (the fakeHealth helper that v0's proxy_test.go held and recorder_test.go still needs)
  • Modify: internal/proxy/proxy.go (split into more files if it passes 400 lines: proxy.go, tee.go, hosts.go), docs/implementer-log.md
  • Keep: internal/proxy/recorder_test.go from v0.1 — it must still pass. Its proxy.New(cfg, h, nil) call no longer compiles, so this is the one given test you edit: change that call to proxy.New(cfg, h, nil, nil, nil, nil) and nothing else; New must accept nils for leases, lim, rec and then behave like v0 (first healthy host, no queue, no recording). Say so in Deviations. The edited file is the plan's reference copy at _files/internal/proxy/recorder_test.go for the reviewer's byte-exact check.

Interfaces

internal/proxy, package proxy (v0 names kept; additions):

const (
    MaxBody     = 16 << 20
    HostHeader  = "X-Crossbar-Host"
    LeaseHeader = "X-Crossbar-Lease"   // "new" or "reused"
    RouteHeader = "X-Crossbar-Route"   // client may name the route here instead of the path
)
type Health interface { Get(name string) (health.Status, bool); MarkDown(name, reason string) }
type Recorder interface { RecordRequest(store.Request) error }   // *store.Store satisfies it

// Hosts adapts the health table and config for the lease table, and holds the drain set.
type Hosts struct { /* private */ }
func HostView(h *health.Table, cfg *config.Config) *Hosts
func (h *Hosts) Healthy(name string) bool
func (h *Hosts) Draining(name string) bool
func (h *Hosts) SetDraining(name string, on bool)

// Chooser adapts config, health and limiter to lease.Chooser using choose.Best:
// Info{Healthy, Draining: false (the lease table already filtered), Loaded: model in Loaded,
// CanServe: cfg.Serves, Free: free slots FOR THIS MODEL on this host =
// cfg.Hosts[host].Models[model].Parallel - lim.InFlight(host, model) (never below 0; 0 when the
// host does not list the model), Queued: lim.Queued(host, model), Weight}.
// (Corrected 2026-09-25: an earlier version said lim.FreeSlots(host), which sums every model's
// slots and let a host win on slots the requested model cannot use.)
func Chooser(cfg *config.Config, h *health.Table, l *limiter.Limiter) lease.Chooser

func New(cfg *config.Config, h Health, leases *lease.Table, lim *limiter.Limiter, rec Recorder, log *slog.Logger) *Handler
func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request)
func SplitRoute(path string) (route, rest string, ok bool)

ServeHTTP, in order (every error answer is JSON {"error":"…"} as in v0):

  1. Route. hdr := r.Header.Get(RouteHeader). If hdr != "": the path is used whole as rest (it must then start with /v1/ or be /health or /props); if the path also starts with a known route name and it differs from hdr → 400 conflicting route. If hdr == "": SplitRoute as in v0 (400 missing route). Unknown route (either source) → 404 unknown route. Disallowed rest → 404 not found.
  2. Peek (v0 rule): body up to MaxBody → 413; model from the body or the route default; fp := fingerprint.Of(body) (GET/HEAD → "").
  3. Lease. host, reused, err := leases.Acquire(lease.Key{route, fp, model}, rt.Hosts, time.Now()). ErrNoHost → 503 no healthy host; ErrPinnedDown → 503 pinned host down.
  4. Slot. release, waited, err := lim.Acquire(r.Context(), host, model). ErrQueueFull → 503 queue full; ctx error → 499-style: just return (the client left; log status 499). defer release().
  5. Forward as in v0 (Rewrite, FlushInterval: -1, ModifyResponse sets HostHeader and LeaseHeader, ErrorHandler marks down + 502 with host). Tee: in ModifyResponse, wrap resp.Body in a reader that passes every byte through unchanged and, when Content-Type starts with text/event-stream, scans complete data: lines for a JSON object with usage and/or timings, remembering the last one seen; for non-streamed JSON answers, remember the whole body's usage/timings (bounded: keep at most 1 MiB for the parse; beyond that, record no tokens). The scanner must not hold data back: Read returns what the upstream returned.
  6. Record, after the upstream body is closed (the tee's Close, or the error handler): store.Request{Route, FP: fp, Model, Host, Started, QueuedMs: waited, TTFBMs (first byte of the response head), TotalMs, Status, Streamed, PromptTokens: usage.prompt_tokens (or timings.prompt_n), CachedTokens: timings.cache_n, CompletionTokens: usage.completion_tokens (or timings.predicted_n), Err}. Also record the 503/502 cases (Status set, no tokens). Do it from the request goroutine after rp.ServeHTTP returns, so tests that read the store right after the response see the row; if the tee cannot tell that the body closed, record what you have when ServeHTTP returns. A rec error is logged, never returned to the client.
  7. When leases == nil (v0.1 compatibility path used by recorder_test.go): choose with the v0 Choose rule, skip the limiter and the store, still set HostHeader.
  8. Log line as v0, adding lease=new|reused, queued_ms, fp (first 8 hex chars only). Never the body.

Steps

  • 1. Copy (replace). git switch v1; cp docs/plans/v1/_files/internal/proxy/proxy_test.go internal/proxy/proxy_test.go; cp docs/plans/v1/_files/internal/proxy/helpers_test.go internal/proxy/. Edit the one call in internal/proxy/recorder_test.go as described above.
  • 2. See it fail (compile). 3. Write the code. gofmt -w internal/proxy/.
  • 4. See it pass. go test -race -count=1 ./internal/proxy/. TestDifferentConversationsSpreadByFreeSlots and TestQueueFullIs503 are timing-based with generous margins; run -count=3.
  • 5. Run the gate. make gate. cmd/crossbar will not compile until task 06 — if go vet ./... fails only in cmd/crossbar/main.go because of the new New signature, change that one call to pass nil, nil, nil for the new arguments (task 06 wires it properly) and say so in Deviations.
  • 6. Log and commit. Row v1/05-proxy.
git add internal/proxy cmd/crossbar docs/implementer-log.md
git commit

Done when

  • go test -race -count=3 ./internal/proxy/ is ok; make gate prints gate: ok; cmp internal/proxy/proxy_test.go docs/plans/v1/_files/internal/proxy/proxy_test.go prints nothing.

Stop and report if

  • TestAccountingRowsFromUsageAndTimings fails on the token sums while the stream test passes: quote the recorded row.