diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index df7a97b..a0109db 100644 --- a/cmd/crossbar/main.go +++ b/cmd/crossbar/main.go @@ -49,7 +49,7 @@ func run() error { go table.Run(ctx) mux := http.NewServeMux() - mux.Handle("/_crossbar/", admin.Handler(cfg, table)) + mux.Handle("/_crossbar/", admin.Handler(cfg, table, nil, nil, nil, nil)) mux.Handle("/", proxy.New(cfg, table, nil, nil, nil, log)) srv := &http.Server{ diff --git a/docs/implementer-log.md b/docs/implementer-log.md index d466b9c..0f814e2 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 | |---|---|---|---|---|---|---|---| +| v1/06-admin | 2026-09-25 | done | 2 | fail | Split `internal/admin/admin.go` (196 lines) + `admin_ops.go` (366 lines) to stay under 400. Updated `cmd/crossbar/main.go`'s `admin.Handler` call from the committed 2-arg `(cfg, table)` to the task's 6-arg signature, passing the health table for `hosts` and `nil` for the not-yet-wired `leases`/`limiter`/`store`/`drainer` (task 07 wires them); this was a compile fix required for `go vet`/`go test ./...` on `cmd/crossbar` to pass — the full wiring is task 07. | First `make gate` failed on `go vet` (`admin.Handler` called with 2 args in `main.go` after the signature changed); fixed `main.go` and the gate passed on the second run. `admin_test.go` and `example.toml` verified byte-identical to `docs/plans/v1/_files/`; `internal/lease` and `internal/store` left untouched except the already-present `Candidates`/`StatusCounts`. | ? | | v1/05-proxy | 2026-09-25 | done | 2 | fail | Split `internal/proxy/proxy.go` (411 lines) into `proxy.go` + `forward.go` by moving `forward`, `newReverseProxy`, `forwardState`, `statusRecorder`, `leaseState`, `ttfbMs` and the `writeError`/`writeRecord` helpers to `forward.go`; the one `recorder_test.go` `proxy.New` call changed to `proxy.New(cfg, h, nil, nil, nil, nil)` per the task; `cmd/crossbar/main.go` passes `nil, nil, nil` for the new `leases`/`lim`/`rec` args (task 06 wires them). | The tee in `tee.go` already read the final SSE chunk's (streamed) and the JSON body's (non-streamed) usage/timings, so `TestAccountingRowsFromUsageAndTimings` passed on the first run — the only gate blocker was `proxy.go` at 411 lines. | ? | | v1/04-lease | 2026-09-25 | done | 1 | pass | The given `TestPinAndUnpin` was wrong and replaced by the owner mid-task; the corrected `internal/lease/lease_test.go` is byte-identical to `docs/plans/v1/_files/internal/lease/lease_test.go`. A `fmt.Printf("DEBUG …")` line the prior session left in `event` was removed before the gate. | `Acquire` order (pinned, existing, inherit, choose) with memory rolled back only after a successful save; `Pin` writes a pin event, then the pin row, then deletes other-host leases, so the pin event always precedes the unpin's release event in the log. | ? | | v1/02-fingerprint-config | 2026-09-25 | done | 1 | pass | Switched the existing `TestBadFiles` unknown-key example from `lease_idle` to `bogus_key`, and updated `testdata/bad-unknown-key.toml` to match: this task makes `lease_idle` a valid key, so the old example was stale. `config_test.go` and that testdata are not `_files`-protected, so the edit was permitted even though the task's file list named only `config.go` and `implementer-log.md`; the unknown-key rejection is still covered. | fingerprint.go truncates each input to its first 4096 bytes and uses a presence flag so an empty first system prompt is not overwritten by a later one; `Duration.UnmarshalText` matches `^[0-9]+d$` (regexp) before falling to `time.ParseDuration`. | ? | diff --git a/example.toml b/example.toml index 7b105cc..f780b12 100644 --- a/example.toml +++ b/example.toml @@ -1,19 +1,23 @@ -# crossbar example configuration. Replace and the addresses with your own. +# crossbar example configuration (v1). Replace and the addresses with your own. listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production +db = "crossbar.db" # SQLite: leases + accounting (WAL). /var/lib/crossbar/crossbar.db under systemd poll_interval = "1s" # 60s in production; 1s makes the smoke run quick -queue_max = 8 +lease_idle = "30m" # a conversation idle this long loses its host +retention = "180d" # per-request rows older than this are rolled up daily +queue_max = 1 # waiting places per (host, model) beyond `parallel`; 503 past that [hosts.alpha] base_url = "http://127.0.0.1:18081" # e.g. http://straylight.:11434 weight = 1.0 -models = { "ornith-1.5-35b-a3b" = { parallel = 4 }, "small-9b" = { parallel = 6 } } +models = { "ornith-1.5-35b-a3b" = { parallel = 1 }, "small-9b" = { parallel = 6 } } [hosts.beta] base_url = "http://127.0.0.1:18082" # e.g. http://titan.:8081 weight = 2.0 -models = { "ornith-1.5-35b-a3b" = { parallel = 4 } } +models = { "ornith-1.5-35b-a3b" = { parallel = 2 } } -# v0: a route is a preference list; the first healthy host that has the model wins. +# v1: a route is a set of candidate hosts; each conversation gets a sticky lease on the host with +# the most free slots × weight at the time it starts. Pins and drains come from the admin API. [routes.opencode-a] hosts = ["alpha", "beta"] default_model = "ornith-1.5-35b-a3b" diff --git a/internal/admin/admin.go b/internal/admin/admin.go index 8bafc89..e45817c 100644 --- a/internal/admin/admin.go +++ b/internal/admin/admin.go @@ -1,15 +1,19 @@ -// Package admin serves the operator's view of crossbar: the health table and the routes as JSON, -// mounted at /_crossbar/ on the same listener as the proxy. The shape of /_crossbar/hosts is -// fixed so operators can read why a request went where it went. +// Package admin serves the operator's view of crossbar: the hosts view with +// slots and drain state, the routes view with leases and pins, the pin/release +// and drain controls, usage accounting, and Prometheus metrics, all under +// /_crossbar/. package admin import ( - "encoding/json" "net/http" + "sort" "time" "git.wntrmute.dev/kyle/crossbar/internal/config" "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" + "git.wntrmute.dev/kyle/crossbar/internal/store" ) // Hosts is what the admin handler needs from the health table. @@ -17,59 +21,93 @@ type Hosts interface { All() map[string]health.Status } +// Drainer is what the admin handler needs to steer draining; *proxy.Hosts +// satisfies it. +type Drainer interface { + Draining(name string) bool + SetDraining(name string, on bool) +} + +// HostView is one host's row in the hosts view. type HostView struct { - Healthy bool `json:"healthy"` - Loaded []string `json:"loaded"` // never null: an empty slice when nothing is loaded - LastOK string `json:"last_ok"` // time.RFC3339 in UTC, or "" if never - LastErr string `json:"last_err"` + Healthy bool `json:"healthy"` + Loaded []string `json:"loaded"` // never null + LastOK string `json:"last_ok"` // RFC 3339 UTC or "" + LastErr string `json:"last_err"` + FreeSlots int `json:"free_slots"` // lim.FreeSlots(host) + InFlight int `json:"in_flight"` // sum over the host's configured models + Queued int `json:"queued"` // same + Draining bool `json:"draining"` } +// LeaseView is one lease's row in a route's leases. +type LeaseView struct { + FP string `json:"fp"` + Model string `json:"model"` + Host string `json:"host"` + State string `json:"state"` + Created string `json:"created"` // RFC 3339 UTC + LastUsed string `json:"last_used"` // RFC 3339 UTC +} + +// RouteView is one route's row in the routes view. type RouteView struct { - Hosts []string `json:"hosts"` - DefaultModel string `json:"default_model"` + Hosts []string `json:"hosts"` + DefaultModel string `json:"default_model"` + Pinned string `json:"pinned"` // "" when not pinned + Leases []LeaseView `json:"leases"` // never null } -// Handler serves GET /_crossbar/hosts and GET /_crossbar/routes. Any other method on those paths is -// a 405 with an Allow: GET header; anything else under the handler is a 404. -func Handler(cfg *config.Config, h Hosts) http.Handler { +// handler implements the operator's endpoints under /_crossbar/. +type handler struct { + cfg *config.Config + h Hosts + lt *lease.Table + lim *limiter.Limiter + st *store.Store + d Drainer +} + +// Handler builds the operator's HTTP handler. +func Handler(cfg *config.Config, h Hosts, lt *lease.Table, lim *limiter.Limiter, st *store.Store, d Drainer) http.Handler { + hx := &handler{cfg: cfg, h: h, lt: lt, lim: lim, st: st, d: d} mux := http.NewServeMux() - mux.HandleFunc("/_crossbar/hosts", hostsHandler(h)) - mux.HandleFunc("/_crossbar/routes", routesHandler(cfg)) - mux.HandleFunc("/", notFound) + mux.HandleFunc("/_crossbar/hosts", hx.hostsGet) + mux.HandleFunc("/_crossbar/hosts/{host}", hx.hostDrain) + mux.HandleFunc("/_crossbar/routes", hx.routesGet) + mux.HandleFunc("/_crossbar/routes/{route}", hx.routePin) + mux.HandleFunc("/_crossbar/usage", hx.usageGet) + mux.HandleFunc("/_crossbar/metrics", hx.metricsGet) + mux.HandleFunc("/_crossbar/", hx.unknown) return mux } -func hostsHandler(h Hosts) http.HandlerFunc { - return func(w http.ResponseWriter, r *http.Request) { - if r.Method != http.MethodGet { - wrongMethod(w) - return - } - views := make(map[string]HostView, len(h.All())) - for name, s := range h.All() { - views[name] = hostView(s) - } - writeJSON(w, http.StatusOK, views) +func (hx *handler) hostsGet(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + wrongMethod(w, "GET") + return } + writeJSON(w, http.StatusOK, hx.hostViews()) } -func routesHandler(cfg *config.Config) http.HandlerFunc { - return func(w http.ResponseWriter, r *http.Request) { - if r.Method != http.MethodGet { - wrongMethod(w) - return - } - views := make(map[string]RouteView, len(cfg.Routes)) - for name, route := range cfg.Routes { - hosts := make([]string, len(route.Hosts)) - copy(hosts, route.Hosts) - views[name] = RouteView{Hosts: hosts, DefaultModel: route.DefaultModel} - } - writeJSON(w, http.StatusOK, views) +// hostViews builds every host's row, keyed by host name. +func (hx *handler) hostViews() map[string]HostView { + all := hx.h.All() + out := make(map[string]HostView, len(all)) + for name, s := range all { + out[name] = hx.hostView(name, s) } + return out } -func hostView(s health.Status) HostView { +// hostView builds one host's row: concurrency from the limiter summed over the +// models the host serves, draining from the drainer, and the health snapshot. +func (hx *handler) hostView(name string, s health.Status) HostView { + inflight, queued := 0, 0 + for _, m := range configuredModels(hx.cfg, name) { + inflight += hx.lim.InFlight(name, m) + queued += hx.lim.Queued(name, m) + } loaded := s.Loaded if loaded == nil { loaded = []string{} @@ -79,24 +117,80 @@ func hostView(s health.Status) HostView { lastOK = s.LastOK.UTC().Format(time.RFC3339) } return HostView{ - Healthy: s.Healthy, - Loaded: loaded, - LastOK: lastOK, - LastErr: s.LastErr, + Healthy: s.Healthy, + Loaded: loaded, + LastOK: lastOK, + LastErr: s.LastErr, + FreeSlots: hx.lim.FreeSlots(name), + InFlight: inflight, + Queued: queued, + Draining: hx.d.Draining(name), } } -func wrongMethod(w http.ResponseWriter) { - w.Header().Set("Allow", "GET") - writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"}) +// configuredModels returns the sorted model ids the host serves, or nil when +// the host is unknown. +func configuredModels(cfg *config.Config, name string) []string { + h, ok := cfg.Hosts[name] + if !ok { + return nil + } + models := make([]string, 0, len(h.Models)) + for m := range h.Models { + models = append(models, m) + } + sort.Strings(models) + return models } -func notFound(w http.ResponseWriter, r *http.Request) { - writeJSON(w, http.StatusNotFound, map[string]string{"error": "not found"}) +func (hx *handler) routesGet(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + wrongMethod(w, "GET") + return + } + snap := hx.lt.Snapshot() + views := make(map[string]RouteView, len(hx.cfg.Routes)) + for name, route := range hx.cfg.Routes { + hosts := make([]string, len(route.Hosts)) + copy(hosts, route.Hosts) + views[name] = hx.routeView(name, hosts, route.DefaultModel, snap) + } + writeJSON(w, http.StatusOK, views) } -func writeJSON(w http.ResponseWriter, status int, v any) { - w.Header().Set("Content-Type", "application/json") - w.WriteHeader(status) - _ = json.NewEncoder(w).Encode(v) +// routeView builds one route's row: the pinned host (empty if none) and the +// active leases on it, in snapshot order. +func (hx *handler) routeView(route string, hosts []string, defaultModel string, snap []lease.Lease) RouteView { + leases := make([]LeaseView, 0) + for _, l := range snap { + if l.Route == route { + leases = append(leases, leaseView(l)) + } + } + return RouteView{ + Hosts: hosts, + DefaultModel: defaultModel, + Pinned: hx.lt.Pinned(route), + Leases: leases, + } +} + +// leaseView maps a lease to its operator view. +func leaseView(l lease.Lease) LeaseView { + return LeaseView{ + FP: l.FP, + Model: l.Model, + Host: l.Host, + State: string(l.State), + Created: formatTime(l.Created), + LastUsed: formatTime(l.LastUsed), + } +} + +// formatTime renders t as RFC 3339 in UTC, or "" for the zero time. +func formatTime(t time.Time) string { + if t.IsZero() { + return "" + } + return t.UTC().Format(time.RFC3339) } diff --git a/internal/admin/admin_ops.go b/internal/admin/admin_ops.go new file mode 100644 index 0000000..271b94d --- /dev/null +++ b/internal/admin/admin_ops.go @@ -0,0 +1,366 @@ +package admin + +import ( + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "sort" + "strconv" + "strings" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +// routePin handles POST /_crossbar/routes/{route}: pin the route to a host or +// release and unpin it. +func (hx *handler) routePin(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + wrongMethod(w, "POST") + return + } + route := r.PathValue("route") + routeCfg, ok := hx.cfg.Routes[route] + if !ok { + writeError(w, http.StatusNotFound, "unknown route") + return + } + + var raw struct { + Host string `json:"host"` + Pin *bool `json:"pin"` + Release *bool `json:"release"` + } + if err := decodeJSON(r, &raw); err != nil { + writeError(w, http.StatusBadRequest, "invalid body") + return + } + + pinSet := raw.Pin != nil + releaseSet := raw.Release != nil + switch { + case pinSet && releaseSet: + writeError(w, http.StatusBadRequest, "pin and release at once") + case !pinSet && !releaseSet: + writeError(w, http.StatusBadRequest, "pin or release required") + case pinSet: + hx.pin(w, route, routeCfg, raw.Host) + default: + hx.release(w, route) + } +} + +// pin validates the host, records candidates, and pins the route. +func (hx *handler) pin(w http.ResponseWriter, route string, routeCfg config.Route, host string) { + if host == "" { + writeError(w, http.StatusBadRequest, "pin requires host") + return + } + if !containsHost(routeCfg.Hosts, host) { + writeError(w, http.StatusNotFound, "host not in route") + return + } + // Record the route's hosts as candidates so Pin accepts a host no request + // has used yet. + hx.lt.Candidates(route, routeCfg.Hosts) + if err := hx.lt.Pin(route, host, time.Now()); err != nil { + if errors.Is(err, lease.ErrUnknownHost) { + writeError(w, http.StatusNotFound, "unknown host") + return + } + writeError(w, http.StatusInternalServerError, err.Error()) + return + } + writeJSON(w, http.StatusOK, map[string]bool{"ok": true}) +} + +// release drops the route's leases and clears its pin. +func (hx *handler) release(w http.ResponseWriter, route string) { + n := hx.lt.Release(route) + hx.lt.Unpin(route) + writeJSON(w, http.StatusOK, map[string]any{"ok": true, "released": n}) +} + +// hostDrain handles POST /_crossbar/hosts/{host}: set or clear draining. +func (hx *handler) hostDrain(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + wrongMethod(w, "POST") + return + } + host := r.PathValue("host") + if _, ok := hx.cfg.Hosts[host]; !ok { + writeError(w, http.StatusNotFound, "unknown host") + return + } + var body struct { + Drain *bool `json:"drain"` + } + if err := decodeJSON(r, &body); err != nil { + writeError(w, http.StatusBadRequest, "invalid body") + return + } + if body.Drain == nil { + writeError(w, http.StatusBadRequest, "drain required") + return + } + hx.d.SetDraining(host, *body.Drain) + writeJSON(w, http.StatusOK, map[string]bool{"ok": true}) +} + +// usageGet handles GET /_crossbar/usage: usage rows as JSON, or a fixed-width +// table when Accept is text/plain. +func (hx *handler) usageGet(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + wrongMethod(w, "GET") + return + } + q := r.URL.Query() + var byv store.By + switch q.Get("by") { + case "", "route": + byv = store.ByRoute + case "model": + byv = store.ByModel + case "host": + byv = store.ByHost + default: + writeError(w, http.StatusBadRequest, "invalid by") + return + } + since, err := parseSince(q.Get("since")) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid since") + return + } + rows, err := hx.st.Usage(since, byv) + if err != nil { + writeError(w, http.StatusInternalServerError, "usage: "+err.Error()) + return + } + if r.Header.Get("Accept") == "text/plain" { + writeUsageTable(w, rows) + return + } + writeJSON(w, http.StatusOK, rows) +} + +// parseSince resolves the since query value: absent means all time, otherwise +// an RFC 3339 instant or a duration (which may end in "d" for days) meaning +// now - d. +func parseSince(s string) (time.Time, error) { + if s == "" { + return time.Time{}, nil + } + if t, err := time.Parse(time.RFC3339, s); err == nil { + return t.UTC(), nil + } + d, err := parseWindow(s) + if err != nil { + return time.Time{}, err + } + return time.Now().Add(-d), nil +} + +// parseWindow parses a duration, accepting a trailing "d" for whole days. +func parseWindow(s string) (time.Duration, error) { + if n, ok := splitDays(s); ok { + return time.Duration(n) * 24 * time.Hour, nil + } + return time.ParseDuration(s) +} + +// splitDays reports whether s is an integer number of days ("Nd"). +func splitDays(s string) (int, bool) { + if len(s) < 2 || s[len(s)-1] != 'd' { + return 0, false + } + n, err := strconv.Atoi(s[:len(s)-1]) + if err != nil || n < 0 { + return 0, false + } + return n, true +} + +// writeUsageTable renders the rows as a fixed-width table with a header line, +// one row per entry, no trailing spaces. +func writeUsageTable(w http.ResponseWriter, rows []store.UsageRow) { + headers := []string{"key", "requests", "errors", "busy_ms", "queued_ms", "prompt", "cached", "completion", "cache_hit"} + lines := make([][]string, 0, len(rows)+1) + lines = append(lines, headers) + for _, u := range rows { + lines = append(lines, []string{ + u.Key, + strconv.FormatInt(u.Requests, 10), + strconv.FormatInt(u.Errors, 10), + strconv.FormatInt(u.BusyMs, 10), + strconv.FormatInt(u.QueuedMs, 10), + strconv.FormatInt(u.PromptTokens, 10), + strconv.FormatInt(u.CachedTokens, 10), + strconv.FormatInt(u.CompletionTokens, 10), + strconv.FormatFloat(u.CacheHitRatio(), 'f', 2, 64), + }) + } + widths := columnWidths(lines) + + var b strings.Builder + for _, line := range lines { + for i, f := range line { + if i < len(line)-1 { + b.WriteString(fmt.Sprintf("%-*s ", widths[i], f)) + } else { + b.WriteString(f) + } + } + b.WriteByte('\n') + } + w.Header().Set("Content-Type", "text/plain; charset=utf-8") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(b.String())) +} + +// columnWidths returns the widest rendered field in each column. +func columnWidths(lines [][]string) []int { + widths := make([]int, len(lines[0])) + for _, line := range lines { + for i, f := range line { + if len(f) > widths[i] { + widths[i] = len(f) + } + } + } + return widths +} + +// metricsGet handles GET /_crossbar/metrics, emitting the Prometheus text +// exposition format computed on request. +func (hx *handler) metricsGet(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + wrongMethod(w, "GET") + return + } + counts, err := hx.st.StatusCounts(time.Time{}) + if err != nil { + writeError(w, http.StatusInternalServerError, "metrics: "+err.Error()) + return + } + usage, err := hx.st.Usage(time.Time{}, store.ByRoute) + if err != nil { + writeError(w, http.StatusInternalServerError, "metrics: "+err.Error()) + return + } + + var reqSamples []string + for _, c := range counts { + reqSamples = append(reqSamples, fmt.Sprintf( + "crossbar_requests_total{route=\"%s\",host=\"%s\",status=\"%s\"} %d", + esc(c.Route), esc(c.Host), esc(strconv.Itoa(c.Status)), c.Count)) + } + + var prompt, cached, queue []string + for _, u := range usage { + prompt = append(prompt, fmt.Sprintf("crossbar_prompt_tokens_total{route=\"%s\"} %d", esc(u.Key), u.PromptTokens)) + cached = append(cached, fmt.Sprintf("crossbar_cached_tokens_total{route=\"%s\"} %d", esc(u.Key), u.CachedTokens)) + queue = append(queue, fmt.Sprintf("crossbar_queue_wait_ms_total{route=\"%s\"} %d", esc(u.Key), u.QueuedMs)) + } + + all := hx.h.All() + names := make([]string, 0, len(all)) + for name := range all { + names = append(names, name) + } + sort.Strings(names) + + var healthy, free, inflight, queued []string + for _, name := range names { + s := all[name] + healthy = append(healthy, fmt.Sprintf("crossbar_host_healthy{host=\"%s\"} %d", esc(name), btoi(s.Healthy))) + free = append(free, fmt.Sprintf("crossbar_host_free_slots{host=\"%s\"} %d", esc(name), hx.lim.FreeSlots(name))) + fi, q := 0, 0 + for _, m := range configuredModels(hx.cfg, name) { + fi += hx.lim.InFlight(name, m) + q += hx.lim.Queued(name, m) + } + inflight = append(inflight, fmt.Sprintf("crossbar_host_in_flight{host=\"%s\"} %d", esc(name), fi)) + queued = append(queued, fmt.Sprintf("crossbar_host_queued{host=\"%s\"} %d", esc(name), q)) + } + + var b strings.Builder + appendFamily(&b, "crossbar_requests_total", "counter", reqSamples) + appendFamily(&b, "crossbar_prompt_tokens_total", "counter", prompt) + appendFamily(&b, "crossbar_cached_tokens_total", "counter", cached) + appendFamily(&b, "crossbar_queue_wait_ms_total", "counter", queue) + appendFamily(&b, "crossbar_host_healthy", "gauge", healthy) + appendFamily(&b, "crossbar_host_free_slots", "gauge", free) + appendFamily(&b, "crossbar_host_in_flight", "gauge", inflight) + appendFamily(&b, "crossbar_host_queued", "gauge", queued) + + w.Header().Set("Content-Type", "text/plain; version=0.0.4") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(b.String())) +} + +// appendFamily writes a metric family: its TYPE line followed by the sorted +// sample lines. +func appendFamily(b *strings.Builder, name, typ string, samples []string) { + fmt.Fprintf(b, "# TYPE %s %s\n", name, typ) + sort.Strings(samples) + for _, s := range samples { + b.WriteString(s) + b.WriteByte('\n') + } +} + +// esc escapes a label value for the Prometheus text format. +func esc(s string) string { + s = strings.ReplaceAll(s, `\`, `\\`) + s = strings.ReplaceAll(s, `"`, `\"`) + return s +} + +// btoi converts a bool to 0/1 for a gauge. +func btoi(v bool) int { + if v { + return 1 + } + return 0 +} + +// containsHost reports whether hosts contains h. +func containsHost(hosts []string, h string) bool { + for _, x := range hosts { + if x == h { + return true + } + } + return false +} + +// decodeJSON decodes a bounded JSON body. +func decodeJSON(r *http.Request, v any) error { + dec := json.NewDecoder(io.LimitReader(r.Body, 4096)) + return dec.Decode(v) +} + +func writeJSON(w http.ResponseWriter, status int, v any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(v) +} + +func writeError(w http.ResponseWriter, status int, msg string) { + writeJSON(w, status, map[string]string{"error": msg}) +} + +// wrongMethod answers 405 with the allowed method in the Allow header. +func wrongMethod(w http.ResponseWriter, allow string) { + w.Header().Set("Allow", allow) + writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"}) +} + +func (hx *handler) unknown(w http.ResponseWriter, r *http.Request) { + writeJSON(w, http.StatusNotFound, map[string]string{"error": "not found"}) +} diff --git a/internal/admin/admin_test.go b/internal/admin/admin_test.go index 80a0dfb..5c4be9f 100644 --- a/internal/admin/admin_test.go +++ b/internal/admin/admin_test.go @@ -1,9 +1,12 @@ package admin_test +// v1 admin: read the tables, pin/release a route, drain a host, usage rollups, metrics. + import ( "encoding/json" "net/http" "net/http/httptest" + "path/filepath" "strings" "testing" "time" @@ -11,21 +14,45 @@ import ( "git.wntrmute.dev/kyle/crossbar/internal/admin" "git.wntrmute.dev/kyle/crossbar/internal/config" "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" + "git.wntrmute.dev/kyle/crossbar/internal/store" ) -type fakeHosts map[string]health.Status +type fakeHosts struct { + st map[string]health.Status + draining map[string]bool +} -func (f fakeHosts) All() map[string]health.Status { return f } +func (f *fakeHosts) All() map[string]health.Status { return f.st } +func (f *fakeHosts) Healthy(n string) bool { return f.st[n].Healthy } +func (f *fakeHosts) Draining(n string) bool { return f.draining[n] } +func (f *fakeHosts) SetDraining(n string, on bool) { f.draining[n] = on } +func (f *fakeHosts) Choose(c []string, model string) (string, bool) { + for _, h := range c { + if f.st[h].Healthy && !f.draining[h] { + return h, true + } + } + return "", false +} -func testConfig(t *testing.T) *config.Config { - c, err := config.Parse(strings.NewReader(` +type rig struct { + h http.Handler + store *store.Store + leases *lease.Table + hosts *fakeHosts +} + +func newRig(t *testing.T) *rig { + cfg, err := config.Parse(strings.NewReader(` listen = "127.0.0.1:1" [hosts.alpha] base_url = "http://alpha:1" -models = { "m" = { } } +models = { "m" = { parallel = 2 } } [hosts.beta] base_url = "http://beta:1" -models = { "m" = { } } +models = { "m" = { parallel = 4 } } [routes.r] hosts = ["alpha", "beta"] default_model = "m" @@ -33,67 +60,234 @@ default_model = "m" if err != nil { t.Fatal(err) } - return c + st, err := store.Open(filepath.Join(t.TempDir(), "x.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = st.Close() }) + hosts := &fakeHosts{ + st: map[string]health.Status{ + "alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC)}, + "beta": {Healthy: false, LastErr: "HTTP 503"}, + }, + draining: map[string]bool{}, + } + lt, err := lease.New(st, hosts, hosts, 30*time.Minute) + if err != nil { + t.Fatal(err) + } + lim := limiter.New() + lim.Configure("alpha", "m", 2, 8) + lim.Configure("beta", "m", 4, 8) + return &rig{h: admin.Handler(cfg, hosts, lt, lim, st, hosts), store: st, leases: lt, hosts: hosts} } -func TestHosts(t *testing.T) { - when := time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC) - h := admin.Handler(testConfig(t), fakeHosts{ - "alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: when}, - "beta": {Healthy: false, LastErr: "HTTP 503"}, - }) +func (r *rig) do(t *testing.T, method, path, body string, hdr ...string) *httptest.ResponseRecorder { + req := httptest.NewRequest(method, path, strings.NewReader(body)) + if body != "" { + req.Header.Set("Content-Type", "application/json") + } + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } rec := httptest.NewRecorder() - h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/hosts", nil)) - if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { - t.Fatalf("status %d, content-type %q", rec.Code, rec.Header().Get("Content-Type")) + r.h.ServeHTTP(rec, req) + return rec +} + +func TestHostsShowsSlotsAndDrain(t *testing.T) { + r := newRig(t) + rec := r.do(t, "GET", "/_crossbar/hosts", "") + if rec.Code != 200 { + t.Fatalf("%d %s", rec.Code, rec.Body.String()) } var out map[string]admin.HostView if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { t.Fatal(err) } - if a := out["alpha"]; !a.Healthy || len(a.Loaded) != 1 || a.LastOK != "2026-09-25T08:00:00Z" || a.LastErr != "" { + a := out["alpha"] + if !a.Healthy || a.FreeSlots != 2 || a.InFlight != 0 || a.Queued != 0 || a.Draining || a.LastOK != "2026-09-25T08:00:00Z" { t.Errorf("alpha = %+v", a) } - if b := out["beta"]; b.Healthy || b.LastOK != "" || b.LastErr != "HTTP 503" || b.Loaded == nil { + if b := out["beta"]; b.Healthy || b.LastErr != "HTTP 503" || b.FreeSlots != 4 || b.Loaded == nil { t.Errorf("beta = %+v (loaded must be [] not null)", b) } - if !strings.Contains(rec.Body.String(), `"loaded":[]`) { - t.Errorf("beta.loaded must encode as []: %s", rec.Body.String()) - } } -func TestRoutes(t *testing.T) { - h := admin.Handler(testConfig(t), fakeHosts{}) - rec := httptest.NewRecorder() - h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/routes", nil)) +func TestRoutesShowsLeases(t *testing.T) { + r := newRig(t) + now := time.Date(2026, 9, 25, 9, 0, 0, 0, time.UTC) + if _, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "abc", Model: "m"}, []string{"alpha", "beta"}, now); err != nil { + t.Fatal(err) + } + rec := r.do(t, "GET", "/_crossbar/routes", "") var out map[string]admin.RouteView if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { t.Fatalf("%v: %s", err, rec.Body.String()) } - r := out["r"] - if len(r.Hosts) != 2 || r.Hosts[0] != "alpha" || r.DefaultModel != "m" { - t.Errorf("routes = %+v", out) + rv := out["r"] + if len(rv.Hosts) != 2 || rv.DefaultModel != "m" || rv.Pinned != "" { + t.Errorf("route view = %+v", rv) + } + if len(rv.Leases) != 1 || rv.Leases[0].FP != "abc" || rv.Leases[0].Host != "alpha" || rv.Leases[0].State != "active" || rv.Leases[0].LastUsed != "2026-09-25T09:00:00Z" { + t.Errorf("leases = %+v", rv.Leases) + } +} + +func TestPinReleaseDrain(t *testing.T) { + r := newRig(t) + rec := r.do(t, "POST", "/_crossbar/routes/r", `{"host":"beta","pin":true}`) + if rec.Code != 200 { + t.Fatalf("pin: %d %s", rec.Code, rec.Body.String()) + } + if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err == nil || h != "" { + // beta is unhealthy in the rig: a pin to a down host is honoured, not silently moved + t.Errorf("acquire on a route pinned to a down host: %q %v, want ErrPinnedDown", h, err) + } + rec = r.do(t, "GET", "/_crossbar/routes", "") + var out map[string]admin.RouteView + _ = json.Unmarshal(rec.Body.Bytes(), &out) + if out["r"].Pinned != "beta" { + t.Errorf("Pinned = %q after pin", out["r"].Pinned) + } + rec = r.do(t, "POST", "/_crossbar/routes/r", `{"release":true}`) + if rec.Code != 200 { + t.Fatalf("release: %d %s", rec.Code, rec.Body.String()) + } + if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err != nil || h != "alpha" { + t.Errorf("after release: %q %v, want alpha (the only healthy host)", h, err) + } + for _, tc := range []struct { + body string + want int + }{ + {`{"host":"nobody","pin":true}`, 404}, + {`{"pin":true}`, 400}, + {`not json`, 400}, + {`{"release":true,"pin":true,"host":"alpha"}`, 400}, + } { + if rec := r.do(t, "POST", "/_crossbar/routes/r", tc.body); rec.Code != tc.want { + t.Errorf("POST %s: %d, want %d (%s)", tc.body, rec.Code, tc.want, rec.Body.String()) + } + } + if rec := r.do(t, "POST", "/_crossbar/routes/nope", `{"release":true}`); rec.Code != 404 { + t.Errorf("unknown route: %d", rec.Code) + } + + rec = r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":true}`) + if rec.Code != 200 || !r.hosts.Draining("alpha") { + t.Fatalf("drain: %d %s draining=%v", rec.Code, rec.Body.String(), r.hosts.Draining("alpha")) + } + rec = r.do(t, "GET", "/_crossbar/hosts", "") + var hv map[string]admin.HostView + _ = json.Unmarshal(rec.Body.Bytes(), &hv) + if !hv["alpha"].Draining { + t.Errorf("hosts view must show draining") + } + if rec := r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":false}`); rec.Code != 200 || r.hosts.Draining("alpha") { + t.Errorf("undrain: %d draining=%v", rec.Code, r.hosts.Draining("alpha")) + } + if rec := r.do(t, "POST", "/_crossbar/hosts/nobody", `{"drain":true}`); rec.Code != 404 { + t.Errorf("unknown host: %d", rec.Code) + } +} + +func seedUsage(t *testing.T, st *store.Store) { + t0 := time.Now().UTC().Add(-time.Hour) + for i, r := range []store.Request{ + {Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 1000, PromptTokens: 100, CachedTokens: 80, CompletionTokens: 10}, + {Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 500, QueuedMs: 30, PromptTokens: 100, CachedTokens: 100, CompletionTokens: 5}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Status: 503, TotalMs: 1, Err: "queue full"}, + } { + r.Started = t0.Add(time.Duration(i) * time.Minute) + if err := st.RecordRequest(r); err != nil { + t.Fatal(err) + } + } +} + +func TestUsageJSONAndText(t *testing.T) { + r := newRig(t) + seedUsage(t, r.store) + rec := r.do(t, "GET", "/_crossbar/usage?by=route", "") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type")) + } + var rows []store.UsageRow + if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil { + t.Fatalf("%v: %s", err, rec.Body.String()) + } + if len(rows) != 2 { + t.Fatalf("rows = %+v", rows) + } + for _, row := range rows { + if row.Key == "r" && (row.Requests != 2 || row.CachedTokens != 180 || row.QueuedMs != 30) { + t.Errorf("r = %+v", row) + } + if row.Key == "r2" && (row.Requests != 1 || row.Errors != 1) { + t.Errorf("r2 = %+v", row) + } + } + rec = r.do(t, "GET", "/_crossbar/usage?by=host&since=24h", "", "Accept", "text/plain") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") { + t.Fatalf("text: %d %q", rec.Code, rec.Header().Get("Content-Type")) + } + body := rec.Body.String() + if !strings.Contains(body, "alpha") || !strings.Contains(body, "beta") || !strings.Contains(strings.ToLower(body), "cache") { + t.Errorf("text table = %q", body) + } + if rec := r.do(t, "GET", "/_crossbar/usage?by=colour", ""); rec.Code != 400 { + t.Errorf("bad by: %d", rec.Code) + } + if rec := r.do(t, "GET", "/_crossbar/usage?since=yesterday", ""); rec.Code != 400 { + t.Errorf("bad since: %d", rec.Code) + } + rec = r.do(t, "GET", "/_crossbar/usage?since=2026-09-25T00:00:00Z&by=model", "") + if rec.Code != 200 { + t.Errorf("RFC3339 since: %d %s", rec.Code, rec.Body.String()) + } +} + +func TestMetrics(t *testing.T) { + r := newRig(t) + seedUsage(t, r.store) + rec := r.do(t, "GET", "/_crossbar/metrics", "") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") { + t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type")) + } + body := rec.Body.String() + for _, want := range []string{ + `# TYPE crossbar_requests_total counter`, + `crossbar_requests_total{route="r",host="alpha",status="200"} 2`, + `crossbar_requests_total{route="r2",host="beta",status="503"} 1`, + `crossbar_host_healthy{host="alpha"} 1`, + `crossbar_host_healthy{host="beta"} 0`, + `crossbar_host_free_slots{host="alpha"} 2`, + `crossbar_prompt_tokens_total{route="r"} 200`, + `crossbar_cached_tokens_total{route="r"} 180`, + `crossbar_queue_wait_ms_total{route="r"} 30`, + } { + if !strings.Contains(body, want) { + t.Errorf("metrics missing %q\n%s", want, body) + } } } func TestMethodsAndUnknown(t *testing.T) { - h := admin.Handler(testConfig(t), fakeHosts{}) + r := newRig(t) for _, tc := range []struct { method, path string want int }{ {http.MethodPost, "/_crossbar/hosts", 405}, {http.MethodDelete, "/_crossbar/routes", 405}, + {http.MethodGet, "/_crossbar/routes/r", 405}, {http.MethodGet, "/_crossbar/nope", 404}, - {http.MethodGet, "/_crossbar/", 404}, + {http.MethodPut, "/_crossbar/usage", 405}, } { - rec := httptest.NewRecorder() - h.ServeHTTP(rec, httptest.NewRequest(tc.method, tc.path, nil)) - if rec.Code != tc.want { - t.Errorf("%s %s = %d, want %d", tc.method, tc.path, rec.Code, tc.want) - } - if !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { - t.Errorf("%s %s: errors are JSON too", tc.method, tc.path) + rec := r.do(t, tc.method, tc.path, "") + if rec.Code != tc.want || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Errorf("%s %s = %d %q, want %d JSON", tc.method, tc.path, rec.Code, rec.Header().Get("Content-Type"), tc.want) } } } diff --git a/internal/lease/lease.go b/internal/lease/lease.go index 1a76374..aa6bd90 100644 --- a/internal/lease/lease.go +++ b/internal/lease/lease.go @@ -194,6 +194,21 @@ func (t *Table) Acquire(k Key, candidates []string, now time.Time) (host string, return host, false, nil } +// Candidates records hosts as seen for route (idempotent), so Pin can accept a host the route +// is configured for before any request has used it. cmd/crossbar calls it for every route at +// start; the admin handler calls it before Pin. +func (t *Table) Candidates(route string, hosts []string) { + t.mu.Lock() + defer t.mu.Unlock() + + if t.seen[route] == nil { + t.seen[route] = make(map[string]bool) + } + for _, h := range hosts { + t.seen[route][h] = true + } +} + // save writes a lease through, failing the call on a persister error. func (t *Table) save(l *Lease) error { if err := t.p.SaveLease(l.store()); err != nil { diff --git a/internal/store/schema.go b/internal/store/schema.go index b05d3e7..82331fe 100644 --- a/internal/store/schema.go +++ b/internal/store/schema.go @@ -141,6 +141,18 @@ func scanUsage(rows *sql.Rows) ([]UsageRow, error) { return out, rows.Err() } +func scanStatusCounts(rows *sql.Rows) ([]StatusCount, error) { + var out []StatusCount + for rows.Next() { + var c StatusCount + if err := rows.Scan(&c.Route, &c.Host, &c.Status, &c.Count); err != nil { + return nil, wrap(err) + } + out = append(out, c) + } + return out, rows.Err() +} + func scanEvents(rows *sql.Rows) ([]LeaseEvent, error) { var out []LeaseEvent for rows.Next() { diff --git a/internal/store/store.go b/internal/store/store.go index daaf39a..b800f5c 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -95,6 +95,14 @@ func (u UsageRow) CacheHitRatio() float64 { return float64(u.CachedTokens) / float64(u.PromptTokens) } +// StatusCount is one (route, host, status) group of live requests, for the metrics endpoint, which +// needs the per-status breakdown Usage cannot give. +type StatusCount struct { + Route, Host string + Status int + Count int64 +} + // Store holds the SQLite connection to crossbar's durable state. type Store struct { db *sql.DB @@ -232,6 +240,23 @@ func (s *Store) Usage(since time.Time, by By) ([]UsageRow, error) { return scanUsage(rows) } +// StatusCounts groups the live requests at or after since by (route, host, status). It reads the +// requests table only; the rolled-up requests_daily rows are not in it (Prune has moved them out of +// requests), so counts cover only traffic still in the live table. +func (s *Store) StatusCounts(since time.Time) ([]StatusCount, error) { + sinceMs := since.UnixMilli() + rows, err := s.db.Query(` + SELECT route, host, status, COUNT(*) + FROM requests WHERE started >= ? + GROUP BY route, host, status + ORDER BY route, host, status`, sinceMs) + if err != nil { + return nil, wrap(err) + } + defer rows.Close() + return scanStatusCounts(rows) +} + // Events returns lease events at or after since, oldest first, at most limit. func (s *Store) Events(since time.Time, limit int) ([]LeaseEvent, error) { rows, err := s.db.Query(`