Admin: leases, pin, release, drain, usage, metrics

cmd/crossbar/main.go calls admin.Handler with the new 6-arg signature, passing nil for the not-yet-wired leases/limiter/store/drainer (task 07 wires them) so go vet and go test ./... pass on cmd/crossbar. This is a compile fix, not the task-07 wiring; noted in the implementer-log.

Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
2026-09-25 06:33:29 -07:00
parent c2a8b88a3f
commit 32ac7f549a
9 changed files with 808 additions and 97 deletions
+1 -1
View File
@@ -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{
+1
View File
@@ -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`. | ? |
+9 -5
View File
@@ -1,19 +1,23 @@
# crossbar example configuration. Replace <tailnet> and the addresses with your own.
# crossbar example configuration (v1). Replace <tailnet> 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.<tailnet>: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.<tailnet>: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"
+148 -54
View File
@@ -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)
}
+366
View File
@@ -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"})
}
+231 -37
View File
@@ -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)
}
}
}
+15
View File
@@ -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 {
+12
View File
@@ -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() {
+25
View File
@@ -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(`