125 lines
7.3 KiB
Markdown
125 lines
7.3 KiB
Markdown
# 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.
|
|
|
|
## Interfaces
|
|
|
|
`internal/proxy`, package `proxy` (v0 names kept; additions):
|
|
|
|
```go
|
|
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: lim.FreeSlots(host), Queued: lim.Queued(host, model), Weight}.
|
|
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`.
|
|
|
|
```sh
|
|
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.
|