7.6 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(thefakeHealthhelper that v0'sproxy_test.goheld andrecorder_test.gostill 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.gofrom v0.1 — it must still pass. Itsproxy.New(cfg, h, nil)call no longer compiles, so this is the one given test you edit: change that call toproxy.New(cfg, h, nil, nil, nil, nil)and nothing else;Newmust accept nils forleases,lim,recand then behave like v0 (first healthy host, no queue, no recording). Say so in Deviations.
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):
- Route.
hdr := r.Header.Get(RouteHeader). Ifhdr != "": the path is used whole asrest(it must then start with/v1/or be/healthor/props); if the path also starts with a known route name and it differs fromhdr→ 400conflicting route. Ifhdr == "":SplitRouteas in v0 (400missing route). Unknown route (either source) → 404unknown route. Disallowedrest→ 404not found. - Peek (v0 rule): body up to
MaxBody→ 413;modelfrom the body or the route default;fp := fingerprint.Of(body)(GET/HEAD →""). - Lease.
host, reused, err := leases.Acquire(lease.Key{route, fp, model}, rt.Hosts, time.Now()).ErrNoHost→ 503no healthy host;ErrPinnedDown→ 503pinned host down. - Slot.
release, waited, err := lim.Acquire(r.Context(), host, model).ErrQueueFull→ 503queue full; ctx error → 499-style: just return (the client left; log status 499).defer release(). - Forward as in v0 (
Rewrite,FlushInterval: -1,ModifyResponsesetsHostHeaderandLeaseHeader,ErrorHandlermarks down + 502 with host). Tee: inModifyResponse, wrapresp.Bodyin a reader that passes every byte through unchanged and, whenContent-Typestarts withtext/event-stream, scans completedata:lines for a JSON object withusageand/ortimings, remembering the last one seen; for non-streamed JSON answers, remember the whole body'susage/timings(bounded: keep at most 1 MiB for the parse; beyond that, record no tokens). The scanner must not hold data back:Readreturns what the upstream returned. - 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 afterrp.ServeHTTPreturns, 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 whenServeHTTPreturns. Arecerror is logged, never returned to the client. - When
leases == nil(v0.1 compatibility path used byrecorder_test.go): choose with the v0Chooserule, skip the limiter and the store, still setHostHeader. - 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 ininternal/proxy/recorder_test.goas 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/.TestDifferentConversationsSpreadByFreeSlotsandTestQueueFullIs503are timing-based with generous margins; run-count=3. - 5. Run the gate.
make gate.cmd/crossbarwill not compile until task 06 — ifgo vet ./...fails only incmd/crossbar/main.gobecause of the newNewsignature, change that one call to passnil, nil, nilfor 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/isok;make gateprintsgate: ok;cmp internal/proxy/proxy_test.go docs/plans/v1/_files/internal/proxy/proxy_test.goprints nothing.
Stop and report if
TestAccountingRowsFromUsageAndTimingsfails on the token sums while the stream test passes: quote the recorded row.