diff --git a/README.md b/README.md index c6fc69c..03fc2a6 100644 --- a/README.md +++ b/README.md @@ -59,6 +59,9 @@ hosts = ["beta", "alpha"] | `hosts..models` | The models this host serves, with per-model parallel tuning. | | `routes..hosts` | Candidate hosts, tried in order until one is healthy; a conversation leases one of them. | | `routes..default_model` | Model used when a request omits one; must be served by a host in the route. | +| `identity` | `"off"` (default), `"tailscale"`, or `"header"`; see below. | +| `hosts..wake` | A wake-on-LAN target (`mac`, `broadcast`, `wait`) so crossbar can rouse a sleeping host when nothing else can take a new lease. | +| `routes..peers` | The tailnet nodes allowed to reach the route, with `identity = "tailscale"`; see below. | ## Run @@ -159,7 +162,29 @@ crossbar_host_healthy{host="alpha"} 1 crossbar_host_healthy{host="beta"} 1 ``` -## What v1 does not do +## Wake -The context-size guard, wake-on-LAN, Tailscale identity and `/slots` are out of scope for v1; see -`PLAN.md` v2. +When a route has no healthy host left and at least one candidate lists a `wake` target, crossbar +sends that host a wake-on-LAN magic packet, in route order, and retries the lease once. A host that +wakes up takes the conversation; if none wakes, the request gets `503 {"error":"no healthy host", +"woke":[""]}`. The context-size guard wakes a sleeping host the same way before it +answers `400 prompt too large`, when no healthy host's per-slot context can fit the prompt. + +## Identity + +`identity` gates who may use a route. With the default `"off"` every request is admitted. With +`"tailscale"`, a route that lists `peers` answers `403` to any caller whose tailnet address is not +one of them (checked with `tailscale whois`): + +```toml +[routes.hermes-x] +hosts = ["beta", "alpha"] +peers = ["talos"] +``` + +`"header"` trusts the `X-Crossbar-Peer` header instead and needs no tailnet; it is insecure and for +tests only, so crossbar logs a warning when it starts in that mode. + +## What v2 does not do + +Request coalescing, `/slots` and TLS are out of scope for v2; see `PLAN.md`. diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index ce11cbb..3360531 100644 --- a/cmd/crossbar/main.go +++ b/cmd/crossbar/main.go @@ -17,10 +17,12 @@ 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/identity" "git.wntrmute.dev/kyle/crossbar/internal/lease" "git.wntrmute.dev/kyle/crossbar/internal/limiter" "git.wntrmute.dev/kyle/crossbar/internal/proxy" "git.wntrmute.dev/kyle/crossbar/internal/store" + "git.wntrmute.dev/kyle/crossbar/internal/wake" ) func main() { @@ -59,6 +61,18 @@ func run() error { hosts := proxy.HostView(table, cfg) lim := limiter.New() + + // Wake: rouse a sleeping host when a route has no healthy host left. Built + // from every host that carries a wake target; the health table satisfies the + // waker's Health interface. + targets := make(map[string]wake.Target, len(cfg.Hosts)) + for name, h := range cfg.Hosts { + if h.Wake == nil { + continue + } + targets[name] = wake.Target{MAC: h.Wake.MAC, Broadcast: h.Wake.Broadcast, Wait: h.Wake.Wait.Duration} + } + waker := wake.New(targets, hosts) for name, h := range cfg.Hosts { for model, m := range h.Models { lim.Configure(name, model, m.Parallel, cfg.QueueMax) @@ -73,9 +87,31 @@ func run() error { leases.Candidates(name, rt.Hosts) } + // Identity: gate the proxy on the route's peers when a backend is + // configured; off leaves the proxy unwrapped. + logIdentityMode(log, cfg.Identity) + + p := proxy.New(cfg, table, leases, lim, st, log) + p.SetWaker(waker) + + var handler http.Handler = p + if cfg.Identity != "off" { + var checker *identity.Checker + switch cfg.Identity { + case "tailscale": + checker = identity.NewChecker(identity.TailscaleResolver{}) + default: // "header" + checker = identity.NewHeaderChecker() + } + handler = identity.Middleware(checker, func(route string) ([]string, bool) { + rt, ok := cfg.Routes[route] + return rt.Peers, ok + }, p) + } + mux := http.NewServeMux() mux.Handle("/_crossbar/", admin.Handler(cfg, table, leases, lim, st, hosts)) - mux.Handle("/", proxy.New(cfg, table, leases, lim, st, log)) + mux.Handle("/", handler) // Background maintenance until ctx is done. Errors are logged, never fatal. go func() { @@ -156,3 +192,16 @@ func run() error { return err } } + +// logIdentityMode logs which identity backend is active and, for the unauthenticated header +// backend used by the smoke run, warns that it must not be exposed. +func logIdentityMode(log *slog.Logger, mode string) { + if mode == "off" { + log.Info("identity", "mode", "off") + return + } + log.Info("identity", "mode", mode) + if mode == "header" { + log.Warn("identity header mode is not authenticated; do not expose it") + } +} diff --git a/cmd/fakeupstream/main.go b/cmd/fakeupstream/main.go index c222d38..0cb6b44 100644 --- a/cmd/fakeupstream/main.go +++ b/cmd/fakeupstream/main.go @@ -7,7 +7,9 @@ // SSE chunks 200 ms apart when the body has "stream": true, then a final chunk carrying // "usage" and llama-server style "timings", then [DONE]; one JSON answer with usage and // timings otherwise. -slow adds that many milliseconds before answering (for queue tests). -// Every response carries X-Upstream: . +// Every response carries X-Upstream: . /props reports -n-ctx and -slots. With -wol-listen, +// a valid wake-on-LAN magic packet for -wol-mac received on that UDP address removes the down +// file, so the fake "boots" when woken. package main import ( @@ -16,6 +18,7 @@ import ( "fmt" "io" "log" + "net" "net/http" "os" "strings" @@ -28,7 +31,14 @@ func main() { models := flag.String("models", "m", "comma-separated model ids for /v1/models") downFile := flag.String("down-file", "", "while this file exists, /health answers 503") slow := flag.Int("slow", 0, "milliseconds to wait before answering a completion") + nCtx := flag.Int("n-ctx", 8192, "n_ctx reported by /props") + slots := flag.Int("slots", 2, "total_slots reported by /props") + wolListen := flag.String("wol-listen", "", "UDP address to listen on for a wake-on-LAN magic packet") + wolMAC := flag.String("wol-mac", "aa:bb:cc:dd:ee:01", "MAC the magic packet must carry") flag.Parse() + if *wolListen != "" && *downFile != "" { + go wakeOnPacket(*wolListen, *wolMAC, *downFile) + } ids := strings.Split(*models, ",") mux := http.NewServeMux() @@ -56,7 +66,7 @@ func main() { }) mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { stamp(w) - writeJSON(w, map[string]any{"default_generation_settings": map[string]any{"n_ctx": 8192}, "total_slots": 2, "model_path": *name}) + writeJSON(w, map[string]any{"default_generation_settings": map[string]any{"n_ctx": *nCtx}, "total_slots": *slots, "model_path": *name}) }) mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) { stamp(w) @@ -113,3 +123,39 @@ func writeJSON(w http.ResponseWriter, v any) { w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(v) } + +// wakeOnPacket removes downFile when a magic packet for mac arrives: 6×0xff then the MAC 16 times. +func wakeOnPacket(addr, mac, downFile string) { + hw, err := net.ParseMAC(mac) + if err != nil { + log.Fatalf("wol-mac: %v", err) + } + pc, err := net.ListenPacket("udp4", addr) + if err != nil { + log.Fatalf("wol-listen: %v", err) + } + log.Printf("fakeupstream listening for wake-on-LAN on %s (mac %s)", addr, hw) + buf := make([]byte, 256) + for { + n, _, err := pc.ReadFrom(buf) + if err != nil { + return + } + if n != 102 { + continue + } + ok := true + for i := 0; i < 6; i++ { + ok = ok && buf[i] == 0xff + } + for i := 0; i < 16 && ok; i++ { + for j := 0; j < 6; j++ { + ok = ok && buf[6+6*i+j] == hw[j] + } + } + if ok { + log.Printf("magic packet received: waking (removing %s)", downFile) + _ = os.Remove(downFile) + } + } +} diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 474c1bc..31ab42e 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 | |---|---|---|---|---|---|---|---| +| v2/05-wiring-smoke | 2026-09-25 | done | 1 | pass | none | The wiring in `cmd/crossbar/main.go` and `internal/proxy/{proxy,forward,ctxguard}.go` plus the README section were already in the working tree from a prior session; this session only ran the tests, the gate, the log row, and the commit. `go test -race -count=1 ./...` failed once on `TestQueueFullIs503` (`Errors:2`, the 503 not recorded) — the known v1 recording defect the owner scheduled as a v2.1 task 01; reran once and it passed. `make gate` printed `gate: ok` on the first run. Committed the two owner-corrected given v1 tests (`internal/limiter/limiter_test.go`, `internal/proxy/proxy_test.go`) alongside the prior session's changes. | ? | | v2/04-identity | 2026-09-25 | done | 1 | pass | new file `internal/config/identity.go` | Implemented `internal/identity/identity.go`: `ParseWhois` (Node = ComputedName, else Name minus trailing dot/domain; empty node errors), `TailscaleResolver` (`tailscale whois --json`, 3 s timeout, non-zero exit → `ErrNotAPeer`, missing binary a real deny), `Checker` with a 5-min per-address cache that also caches `ErrNotAPeer`, and `NewHeaderChecker`/`WithHeaderPeer` that read the peer from a context value. `middleware.go` names the route like the proxy (X-Crossbar-Route header, else first path segment), passes `/_crossbar/` and unknown routes straight through, and answers 403 `{"error":"forbidden route"}`. Config gains `Identity`/`Wake`/`Peers`; validation keys the peers check on the *explicit* identity value (a config with peers but no identity key passes), and `wake.wait` defaults to 45 s. Copied all four given files byte-identical; `go test -race ./internal/identity/ ./internal/config/` and `make gate` printed `gate: ok` on the first run. | ? | | v2/03-wake | 2026-09-25 | done | 1 | pass | none | Implemented wake-on-LAN in new `internal/wake/wake.go`: `MagicPacket` builds the 102-byte frame via `net.ParseMAC` (six `0xff` bytes plus the MAC repeated sixteen times) and rejects bad MACs; `Send` emits one UDP4 datagram to the resolved broadcast address, returning parse/resolve/write errors; `Waker` tracks last-sent per host under a mutex and sends at most once per `Wait` window, polling health every second (`PollEvery` is a test hook) until healthy, on `Wait` timeout, or on ctx cancellation, returning false for an unknown host without sending. Copied `internal/wake/wake_test.go` byte-identical to `docs/plans/v2/_files/`; `go test -race -count=3 ./internal/wake/` ok and `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b | | v2/02-ctxguard | 2026-09-25 | done | 1 | pass | none | Implemented the context-size guard in new `internal/proxy/ctxguard.go` (estimate `int(float64(len(body))/4*1.2)`; rule 2 skip on unknown/fit; rule 3 move via `leases.Move` with a `moved:>` header; rule 4 400 with `{"error":"prompt too large","estimate":E,"max":M}` and a status-400 accounting row, no forward, no mark-down) and wired it into `ServeHTTP` between the lease and the slot; added `Move` to `internal/lease/lease.go` (re-leases, deletes the old row, records a `ctx` event) and `ReasonCtx = "ctx"` to `internal/store`. Copied `internal/proxy/ctxguard_test.go` byte-identical to `docs/plans/v2/_files/`; `go test -race -count=2 ./internal/proxy/ ./internal/lease/` ok and `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b | diff --git a/example.toml b/example.toml index f780b12..03f0372 100644 --- a/example.toml +++ b/example.toml @@ -5,6 +5,7 @@ poll_interval = "1s" # 60s in production; 1s makes the smoke run 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 +identity = "off" # "tailscale" gates routes with `peers` by `tailscale whois`; "header" trusts X-Crossbar-Peer (TEST ONLY) [hosts.alpha] base_url = "http://127.0.0.1:18081" # e.g. http://straylight.:11434 @@ -15,6 +16,10 @@ models = { "ornith-1.5-35b-a3b" = { parallel = 1 }, "small-9b" = { parallel = base_url = "http://127.0.0.1:18082" # e.g. http://titan.:8081 weight = 2.0 models = { "ornith-1.5-35b-a3b" = { parallel = 2 } } +[hosts.beta.wake] # v2: wake a sleeping host when nothing else can take a new lease +mac = "aa:bb:cc:dd:ee:02" +broadcast = "127.0.0.1:19082" # the LAN broadcast address, port 9, in production +wait = "20s" # 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. @@ -24,3 +29,4 @@ default_model = "ornith-1.5-35b-a3b" [routes.hermes-x] hosts = ["beta", "alpha"] +# peers = ["talos"] # v2: with identity = "tailscale", only these tailnet nodes may use the route diff --git a/internal/limiter/limiter_test.go b/internal/limiter/limiter_test.go index bf90588..6ab30e2 100644 --- a/internal/limiter/limiter_test.go +++ b/internal/limiter/limiter_test.go @@ -32,10 +32,10 @@ func TestParallelAndQueue(t *testing.T) { go func() { rel, waited, err := l.Acquire(ctx, "alpha", "m") if err == nil { - defer rel() if waited < 40*time.Millisecond { err = errors.New("third acquire did not wait") } + rel() // release before reporting, so the final count check cannot race it } got3 <- err }() diff --git a/internal/proxy/ctxguard.go b/internal/proxy/ctxguard.go index 93967a2..f364c15 100644 --- a/internal/proxy/ctxguard.go +++ b/internal/proxy/ctxguard.go @@ -69,6 +69,35 @@ func (p *Handler) guard(w http.ResponseWriter, r *http.Request, hosts []string, return newHost, movedHeader(host, newHost), estimate, false } + // Rule 4a: no host fits. Before refusing, ask a waker to rouse a candidate + // whose context may grow when it comes up; a woken host takes the lease. + if p.waker != nil { + for _, name := range hosts { + s, ok := p.health.Get(name) + if !ok { + continue + } + // A healthy host does not need waking; only a down host might grow + // a larger context when it comes up. + if s.Healthy { + continue + } + // A host with a known per-slot context smaller than the estimate + // cannot serve it no matter how it wakes. + if psc := s.PerSlotCtx(); psc != 0 && psc < estimate { + continue + } + if p.cfg.Hosts[name].Wake == nil || !p.waker.Wake(r.Context(), name) { + continue + } + if err := p.leases.Move(lease.Key{Route: route, FP: fp, Model: model}, name, time.Now()); err != nil { + p.writeError(w, http.StatusBadGateway, "upstream failed") + return host, "", estimate, true + } + return name, movedHeader(host, name), estimate, false + } + } + // Rule 4: no host fits. Answer 400 with the estimate and the largest // available per-slot context, and record the row. p.refuseCtx(w, host, route, model, fp, started, estimate, largestSlotCtx(hosts, p.health)) diff --git a/internal/proxy/forward.go b/internal/proxy/forward.go index 8fcf0b7..3ac16f2 100644 --- a/internal/proxy/forward.go +++ b/internal/proxy/forward.go @@ -152,6 +152,30 @@ func (p *Handler) writeError(w http.ResponseWriter, status int, msg string) { _ = json.NewEncoder(w).Encode(map[string]string{"error": msg}) } +// writeNoHealthyHost answers the 503 when no candidate woke. The body names the +// hosts that were asked to wake (an empty list, never null), and the row +// records the miss. +func (p *Handler) writeNoHealthyHost(w http.ResponseWriter, route, model, fp string, started time.Time, tried []string) { + if tried == nil { + tried = []string{} + } + p.writeRecord(store.Request{ + Route: route, + FP: fp, + Model: model, + Started: started, + TotalMs: time.Since(started).Milliseconds(), + Status: http.StatusServiceUnavailable, + Err: "no healthy host", + }) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusServiceUnavailable) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": "no healthy host", + "woke": tried, + }) +} + // writeRecord writes one accounting row, logging (never returning) a recorder error. func (p *Handler) writeRecord(req store.Request) { if p.rec == nil { diff --git a/internal/proxy/proxy.go b/internal/proxy/proxy.go index 23619e1..60feb29 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -7,6 +7,7 @@ package proxy import ( "bytes" + "context" "encoding/json" "errors" "io" @@ -45,6 +46,11 @@ type Recorder interface { RecordRequest(store.Request) error } +// Waker rouses a sleeping host. *wake.Waker satisfies it. +type Waker interface { + Wake(ctx context.Context, host string) bool +} + // Handler forwards requests for a route to one of the route's healthy hosts, choosing by lease when // one is configured and by health alone otherwise. type Handler struct { @@ -54,6 +60,13 @@ type Handler struct { lim *limiter.Limiter rec Recorder log *slog.Logger + waker Waker +} + +// SetWaker installs the waker the consults when a route has no healthy host left. A nil waker +// (the default) leaves the ErrNoHost answer as it was in v0: a plain 503. +func (p *Handler) SetWaker(w Waker) { + p.waker = w } // New builds a Handler. A nil logger becomes slog.Default(). With a nil lease table it behaves like @@ -223,6 +236,12 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { if err != nil { switch { case errors.Is(err, lease.ErrNoHost): + // No host healthy. Ask a waker to rouse a sleeping one; it answers + // (served or 503) when it has had a turn, else falls through to the + // plain 503. + if p.waker != nil && p.wakeOnErrNoHost(w, r, route, routeCfg, rest, model, fp, started, lease.Key{Route: route, FP: fp, Model: model}) { + return + } p.writeError(w, http.StatusServiceUnavailable, "no healthy host") case errors.Is(err, lease.ErrPinnedDown): p.writeError(w, http.StatusServiceUnavailable, "pinned host down") @@ -232,6 +251,14 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { return } + // Slot, context guard and forward, holding the slot for the leased host. + p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, host, reused) +} + +// serveLeased queues the request against the leased host's limiter, runs the +// context guard, and forwards. The slot is held for the originally leased host +// even if the guard relocates the lease: the guard already moved it. +func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, host string, reused bool) { // Slot. A full queue is a 503; a context done while waiting means the client left. release, waited, err := p.lim.Acquire(r.Context(), host, model) if err != nil { @@ -273,3 +300,26 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { } p.forward(w, r, route, host, leaseState(reused), rest, fp, model, now, waited, 0, header) } + +// wakeOnErrNoHost answers the request when no host was healthy. It asks, in +// route order, each candidate with a wake target to rouse itself; a host that +// wakes is leased once more and then served. When none wakes, it answers 503 +// with the hosts it tried. It returns true when the request has been answered. +func (p *Handler) wakeOnErrNoHost(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, key lease.Key) bool { + var tried []string + for _, name := range routeCfg.Hosts { + if p.cfg.Hosts[name].Wake == nil { + continue + } + tried = append(tried, name) + if !p.waker.Wake(r.Context(), name) { + continue + } + if newHost, _, err := p.leases.Acquire(key, routeCfg.Hosts, time.Now()); err == nil { + p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, newHost, true) + return true + } + } + p.writeNoHealthyHost(w, route, model, fp, started, tried) + return true +} diff --git a/internal/proxy/proxy_test.go b/internal/proxy/proxy_test.go index c1bb26a..39980aa 100644 --- a/internal/proxy/proxy_test.go +++ b/internal/proxy/proxy_test.go @@ -85,6 +85,19 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) { } } +// waitUntil polls cond every 5 ms for up to two seconds and fails the test if it never holds. +func waitUntil(t *testing.T, cond func() bool) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(5 * time.Millisecond) + } + t.Fatal("condition not reached within two seconds") +} + func TestQueueFullIs503(t *testing.T) { alpha := newUpstream(t, "alpha") alpha.delay = 400 * time.Millisecond @@ -99,14 +112,20 @@ hosts = ["alpha"] default_model = "shared" `, alpha) codes := make(chan int, 3) - for i := 1; i <= 3; i++ { - go func(i int) { + fire := func(i int) { + go func() { resp := r.post("/r/v1/chat/completions", conversation(i, 1)) drain(resp) codes <- resp.StatusCode - }(i) - time.Sleep(30 * time.Millisecond) // arrival order: 1 runs, 2 queues, 3 finds the queue full + }() } + // Arrival order is enforced by watching the limiter, not by sleeping: 1 runs, 2 queues, + // 3 finds the queue full. + fire(1) + waitUntil(t, func() bool { return r.lim.InFlight("alpha", "shared") == 1 }) + fire(2) + waitUntil(t, func() bool { return r.lim.Queued("alpha", "shared") == 1 }) + fire(3) got := map[int]int{} for i := 0; i < 3; i++ { got[<-codes]++ diff --git a/tools/smoke.sh b/tools/smoke.sh index 7c2bcb4..420d546 100755 --- a/tools/smoke.sh +++ b/tools/smoke.sh @@ -1,80 +1,60 @@ #!/bin/sh -# Smoke run (v1): two fake upstreams, one crossbar with a fresh SQLite file, real HTTP. -# Checks routing, leases (sticky + header), failover, recovery, streaming, queueing, pin, drain, -# usage and metrics. Prints "smoke: ok" or fails with the crossbar log. +# Smoke run (v2): everything v1 checked, plus the context guard, wake-on-LAN and identity gating. +# Prints "smoke: ok" or fails with the crossbar log. set -eu cd "$(dirname "$0")/.." tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM pids="" -sed "s#^db .*#db = \"$tmp/crossbar.db\"#" example.toml > "$tmp/crossbar.toml" -bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" -slow 600 >"$tmp/alpha.log" 2>&1 & pids="$pids $!" -bin/fakeupstream -listen 127.0.0.1:18082 -name beta -models ornith-1.5-35b-a3b -down-file "$tmp/beta.down" >"$tmp/beta.log" 2>&1 & pids="$pids $!" +sed -e "s#^db .*#db = \"$tmp/crossbar.db\"#" -e 's#^identity .*#identity = "header"#' -e 's#^\# peers = \["talos"\]#peers = ["talos"]#' example.toml > "$tmp/crossbar.toml" +# alpha: small context (4096 per slot = 8192/2); beta: large, sleeps until woken +bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" -slow 600 -n-ctx 8192 -slots 2 >"$tmp/alpha.log" 2>&1 & pids="$pids $!" +bin/fakeupstream -listen 127.0.0.1:18082 -name beta -models ornith-1.5-35b-a3b -down-file "$tmp/beta.down" -n-ctx 131072 -slots 2 -wol-listen 127.0.0.1:19082 -wol-mac aa:bb:cc:dd:ee:02 >"$tmp/beta.log" 2>&1 & pids="$pids $!" +touch "$tmp/beta.down" # beta starts "asleep" bin/crossbar -config "$tmp/crossbar.toml" >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" -sleep 1.5 +sleep 2.5 # two polls: alpha healthy, beta down fail() { echo "smoke: FAIL: $*" >&2; echo "--- crossbar.log"; cat "$tmp/crossbar.log"; exit 1; } base=http://127.0.0.1:17777 conv() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"system","content":"smoke"},{"role":"user","content":"conversation %s"}]}' "$1"; } +big() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"user","content":"%s"}]}' "$(head -c 40000 /dev/zero | tr '\0' 'x')"; } hdrs() { curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host} %header{X-Crossbar-Lease}' "$@"; } -# 1. a conversation gets a lease and keeps it; beta wins (2 slots × weight 2 vs 1 × 1) +# 1. v1 behaviour: with beta asleep, opencode-a goes to alpha h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") -[ "$h" = "200 beta new" ] || fail "first turn should be '200 beta new', got '$h'" -h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") -[ "$h" = "200 beta reused" ] || fail "second turn should reuse beta, got '$h'" +[ "$h" = "200 alpha new" ] || fail "with beta asleep conversation A should be '200 alpha new', got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{[^}]*"n_ctx":8192' || fail "hosts view does not show alpha n_ctx 8192: $(curl -s $base/_crossbar/hosts)" -# 2. header route -h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Route: hermes-x' -d "$(conv B)" "$base/v1/chat/completions") -case "$h" in "200 beta new") ;; *) fail "header route hermes-x should be '200 beta new', got '$h'";; esac -h=$(curl -s -o /dev/null -w '%{http_code}' "$base/nope/v1/models"); [ "$h" = "404" ] || fail "unknown route 404, got $h" +# 2. context guard: a ~12k-token prompt does not fit alpha's 4096-token slot; beta is asleep and +# wakeable, so crossbar must send the magic packet, wait for beta, and place the prompt there. +start=$(date +%s) +h=$(hdrs -m 40 -X POST -H 'Content-Type: application/json' -d "$(big)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta new" ] || fail "oversized prompt should wake beta and land there, got '$h' after $(( $(date +%s) - start ))s" +grep -q "magic packet received" "$tmp/beta.log" || fail "beta never saw a magic packet" +curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":true' || fail "beta not healthy after wake" -# 3. pin opencode-a to alpha: conversation A's next turn moves (an operator pin outranks the lease) -h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d '{"host":"alpha","pin":true}' "$base/_crossbar/routes/opencode-a") -[ "$h" = "200" ] || fail "pin returned $h" -h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") -[ "$h" = "200 alpha new" ] || fail "after pin, conversation A should be '200 alpha new', got '$h'" -curl -s "$base/_crossbar/routes" | grep -q '"pinned":"alpha"' || fail "routes view does not show the pin: $(curl -s $base/_crossbar/routes)" +# 3. with beta awake, a prompt that fits nowhere is a 400 (both slots too small? no — beta fits): +# check the guard's refusal with a prompt beyond beta's 65536-per-slot too +# (the body goes through a file: a 300 KB string cannot be a single argv element on Linux) +{ printf '{"model":"ornith-1.5-35b-a3b","messages":[{"role":"user","content":"'; head -c 300000 /dev/zero | tr '\0' 'x'; printf '"}]}'; } > "$tmp/toolarge-req.json" +h=$(curl -s -o "$tmp/toolarge.json" -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d @"$tmp/toolarge-req.json" "$base/opencode-a/v1/chat/completions") +[ "$h" = "400" ] && grep -q '"prompt too large"' "$tmp/toolarge.json" || fail "300 KB prompt should be 400 prompt too large, got $h $(cat "$tmp/toolarge.json")" -# 4. queue: alpha has parallel 1, queue_max 1, and answers in 600 ms → of three concurrent, one is 503 -for i in 1 2 3; do (curl -s -o /dev/null -w '%{http_code}\n' -X POST -H 'Content-Type: application/json' -d "$(conv Q$i)" "$base/opencode-a/v1/chat/completions" >> "$tmp/codes") & sleep 0.1; done; wait $! 2>/dev/null || true -sleep 2.5 -sort "$tmp/codes" | uniq -c | tr -s ' ' > "$tmp/counts" -grep -q '2 200' "$tmp/counts" && grep -q '1 503' "$tmp/counts" || fail "queue test wanted two 200 and one 503, got: $(cat "$tmp/counts")" +# 4. identity: hermes-x is locked to peer talos (header mode) +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x without a peer header should be 403, got $h" +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: titan' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x as titan should be 403, got $h" +h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +case "$h" in 200*) ;; *) fail "hermes-x as talos should be 200, got '$h'";; esac +h=$(curl -s -o /dev/null -w '%{http_code}' "$base/_crossbar/hosts"); [ "$h" = "200" ] || fail "admin must not be gated, got $h" -# 5. release the pin, drain alpha: new conversations go to beta, A stays on alpha -curl -s -o /dev/null -X POST -H 'Content-Type: application/json' -d '{"release":true}' "$base/_crossbar/routes/opencode-a" -h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d '{"drain":true}' "$base/_crossbar/hosts/alpha"); [ "$h" = "200" ] || fail "drain returned $h" -h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$base/opencode-a/v1/chat/completions") -[ "$h" = "200 beta new" ] || fail "with alpha draining a new conversation should go to beta, got '$h'" -curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{[^}]*"draining":true' || fail "hosts view does not show alpha draining" -curl -s -o /dev/null -X POST -H 'Content-Type: application/json' -d '{"drain":false}' "$base/_crossbar/hosts/alpha" - -# 6. failover + recovery -touch "$tmp/beta.down"; sleep 2.5 -h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv C)" "$base/opencode-a/v1/chat/completions") -[ "$h" = "200 alpha new" ] || fail "with beta down conversation C should move to alpha, got '$h'" -curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":false' || fail "hosts view does not show beta unhealthy" -rm "$tmp/beta.down"; sleep 3.5 -curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":true' || fail "beta did not recover after two good polls" - -# 7. streaming still arrives incrementally, and the final usage chunk is untouched +# 5. v1 regression: streaming still incremental, usage and metrics present start=$(date +%s%N) -curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \ - "$base/opencode-a/v1/chat/completions" | while IFS= read -r line; do - [ -n "$line" ] || continue - now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line" - done > "$tmp/stream.txt" +curl -sN -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \ + "$base/hermes-x/v1/chat/completions" | while IFS= read -r line; do [ -n "$line" ] || continue; now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line"; done > "$tmp/stream.txt" firstms=$(head -1 "$tmp/stream.txt" | cut -d' ' -f1); lastms=$(tail -1 "$tmp/stream.txt" | cut -d' ' -f1) -[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst: $(cat "$tmp/stream.txt")" -grep -q '"usage"' "$tmp/stream.txt" && grep -q 'DONE' "$tmp/stream.txt" || fail "stream lost the usage chunk or DONE" - -# 8. accounting and metrics +[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst" sleep 1 -u=$(curl -s "$base/_crossbar/usage?by=host") -echo "$u" | grep -q '"key":"alpha"' && echo "$u" | grep -q '"key":"beta"' || fail "usage by host: $u" -echo "$u" | grep -q '"cached_tokens":[1-9]' || fail "usage has no cached tokens (SSE/JSON usage not captured): $u" -curl -s -H 'Accept: text/plain' "$base/_crossbar/usage?by=route" | grep -qi 'cache' || fail "text usage table missing" -m=$(curl -s "$base/_crossbar/metrics") -echo "$m" | grep -q 'crossbar_requests_total{route="opencode-a",host="alpha",status="503"} 1' || fail "metrics missing the 503: $m" -echo "$m" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing host health" -grep -q 'route=opencode-a host=' "$tmp/crossbar.log" || fail "no request log line" +curl -s "$base/_crossbar/usage?by=host" | grep -q '"cached_tokens":[1-9]' || fail "usage has no cached tokens" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_requests_total{route="hermes-x",host="beta",status="403"}' && fail "403s are refused before a lease and must not be counted as requests" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing beta health" echo "smoke: ok (stream spread $((lastms - firstms)) ms)"