5 Commits
Author SHA1 Message Date
kyleandClaude Fable 5.1 6bdcf437f4 v1.1 review: checklist, probes, three findings
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-25 09:02:13 -07:00
kyle 9f5b50acf6 Review fixes: record cancelled requests as 499; empty usage is an array
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
2026-09-25 09:00:31 -07:00
kyleandClaude Fable 5.1 74b8af4ce0 AGENTS.md: a refused tool call is not a reason to end the turn; v2 plan (props, ctx guard, wake, identity, wiring) as acceptance tests; v1.1 run note
v2 given tests compiled against a panic-only skeleton (go vet clean); no reference
implementation.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-25 08:37:38 -07:00
kyleandClaude Fable 5.1 9ba04b16e0 v1.1 plan: cancelled requests recorded as 499, empty usage is []; tests proven failing on master
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-25 08:28:07 -07:00
kyle bebee332a5 Merge v1: leases, limiter, chooser, fingerprint, SQLite store, accounting, admin (8/8 tasks by Ornith) 2026-09-25 07:00:34 -07:00
27 changed files with 1708 additions and 17 deletions
+5
View File
@@ -54,6 +54,11 @@ These come from defects found in review; the evidence is in `docs/implementer-lo
- When a rule says "every" or "everywhere", finish by listing each place it applies and checking - When a rule says "every" or "everywhere", finish by listing each place it applies and checking
them one by one. The task shows one place; the rule covers all of them. them one by one. The task shows one place; the rule covers all of them.
- Never end a turn by describing what you are about to do. Do it, then report. - Never end a turn by describing what you are about to do. Do it, then report.
- Work only inside this repository. Scratch programs under `/tmp` or anywhere else are refused
by the sandbox, and **a refused tool call is not a reason to end the turn**: write the
experiment as a `_test.go` file inside the repository (delete it before committing), or reason
it out. Two sessions have ended with a plan and no tool call right after a refusal; that
leaves the owner with no commit and no `stopped` row, the worst outcome.
## The gate ## The gate
+18
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 | | Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model |
|---|---|---|---|---|---|---|---| |---|---|---|---|---|---|---|---|
| v1.1/01-review-fixes | 2026-09-25 | done | 1 | pass | none | Copied `cancel_test.go` and `usage_empty_test.go` byte-identical from `docs/plans/v1.1/_files/`; the earlier session's fixes in `internal/proxy/proxy.go`, `internal/proxy/forward.go` and `internal/admin/admin_ops.go` were already in the working tree. `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b |
| v1/08-smoke-readme | 2026-09-25 | done | 1 | pass | owner-directed fix to `Free` in `proxy.Chooser` | Changed `Free` from `c.lim.FreeSlots(host)` (sum over every model) to per-model free slots, `freeForModel(cfg.Hosts[host], model, c.lim.InFlight(host, model))`, floored at 0 and 0 when the host does not list the model (new helper in hosts.go); the one code change the task directs. `go test -race ./internal/proxy/` and `make gate` pass on the first run; `make smoke` → `smoke: ok (stream spread 1006 ms)`. README intro, `## Configure` (added db/lease_idle/retention, rewrote queue_max and hosts.<name>.hosts) and `## Inspect`→`## Operate` (all six endpoints, examples taken from the smoke run) updated. | llama.cpp/ornith-1.5-35b-a3b | | v1/08-smoke-readme | 2026-09-25 | done | 1 | pass | owner-directed fix to `Free` in `proxy.Chooser` | Changed `Free` from `c.lim.FreeSlots(host)` (sum over every model) to per-model free slots, `freeForModel(cfg.Hosts[host], model, c.lim.InFlight(host, model))`, floored at 0 and 0 when the host does not list the model (new helper in hosts.go); the one code change the task directs. `go test -race ./internal/proxy/` and `make gate` pass on the first run; `make smoke` → `smoke: ok (stream spread 1006 ms)`. README intro, `## Configure` (added db/lease_idle/retention, rewrote queue_max and hosts.<name>.hosts) and `## Inspect`→`## Operate` (all six endpoints, examples taken from the smoke run) updated. | llama.cpp/ornith-1.5-35b-a3b |
| v1/07-main | 2026-09-25 | done | 1 | pass | none | Wired store, limiter and lease table into `cmd/crossbar/main.go`: `store.Open` before the health table, `limiter.Configure` per (host, model) from `cfg.Hosts`, `lease.New` with `proxy.Chooser`, `Candidates` for every route, three background goroutines (idle expiry per minute, prune per hour logging the count, host-health recording per `poll_interval`), and `st.Close` via `defer`. The 3s SIGTERM run exits 0 with `listening`/`shutting down`; the missing-config run exits 1. | llama.cpp/ornith-1.5-35b-a3b | | v1/07-main | 2026-09-25 | done | 1 | pass | none | Wired store, limiter and lease table into `cmd/crossbar/main.go`: `store.Open` before the health table, `limiter.Configure` per (host, model) from `cfg.Hosts`, `lease.New` with `proxy.Chooser`, `Candidates` for every route, three background goroutines (idle expiry per minute, prune per hour logging the count, host-health recording per `poll_interval`), and `st.Close` via `defer`. The 3s SIGTERM run exits 0 with `listening`/`shutting down`; the missing-config run exits 1. | llama.cpp/ornith-1.5-35b-a3b |
| 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`. | llama.cpp/ornith-1.5-35b-a3b | | 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`. | llama.cpp/ornith-1.5-35b-a3b |
@@ -80,3 +81,20 @@ Follow-ups for `v1.1`: fix 1 (record the row on the cancel path with status 499
(`[]`), and an acceptance test for each; consider `lease_idle` expiry while a request is in flight (`[]`), and an acceptance test for each; consider `lease_idle` expiry while a request is in flight
and `Prune` under concurrent writes, which this review did not probe. and `Prune` under concurrent writes, which this review did not probe.
### v1.1 review — 2026-09-25 (reviewer: claude, as owner)
Checked: one task commit `9f5b50a` with the trailer; both given tests byte-identical; protected
files untouched; `make gate` → `gate: ok`; `make smoke` → `smoke: ok (stream spread 1006 ms)`.
Probed: `/_crossbar/usage` on an empty store answers `[]`; a stream cut by the client after
0.4 s appears in `/_crossbar/metrics` as `crossbar_requests_total{…,status="499"} 1`.
Process: four sessions for one task. Sessions 1–3 each ended their turn right after the sandbox
refused a write or read outside the repository (the I9 pattern) — after the fix was already
correct, in sessions 2 and 3. Session 3 also chased test flakes caused by its own inference
loading the host (the owner measured 12/12 passes idle). Findings: (a) model — five
refusal-endings tonight in total; `AGENTS.md` now names the rule, and the fourth session obeyed
it; (b) task — the task text did not state that `httputil.ReverseProxy` aborts the handler with
`http.ErrAbortHandler` on client disconnect, the fact the fix depends on (added mid-run); (c)
test design — timing-based tests (limiter, queue, spread, cancel) have margins tuned for an idle
host; widen or retry in a later plan.
+66
View File
@@ -0,0 +1,66 @@
# v1.1 task 01: review fixes — cancelled clients are recorded; empty usage is `[]`
**Branch:** `v1.1` (create it from `master`: `git switch master && git switch -c v1.1`; `git status --short` must be empty first, otherwise stop)
**Commit subject:** `Review fixes: record cancelled requests as 499; empty usage is an array`
## What the reviewer observed
1. A client that disconnects mid-stream leaves **no accounting row**: after `curl -m 0.4 -N …`
against a streaming completion, `/_crossbar/usage` stayed empty. Task 05's rule 5 said the
reverse proxy's `ErrorHandler` does nothing on `context.Canceled`; rule 6 said "record what
you have when `ServeHTTP` returns". The second rule was not applied on that path, and the
same gap exists for a client that gives up while waiting in the limiter queue (rule 4 said
"just return, log 499"). Cancelled requests held a slot and cost prefill; usage and error
rate must see them. The task text was ambiguous (owner's fault); the fix is still needed.
2. `GET /_crossbar/usage` with no rows answers `null`. The spec said a JSON array. Clients iterate
the result; `null` is not iterable.
## Files
- Copy (never edit afterwards): `internal/proxy/cancel_test.go`, `internal/admin/usage_empty_test.go`
- Modify: files under `internal/proxy/` as needed (`forward.go`, `proxy.go`), `internal/admin/admin_ops.go` (or wherever the usage handler lives), `docs/implementer-log.md`
## Rules
1. **Every request that reached step 3 of `ServeHTTP` (a lease was acquired) writes exactly one
`store.Request` row**, on every exit path: normal completion, upstream error (502), queue full
(503), client cancelled while queued (**499**, `Err: "client cancelled while queued"`), client
cancelled during the forward (**499**, `Err: "client cancelled"`, with whatever tokens the tee
had seen). Detect the forward case with `r.Context().Err() != nil` after `rp.ServeHTTP`
returns, or in the `ErrorHandler` when `errors.Is(err, context.Canceled)`; do not write to the
client in that case, do not mark the host down, but do record. Status 499 is not an HTTP
status the client sees; it is the row's status (and the log line's), as nginx does.
2. **`/_crossbar/usage` JSON** encodes an empty result as `[]`: initialise the slice
(`rows := []store.UsageRow{}` / `make(..., 0)`) before encoding, on every `by` value and with
or without `since`. The text form prints its header line even with no rows.
3. Nothing else changes. Existing tests must keep passing; the two new ones must pass.
## Steps
- [ ] **1. Branch and copy.**
```sh
git switch master && git switch -c v1.1
cp docs/plans/v1.1/_files/internal/proxy/cancel_test.go internal/proxy/
cp docs/plans/v1.1/_files/internal/admin/usage_empty_test.go internal/admin/
```
- [ ] **2. See them fail.** `go test -run 'Cancel|UsageEmpty' ./internal/proxy/ ./internal/admin/`.
Expected: all three tests fail (`no 499 row`, `body "null"`). If one passes already, stop and report.
- [ ] **3. Fix.** `gofmt -w internal/`.
- [ ] **4. See everything pass.** `go test -race -count=2 ./...`. The cancel tests are timing-based with generous margins.
- [ ] **5. Run the gate.** `make gate`. Expected last line: `gate: ok`.
- [ ] **6. Log and commit.** Row `v1.1/01-review-fixes`.
```sh
git add internal/proxy internal/admin docs/implementer-log.md
git commit
```
## Done when
- Step 2 failed before the fix and `go test -race -count=2 ./...` passes after; `make gate` prints `gate: ok`; both copied tests byte-identical to `_files/`.
## Stop and report if
- Step 2 passes before any change, or the cancel tests fail intermittently after the fix (report the failure text).
+34
View File
@@ -0,0 +1,34 @@
# v1.1 implementation plan: review follow-ups
> **For the implementing model:** do not work from this file. The owner gives you one task file at
> a time. This file is the index for the owner and the reviewer.
**Goal:** close findings 1 and 2 of the v1 review (`docs/implementer-log.md`): a request whose
client disconnects — mid-stream or while queued — must still write its accounting row (status
499), and `/_crossbar/usage` with no rows must answer `[]`, not `null`.
**How this plan was made:** acceptance tests first, from the findings; no reference
implementation. Both given tests were run against `master` at the merge of `v1`: all three fail
there (the two cancel tests find no 499 row; the empty-usage test gets `null`).
## Tasks
| # | File | Delivers | Tests that define it |
|---|---|---|---|
| 01 | `01-review-fixes.md` | 499 rows on both cancel paths; `[]` for empty usage | `internal/proxy/cancel_test.go`, `internal/admin/usage_empty_test.go` |
Branch `v1.1`. One task, one fresh OpenCode session, one commit.
## For the reviewer
1. `git log --oneline master..v1.1`: one commit with the trailer.
2. `cmp` both copied tests; `git diff master..v1.1 --stat -- PLAN.md AGENTS.md docs/plans` empty.
3. `make gate`, `make smoke`.
4. Probe: cut a stream with `curl -m 0.4 -N …` against the smoke rig and confirm one `status="499"` line in `/_crossbar/metrics`.
## Changes during the run
- 2026-09-25, task 01, first session: ended after ~8 min with no commit and no row, right after
the sandbox refused a `/tmp` scratch program (the I9 pattern, third time tonight). The rule
against ending a turn on a refusal lived only in v1's task 05; it is now in `AGENTS.md`, so every
task carries it. Resumed from the working tree.
@@ -0,0 +1,31 @@
package admin_test
import (
"encoding/json"
"strings"
"testing"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
// An empty usage table is an empty JSON array, not null: clients iterate it.
func TestUsageEmptyIsAnArray(t *testing.T) {
r := newRig(t)
for _, q := range []string{"/_crossbar/usage", "/_crossbar/usage?by=host", "/_crossbar/usage?by=model&since=1h"} {
rec := r.do(t, "GET", q, "")
if rec.Code != 200 {
t.Fatalf("%s: %d", q, rec.Code)
}
if strings.TrimSpace(rec.Body.String()) != "[]" {
t.Errorf("%s: body %q, want []", q, rec.Body.String())
}
var rows []store.UsageRow
if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil || rows == nil || len(rows) != 0 {
t.Errorf("%s: decoded %v %v, want an empty non-nil slice", q, rows, err)
}
}
rec := r.do(t, "GET", "/_crossbar/usage?by=route", "", "Accept", "text/plain")
if rec.Code != 200 || !strings.Contains(rec.Body.String(), "key") {
t.Errorf("text form with no rows must still print the header: %d %q", rec.Code, rec.Body.String())
}
}
@@ -0,0 +1,109 @@
package proxy_test
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
// A client that goes away mid-stream is still a request that happened: it held a slot, it cost
// prefill, and it belongs in the accounting. The row records status 499 and a non-empty err.
func TestClientCancelMidStreamIsRecorded(t *testing.T) {
slow := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/health":
fmt.Fprint(w, `{"status":"ok"}`)
case "/v1/models":
fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"}]}`)
default:
w.Header().Set("Content-Type", "text/event-stream")
w.WriteHeader(200)
fmt.Fprint(w, "data: {\"choices\":[{\"delta\":{\"content\":\"first\"}}]}\n\n")
w.(http.Flusher).Flush()
select {
case <-r.Context().Done():
case <-time.After(3 * time.Second):
}
}
}))
t.Cleanup(slow.Close)
beta := newUpstream(t, "beta")
r := newRig(t, twoHosts, &upstream{name: "alpha", srv: slow}, beta)
ctx, cancel := context.WithCancel(context.Background())
body := `{"model":"alpha-only","stream":true,"messages":[{"role":"user","content":"cancel me"}]}`
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(body))
req.Header.Set("Content-Type", "application/json")
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
buf := make([]byte, 64)
if _, err := resp.Body.Read(buf); err != nil {
t.Fatalf("first chunk: %v", err)
}
cancel()
resp.Body.Close()
deadline := time.Now().Add(3 * time.Second)
var counts []store.StatusCount
for time.Now().Before(deadline) {
counts, _ = r.store.StatusCounts(time.Time{})
if len(counts) > 0 {
break
}
time.Sleep(25 * time.Millisecond)
}
if len(counts) != 1 || counts[0].Status != 499 || counts[0].Route != "r" || counts[0].Count != 1 {
t.Fatalf("status counts after a cancelled stream = %+v, want one row: route r, status 499", counts)
}
rows, _ := r.store.Usage(time.Time{}, store.ByRoute)
if len(rows) != 1 || rows[0].Requests != 1 || rows[0].Errors != 1 {
t.Errorf("usage = %+v, want 1 request counted as an error", rows)
}
}
// The same when the client gives up while waiting in the queue: a 499 row, no slot leaked.
func TestClientCancelWhileQueuedIsRecorded(t *testing.T) {
alpha := newUpstream(t, "alpha")
alpha.delay = 800 * time.Millisecond
r := newRig(t, `
listen = "127.0.0.1:1"
queue_max = 2
[hosts.alpha]
base_url = %q
models = { "shared" = { parallel = 1 } }
[routes.r]
hosts = ["alpha"]
default_model = "shared"
`, alpha)
go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot
time.Sleep(100 * time.Millisecond)
ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1)))
req.Header.Set("Content-Type", "application/json")
if _, err := http.DefaultClient.Do(req); err == nil {
t.Fatal("the queued request should have been cancelled by its context")
}
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
counts, _ := r.store.StatusCounts(time.Time{})
for _, c := range counts {
if c.Status == 499 {
if r.lim.Queued("alpha", "shared") != 0 {
t.Errorf("queued = %d after the waiter cancelled", r.lim.Queued("alpha", "shared"))
}
return
}
}
time.Sleep(25 * time.Millisecond)
}
t.Fatal("no 499 row recorded for the request cancelled while queued")
}
+64
View File
@@ -0,0 +1,64 @@
# v2 task 01: learn each host's context size from `/props`
**Branch:** `v2` (create it from `master`: `git switch master && git switch -c v2`; `git status --short` must be empty first, otherwise stop)
**Commit subject:** `Health: learn n_ctx and total_slots from /props`
## Goal
The poller already asks each host `/health` and `/v1/models`. It now also reads `/props` and
remembers the context size and slot count, so the proxy (task 02) can tell whether a prompt fits.
A missing or malformed `/props` is **not** a health failure: context is then simply unknown (0).
## Context
`llama-server` answers `GET /props` with a JSON object containing
`default_generation_settings.n_ctx` (the total context the server was started with) and
`total_slots` (how many parallel slots share it). With unified KV, one slot can use up to
`n_ctx / total_slots` tokens. Some builds omit fields; old ones 404. Read at most `MaxModelsBody`
bytes as for the other endpoints.
## Files
- Copy: `internal/health/props_test.go`
- Modify: `internal/health/health.go`, `internal/admin/admin.go` (or wherever `HostView` is built), `docs/implementer-log.md`
## Interfaces
`internal/health`, additions:
```go
type Status struct {
// …existing fields…
NCtx int `json:"n_ctx"` // total context from /props; 0 = unknown
Slots int `json:"slots"` // total_slots from /props; 0 = unknown
}
// PerSlotCtx is the context one request may use: NCtx / Slots, or NCtx when Slots is 0.
func (s Status) PerSlotCtx() int
```
Rules the tests check:
1. A poll is `/health`, `/v1/models` (as before), then `GET <base>/props`. If that request fails,
returns non-200, is not JSON, or lacks the fields, set `NCtx = 0`, `Slots = 0` and **do not
count the poll as failed**. Otherwise `NCtx = default_generation_settings.n_ctx`,
`Slots = total_slots` (negative values → 0).
2. `PerSlotCtx()` is integer division; `NCtx` when `Slots == 0`; 0 when `NCtx == 0`.
3. `admin.HostView` gains `NCtx int \`json:"n_ctx"\`` and `Slots int \`json:"slots"\`` copied
from the status (the v2 smoke reads `"n_ctx":8192` from `/_crossbar/hosts`). `admin_test.go`
must keep passing unchanged.
## Steps
- [ ] **1.** `git switch master && git switch -c v2`; `cp docs/plans/v2/_files/internal/health/props_test.go internal/health/`.
- [ ] **2. See it fail** (compile: `NCtx` undefined). **3. Write the code.** `gofmt -w internal/`.
- [ ] **4.** `go test -race -count=1 ./internal/health/ ./internal/admin/` → both `ok`.
- [ ] **5.** `make gate` → `gate: ok`. **6.** Row `v2/01-props`; commit.
```sh
git add internal/health internal/admin docs/implementer-log.md
git commit
```
## Done when
- Both packages pass; gate ok; `props_test.go` byte-identical to `_files/`.
+67
View File
@@ -0,0 +1,67 @@
# v2 task 02: the context-size guard
**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop)
**Commit subject:** `Proxy: move or refuse prompts that do not fit the leased host's context`
## Goal
The single worst failure a client sees today is the upstream's "Context size has been exceeded"
after a long wait. With `PerSlotCtx` known (task 01), crossbar can estimate a prompt's size from
its body and act before forwarding: move the conversation to a host where it fits, or answer 400
with the estimate and the largest slot available. `PLAN.md` §4b.
## Files
- Copy: `internal/proxy/ctxguard_test.go`
- Modify: `internal/proxy/proxy.go` and/or `forward.go` (new file `ctxguard.go` if that keeps files under 400 lines), `internal/lease/lease.go` (one addition, below), `docs/implementer-log.md`
## Interfaces
```go
// internal/proxy
const CtxHeader = "X-Crossbar-Ctx" // set only when the guard moved a conversation: "moved:<from>>><to>" e.g. "moved:small>big"
// internal/lease — one addition, the only change allowed there:
// Move re-leases k onto host (deleting any existing lease for k), records a LeaseEvent with
// Reason "ctx", FromHost the previous host ("" if none), ToHost host. Returns an error only from
// the persister.
func (t *Table) Move(k Key, host string, now time.Time) error
```
Rules the tests check (`ServeHTTP`, between acquiring the lease and taking a slot):
1. `estimate := int(float64(len(body)) / 4 * 1.2)` tokens (body = the bytes already peeked; GET/HEAD → 0).
2. `limit := PerSlotCtx` of the leased host's `health.Status`. If `limit == 0` (unknown) or
`estimate <= limit`: no action, no header.
3. Otherwise find, among the route's candidate hosts (in order), the healthy, non-draining hosts
whose `PerSlotCtx() >= estimate` — prefer one that lists the model as loaded, else one that
can serve it (`cfg.Serves`). If one exists: `leases.Move(key, host, now)`, set
`CtxHeader` to `moved:<old>><new>`, and continue with the new host (**this request and the
following turns**: the lease moved).
4. If none exists: **400** `{"error":"prompt too large","estimate":E,"max":M}` where `M` is the
largest `PerSlotCtx()` among the route's healthy hosts (0 if all unknown — but then rule 2
already let the request through). Record an accounting row with status 400 and
`Err: "prompt too large"`; do not mark anything down; do not forward.
5. The estimate is never logged with the body; the log line gains `ctx_est=E` only.
Add `ReasonCtx = "ctx"` to `internal/store` constants (one-line change, allowed).
## Steps
- [ ] **1.** `git switch v2`; `cp docs/plans/v2/_files/internal/proxy/ctxguard_test.go internal/proxy/`.
- [ ] **2. See it fail** (compile: `CtxHeader`). **3. Write the code.** `gofmt -w internal/`.
- [ ] **4.** `go test -race -count=2 ./internal/proxy/ ./internal/lease/` → `ok`.
- [ ] **5.** `make gate` → `gate: ok`. **6.** Row `v2/02-ctxguard`; commit.
```sh
git add internal/proxy internal/lease internal/store docs/implementer-log.md
git commit
```
## Done when
- Tests pass with `-race -count=2`; gate ok; the copied test is byte-identical.
## Stop and report if
- `TestStickyLeaseSurvivesGrowthUntilItDoesNotFit` fails on the *second* request after the move (the lease did not actually move): quote the lease table.
+65
View File
@@ -0,0 +1,65 @@
# v2 task 03: wake-on-LAN
**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop)
**Commit subject:** `Add the wake package: magic packets and a waiter`
## Goal
A pure package. `wake.MagicPacket` builds the 102-byte wake-on-LAN frame, `wake.Send` puts it on
the wire as UDP, and `wake.Waker` wakes a named host at most once per wait window and waits for
the health table to report it healthy. Task 05 uses it when a route has no healthy host.
## Context
A magic packet is six `0xff` bytes followed by the target MAC sixteen times, sent as a UDP
datagram to the LAN broadcast address (port 9 by convention). The sleeping Mac (titan) has
wake-on-magic-packet enabled; it takes 20–40 s to be reachable. Waking twice inside that window
is harmless but pointless, so the waker remembers when it last sent.
## Files
- Copy: `internal/wake/wake_test.go`
- Create: `internal/wake/wake.go`
- Modify: `docs/implementer-log.md`
## Interfaces
```go
package wake
type Target struct {
MAC, Broadcast string // "aa:bb:cc:dd:ee:ff" (also "-" separated, any case); "host:port"
Wait time.Duration
}
type Health interface{ Healthy(name string) bool }
func MagicPacket(mac string) ([]byte, error) // net.ParseMAC; must be 6 bytes; 102-byte frame
func Send(mac, broadcast string) error // one UDP datagram via net.DialUDP("udp4", …); errors from parse/resolve/write
type Waker struct { /* private: targets, health, mutex, last-sent per host, poll interval (default 1s) */ }
func New(targets map[string]Target, h Health) *Waker
func (w *Waker) PollEvery(d time.Duration) // test hook; production keeps the 1 s default
// Wake returns true as soon as h.Healthy(host) is true, false if host is unknown, if Wait passes,
// or if ctx ends first. It sends the packet only if none was sent for host in the last Wait.
func (w *Waker) Wake(ctx context.Context, host string) bool
```
Rules the tests check: packet layout; separators; errors for bad MACs and unresolvable
addresses; one packet per window; return within about `Wait` when the host never comes up;
early return on a cancelled context; `false` for an unknown host without sending anything.
Never panic; safe for concurrent `Wake` calls on different hosts.
## Steps
- [ ] **1.** `git switch v2`; `mkdir -p internal/wake`; copy the test.
- [ ] **2. See it fail** (compile). **3. Write `wake.go`.** `gofmt -w internal/wake/`.
- [ ] **4.** `go test -race -count=3 ./internal/wake/` → `ok` (timing tests; three runs).
- [ ] **5.** `make gate`. **6.** Row `v2/03-wake`; commit.
```sh
git add internal/wake docs/implementer-log.md
git commit
```
## Done when
- `-race -count=3` passes; gate ok; the copied test is byte-identical.
+92
View File
@@ -0,0 +1,92 @@
# v2 task 04: tailnet identity and the route gate; config additions
**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop)
**Commit subject:** `Add identity: whois resolver, checker, header mode, middleware; config for wake, peers, identity`
## Goal
Some routes should be usable only from particular tailnet nodes ("`hermes-talos` only from
talos"). `identity` resolves a caller's address to a tailnet node name — in production through
`tailscale whois --json <ip>`, in tests through a fake, in the smoke run through a header — and a
middleware in front of the proxy refuses other callers with 403. Config gains the keys the rest
of v2 needs.
## Files
- Copy: `internal/identity/identity_test.go`, `internal/identity/middleware_test.go`, `internal/identity/testdata/whois.json`, `internal/config/config_v2_test.go`
- Create: `internal/identity/identity.go`, `internal/identity/middleware.go`
- Modify: `internal/config/config.go`, `docs/implementer-log.md`
## Interfaces
```go
package identity
var (
ErrNotAPeer = errors.New("identity: not a tailnet peer")
ErrForbidden = errors.New("identity: forbidden route")
)
type ID struct{ Node, Login string }
type Resolver interface { Identity(ctx context.Context, ip string) (ID, error) }
// ParseWhois reads `tailscale whois --json` output: Node = Node.ComputedName (else Node.Name
// without its trailing dot and domain), Login = UserProfile.LoginName. Empty node → error.
func ParseWhois(raw []byte) (ID, error)
// TailscaleResolver runs `tailscale whois --json <ip>` (exec, 3 s timeout) and parses it; a
// non-zero exit is ErrNotAPeer; a missing binary is an error that the Checker treats as "deny".
type TailscaleResolver struct{ Bin string } // Bin default "tailscale"
func (TailscaleResolver) Identity(ctx context.Context, ip string) (ID, error)
type Checker struct { /* private: resolver, cache map[ip]ID with a 5-minute TTL, mutex */ }
func NewChecker(r Resolver) *Checker
// NewHeaderChecker trusts the X-Crossbar-Peer request header as the node name. TEST/SMOKE ONLY.
func NewHeaderChecker() *Checker
// Allow: nil when peers is empty (open route); otherwise the caller's node (from remoteAddr's
// IP, or the header in header mode) must be in peers, else ErrForbidden. Any resolver error,
// unparsable address or loopback → ErrForbidden.
func (c *Checker) Allow(ctx context.Context, peers []string, remoteAddr string) error
// Middleware names the route like the proxy (X-Crossbar-Route header, else first path segment),
// asks peersFor(route), and answers 403 {"error":"forbidden route"} when Allow refuses. Paths
// under /_crossbar/ and routes peersFor does not know pass straight through.
func Middleware(c *Checker, peersFor func(route string) ([]string, bool), next http.Handler) http.Handler
```
Header mode: `Allow` needs the request to read the header, but its signature takes an address.
Make the header checker's resolver read from a `context.Context` value that `Middleware` sets
(`identity.WithHeaderPeer(ctx, r.Header.Get("X-Crossbar-Peer"))`); the tests only observe the
behaviour. Cache: per address, 5 minutes, for both hit and `ErrNotAPeer`.
`internal/config` gains:
```go
Identity string `toml:"identity"` // "off" (default) | "tailscale" | "header"; anything else → *Error field "identity"
// on Host:
Wake *Wake `toml:"wake"` // nil when absent
type Wake struct { MAC string `toml:"mac"`; Broadcast string `toml:"broadcast"`; Wait Duration `toml:"wait"` }
// on Route:
Peers []string `toml:"peers"`
```
Validation (after the existing host/route checks): `wake.mac` must parse (`net.ParseMAC`, 6
bytes) → field `hosts.<h>.wake.mac`; `wake.broadcast` non-empty `host:port` → `hosts.<h>.wake.broadcast`;
`wake.wait` default 45 s, less than 5 s → `hosts.<h>.wake.wait`. `routes.<r>.peers` non-empty
while `identity == "off"` → `routes.<r>.peers` ("peers need identity = tailscale or header");
`peers = []` (present but empty) with identity on → same field ("empty peers list").
## Steps
- [ ] **1.** `git switch v2`; `mkdir -p internal/identity/testdata`; copy the four given files.
- [ ] **2. See them fail** (compile). **3. Write the code.** `gofmt -w internal/`.
- [ ] **4.** `go test -race -count=1 ./internal/identity/ ./internal/config/` → `ok` (v0/v1 config tests included).
- [ ] **5.** `make gate`. **6.** Row `v2/04-identity`; commit.
```sh
git add internal/identity internal/config docs/implementer-log.md
git commit
```
## Done when
- Both packages pass; gate ok; all four copied files byte-identical.
+58
View File
@@ -0,0 +1,58 @@
# v2 task 05: wiring, the smoke run, README
**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop)
**Commit subject:** `Wire wake and identity into crossbar; v2 smoke and README`
## Goal
Put the pieces together: the proxy wakes a sleeping host when a route has no healthy host left,
`main` builds the waker from config and wraps the proxy in the identity middleware when
`identity` is on, the given fake upstream can be woken, and `tools/smoke.sh` proves the whole of
v2 over real HTTP.
## Files
- Copy (**replaces** v1's): `cmd/fakeupstream/main.go`, `tools/smoke.sh`, `example.toml`
- Modify: `internal/proxy/proxy.go` (or a new file), `cmd/crossbar/main.go`, `README.md`, `docs/implementer-log.md`
## Rules
1. `proxy.Handler` gains `func (p *Handler) SetWaker(w Waker)` where
`type Waker interface{ Wake(ctx context.Context, host string) bool }` (defined in `proxy`).
When `leases.Acquire` returns `ErrNoHost` and a waker is set: for each candidate host of the
route, in order, that has a wake target (ask `cfg.Hosts[h].Wake != nil`), call
`Wake(r.Context(), h)`; on `true`, retry `Acquire` once; on `false` for every candidate, 503
`{"error":"no healthy host","woke":["<hosts tried>"]}`. The context-guard's "no host fits"
path (task 02) also tries waking a host whose `PerSlotCtx` is unknown or large enough, before
answering 400.
2. `cmd/crossbar`: build `wake.New(targets, hosts)` from every host with `Wake != nil`
(`hosts` is the `proxy.HostView`, which has `Healthy`), call `p.SetWaker(w)`. When
`cfg.Identity != "off"`: `checker := identity.NewChecker(identity.TailscaleResolver{})` or
`identity.NewHeaderChecker()`; wrap the proxy handler:
`mux.Handle("/", identity.Middleware(checker, func(route string) ([]string, bool) { rt, ok := cfg.Routes[route]; return rt.Peers, ok }, p))`.
Log at start which mode is active; with `"header"` log a warning that it is insecure.
3. Copy the three given files; `make build`; `make smoke` → `smoke: ok (…)`. The smoke's check 2
waits up to 40 s for the wake; the fake wakes in ~1 s.
4. README: sections stay; add under `## Operate` the wake behaviour and the identity modes with
the `peers` example; under `## Configure` the three new keys; `## What v2 does not do`:
`/slots`, request coalescing, TLS — `PLAN.md`.
## Steps
- [ ] **1.** `git switch v2`; copy the three given files.
- [ ] **2. Write the code** (proxy waker path, `main`). `gofmt -w .`
- [ ] **3.** `go test -race -count=1 ./...` → all `ok`. **4.** `make smoke` → `smoke: ok`.
- [ ] **5.** README. **6.** `make gate`. **7.** Row `v2/05-wiring-smoke`; commit.
```sh
git add internal/proxy cmd/crossbar cmd/fakeupstream tools/smoke.sh example.toml README.md docs/implementer-log.md
git commit
```
## Done when
- `make smoke` prints `smoke: ok (…)`; gate ok; the three copied files byte-identical.
## Stop and report if
- `make smoke` fails twice in the same way; quote the failing check and the crossbar log.
+54
View File
@@ -0,0 +1,54 @@
# v2 implementation plan: learned context, the context guard, wake-on-LAN, identity
> **For the implementing model:** do not work from this file. The owner gives you one task file at
> a time (`01-…` to `05-…`). This file is the index for the owner and the reviewer.
**Goal:** `PLAN.md` §4b and §10 v2. The poller learns each host's context size from `/props`; a
prompt that cannot fit the leased host's per-slot context moves to one where it fits or is
refused with a clear 400; a route whose hosts are all down can wake a sleeping host by
wake-on-LAN and wait for it; a route can be restricted to named tailnet peers.
**Architecture:** `health.Status` gains `NCtx`/`Slots` (task 01); the proxy gains the guard
(task 02); two new small packages, `wake` (magic packets + a waiter, task 03) and `identity`
(whois resolver, checker, middleware, task 04); config gains `identity`, `[hosts.x.wake]`,
`routes.x.peers`; `main` wires the waker and the middleware (task 05).
**How this plan was made:** acceptance tests first, from `PLAN.md`; no reference implementation.
Every given test compiled against a panic-only skeleton of the names in the tasks (`go vet`
clean). The given tests were walked against the task rules and against the other given files
(helpers, line limits, `main.go` call sites) before handover — the v1 findings list is the
reason.
**Tech stack:** as v1; no new module. `identity`'s production resolver shells out to
`tailscale whois --json`, which exists on every fleet host.
## Global constraints
- Everything in `AGENTS.md`. Branch `v2`. One task, one fresh OpenCode session, one commit.
- Bodies never logged. Type assertions two-valued. Files under 400 lines.
- Given files are copied and never edited; some **replace** earlier ones (the task says so).
## Tasks
| # | File | Delivers | Tests that define it |
|---|---|---|---|
| 01 | `01-props.md` | `Status.NCtx`, `Status.Slots`, `PerSlotCtx()`; `/props` in the poll; hosts view shows them | `health/props_test.go` |
| 02 | `02-ctxguard.md` | prompt-size estimate; move or 400; `X-Crossbar-Ctx` | `proxy/ctxguard_test.go` |
| 03 | `03-wake.md` | `internal/wake`: magic packet, `Send`, `Waker` | `wake/wake_test.go` |
| 04 | `04-identity.md` | `internal/identity`: whois parse, checker, header mode, middleware; config `identity`/`peers`/`wake` | `identity/*_test.go`, `config/config_v2_test.go` |
| 05 | `05-wiring-smoke.md` | proxy wakes on no-host; `main` wires waker + middleware; given fakeupstream/smoke/example; README | `make smoke` |
## For the owner
`tools/run-plan.sh docs/plans/v2` from a clean checkout on `master`.
## For the reviewer: after task 05
1. Five task commits with the trailer; given files byte-identical; protected files untouched.
2. `make gate`, `make smoke`.
3. Probe: a `/props` that returns 200 with a huge body (bounded read); a MAC with an unusual
separator in config; `identity = "tailscale"` on a host where `tailscale` is not on PATH
(must log and refuse the gated routes, never allow); two routes, one gated one open, from
the same peer; a wake target whose broadcast address is unroutable (503 within `wait`, no
hang); the guard with a body of exactly `MaxBody`.
4. Findings under "Reviews" in `docs/implementer-log.md`, by fault.
@@ -0,0 +1,161 @@
// fakeupstream stands in for a llama-server router in tests and the smoke run. Do not edit.
//
// fakeupstream -listen 127.0.0.1:18081 -name alpha -models a,b -down-file /tmp/alpha.down -slow 0
//
// /health answers 503 while the down file exists, 200 otherwise. /v1/models lists -models.
// /props answers a small JSON object. /v1/chat/completions echoes: a streamed answer of five
// 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: <name>. /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 (
"encoding/json"
"flag"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"strings"
"time"
)
func main() {
listen := flag.String("listen", "127.0.0.1:18081", "address to listen on")
name := flag.String("name", "fake", "name reported in X-Upstream and answers")
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()
stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) }
usage := map[string]any{"prompt_tokens": 100, "completion_tokens": 10, "total_tokens": 110}
timings := map[string]any{"prompt_n": 100, "cache_n": 90, "predicted_n": 10, "predicted_ms": 50.0}
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
stamp(w)
if *downFile != "" {
if _, err := os.Stat(*downFile); err == nil {
http.Error(w, `{"error":{"message":"Loading model"}}`, http.StatusServiceUnavailable)
return
}
}
writeJSON(w, map[string]string{"status": "ok"})
})
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) {
stamp(w)
data := []map[string]any{}
for _, id := range ids {
data = append(data, map[string]any{"id": id, "object": "model", "owned_by": *name})
}
writeJSON(w, map[string]any{"object": "list", "data": data})
})
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": *nCtx}, "total_slots": *slots, "model_path": *name})
})
mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) {
stamp(w)
body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20))
var req struct {
Model string `json:"model"`
Stream bool `json:"stream"`
}
_ = json.Unmarshal(body, &req)
time.Sleep(time.Duration(*slow) * time.Millisecond)
if !req.Stream {
writeJSON(w, map[string]any{
"id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model,
"choices": []map[string]any{{"index": 0, "message": map[string]string{"role": "assistant", "content": "hello from " + *name}, "finish_reason": "stop"}},
"usage": usage, "timings": timings,
})
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.WriteHeader(http.StatusOK)
fl, _ := w.(http.Flusher)
flush := func() {
if fl != nil {
fl.Flush()
}
}
for i := 1; i <= 5; i++ {
chunk := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model,
"choices": []map[string]any{{"index": 0, "delta": map[string]string{"content": fmt.Sprintf("%s chunk %d ", *name, i)}}}}
b, _ := json.Marshal(chunk)
fmt.Fprintf(w, "data: %s\n\n", b)
flush()
time.Sleep(200 * time.Millisecond)
}
final := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model,
"choices": []map[string]any{}, "usage": usage, "timings": timings}
b, _ := json.Marshal(final)
fmt.Fprintf(w, "data: %s\n\n", b)
flush()
fmt.Fprint(w, "data: [DONE]\n\n")
})
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
stamp(w)
http.Error(w, `{"error":"not found"}`, http.StatusNotFound)
})
log.Printf("fakeupstream %s listening on %s models=%v slow=%dms", *name, *listen, ids, *slow)
srv := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second}
log.Fatal(srv.ListenAndServe())
}
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)
}
}
}
+32
View File
@@ -0,0 +1,32 @@
# 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
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.<tailnet>:11434
weight = 1.0
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 = 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.
[routes.opencode-a]
hosts = ["alpha", "beta"]
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
@@ -0,0 +1,75 @@
package config_test
import (
"strings"
"testing"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/config"
)
const v2Base = `
listen = "127.0.0.1:1"
[hosts.a]
base_url = "http://a:1"
models = { "m" = { } }
[hosts.b]
base_url = "http://b:1"
models = { "m" = { } }
[hosts.b.wake]
mac = "aa:bb:cc:dd:ee:ff"
broadcast = "192.168.1.255:9"
wait = "45s"
[routes.r]
hosts = ["a", "b"]
peers = ["talos", "imladris"]
`
func TestV2Defaults(t *testing.T) {
c, err := config.Parse(strings.NewReader(v2Base))
if err != nil {
t.Fatal(err)
}
if c.Identity != "off" {
t.Errorf("identity default = %q, want off", c.Identity)
}
if c.Hosts["a"].Wake != nil {
t.Errorf("host without [wake] must have nil Wake")
}
w := c.Hosts["b"].Wake
if w == nil || w.MAC != "aa:bb:cc:dd:ee:ff" || w.Broadcast != "192.168.1.255:9" || w.Wait.Duration != 45*time.Second {
t.Errorf("wake = %+v", w)
}
if p := c.Routes["r"].Peers; len(p) != 2 || p[0] != "talos" {
t.Errorf("peers = %v", p)
}
}
func TestV2Validation(t *testing.T) {
good := v2Base
for _, tc := range []struct{ name, text, field string }{
{"bad identity", "identity = \"maybe\"\n" + good, "identity"},
{"peers without identity", "identity = \"off\"\n" + good, "routes.r.peers"},
{"bad mac", strings.Replace(good, `mac = "aa:bb:cc:dd:ee:ff"`, `mac = "nope"`, 1), "hosts.b.wake.mac"},
{"no broadcast", strings.Replace(good, `broadcast = "192.168.1.255:9"`, `broadcast = ""`, 1), "hosts.b.wake.broadcast"},
{"wait too short", strings.Replace(good, `wait = "45s"`, `wait = "2s"`, 1), "hosts.b.wake.wait"},
{"peers on unknown route field", "identity = \"tailscale\"\n" + strings.Replace(good, `peers = ["talos", "imladris"]`, `peers = []`, 1), "routes.r.peers"},
} {
_, err := config.Parse(strings.NewReader(tc.text))
e, ok := config.IsError(err)
if !ok || e.Field != tc.field {
t.Errorf("%s: %v, want *Error on %s", tc.name, err, tc.field)
}
}
// identity = "header" is the test/smoke mode; "tailscale" the real one; both accept peers.
for _, mode := range []string{"header", "tailscale"} {
if _, err := config.Parse(strings.NewReader("identity = \"" + mode + "\"\n" + good)); err != nil {
t.Errorf("identity=%s with peers: %v", mode, err)
}
}
// wait defaults to 45s when the [wake] table omits it
c, err := config.Parse(strings.NewReader("identity = \"header\"\n" + strings.Replace(good, "wait = \"45s\"\n", "", 1)))
if err != nil || c.Hosts["b"].Wake == nil || c.Hosts["b"].Wake.Wait.Duration != 45*time.Second {
t.Errorf("wake.wait default: %v %+v", err, c.Hosts["b"].Wake)
}
}
@@ -0,0 +1,75 @@
package health_test
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/health"
)
// propsFake answers /health, /v1/models and a configurable /props.
func propsFake(t *testing.T, props string, status int) *httptest.Server {
mux := http.NewServeMux()
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) })
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"m"}]}`) })
mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(status)
fmt.Fprint(w, props)
})
srv := httptest.NewServer(mux)
t.Cleanup(srv.Close)
return srv
}
func TestPropsLearned(t *testing.T) {
srv := propsFake(t, `{"default_generation_settings":{"n_ctx":131072,"params":{}},"total_slots":4,"model_path":"/x/m.gguf","chat_template":"..."}`, 200)
tbl := health.New(map[string]string{"a": srv.URL}, time.Hour, nil)
tbl.PollOnce(context.Background())
s, _ := tbl.Get("a")
if !s.Healthy || s.NCtx != 131072 || s.Slots != 4 {
t.Fatalf("status = %+v, want healthy with NCtx 131072 and Slots 4", s)
}
if got := s.PerSlotCtx(); got != 32768 {
t.Errorf("PerSlotCtx = %d, want 131072/4", got)
}
}
func TestPropsAbsentOrBrokenIsNotAFailure(t *testing.T) {
for name, tc := range map[string]struct {
props string
status int
}{
"404": {`not found`, 404},
"not json": {`<html>`, 200},
"no fields": {`{"model_path":"/x"}`, 200},
"zero ctx": {`{"default_generation_settings":{"n_ctx":0},"total_slots":0}`, 200},
} {
t.Run(name, func(t *testing.T) {
srv := propsFake(t, tc.props, tc.status)
tbl := health.New(map[string]string{"a": srv.URL}, time.Hour, nil)
tbl.PollOnce(context.Background())
s, _ := tbl.Get("a")
if !s.Healthy {
t.Fatalf("a bad /props must not make the host unhealthy: %+v", s)
}
if s.NCtx != 0 || s.Slots != 0 || s.PerSlotCtx() != 0 {
t.Errorf("unknown context must read as 0: %+v", s)
}
})
}
}
func TestPerSlotCtxWithUnknownSlots(t *testing.T) {
s := health.Status{NCtx: 8192, Slots: 0}
if s.PerSlotCtx() != 8192 {
t.Errorf("with Slots unknown the whole context is the per-slot value; got %d", s.PerSlotCtx())
}
s = health.Status{NCtx: 8192, Slots: 3}
if s.PerSlotCtx() != 2730 {
t.Errorf("integer division: got %d, want 2730", s.PerSlotCtx())
}
}
@@ -0,0 +1,90 @@
package identity_test
import (
"context"
"errors"
"os"
"path/filepath"
"testing"
"git.wntrmute.dev/kyle/crossbar/internal/identity"
)
func TestParseWhois(t *testing.T) {
raw, err := os.ReadFile(filepath.Join("testdata", "whois.json"))
if err != nil {
t.Fatal(err)
}
id, err := identity.ParseWhois(raw)
if err != nil {
t.Fatal(err)
}
if id.Node != "titan" || id.Login == "" {
t.Errorf("parsed %+v, want Node titan and a login", id)
}
if _, err := identity.ParseWhois([]byte(`{"Node":{}}`)); err == nil {
t.Error("a whois answer without a node name must be an error")
}
if _, err := identity.ParseWhois([]byte(`nope`)); err == nil {
t.Error("non-JSON must be an error")
}
}
// fakeResolver answers from a map; "" means not a tailnet peer.
type fakeResolver map[string]string
func (f fakeResolver) Identity(ctx context.Context, ip string) (identity.ID, error) {
n, ok := f[ip]
if !ok {
return identity.ID{}, identity.ErrNotAPeer
}
return identity.ID{Node: n, Login: n + "@example"}, nil
}
func TestChecker(t *testing.T) {
c := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"})
for _, tc := range []struct {
name string
peers []string
addr string
want error
}{
{"open route", nil, "203.0.113.7:1", nil},
{"allowed peer", []string{"talos", "titan"}, "100.64.0.5:44444", nil},
{"other peer", []string{"talos"}, "100.64.0.9:1", identity.ErrForbidden},
{"not a peer", []string{"talos"}, "203.0.113.7:1", identity.ErrForbidden},
{"loopback", []string{"talos"}, "127.0.0.1:1", identity.ErrForbidden},
{"garbage addr", []string{"talos"}, "nonsense", identity.ErrForbidden},
} {
t.Run(tc.name, func(t *testing.T) {
got := c.Allow(context.Background(), tc.peers, tc.addr)
if !errors.Is(got, tc.want) && !(got == nil && tc.want == nil) {
t.Errorf("Allow(%v, %q) = %v, want %v", tc.peers, tc.addr, got, tc.want)
}
})
}
}
func TestCheckerCachesPerAddress(t *testing.T) {
calls := 0
r := countingResolver{f: fakeResolver{"100.64.0.5": "talos"}, calls: &calls}
c := identity.NewChecker(r)
for i := 0; i < 5; i++ {
if err := c.Allow(context.Background(), []string{"talos"}, "100.64.0.5:1"); err != nil {
t.Fatal(err)
}
}
if calls != 1 {
t.Errorf("resolver called %d times for one address, want 1 (cache)", calls)
}
}
type countingResolver struct {
f fakeResolver
calls *int
}
func (c countingResolver) Identity(ctx context.Context, ip string) (identity.ID, error) {
*c.calls++
return c.f.Identity(ctx, ip)
}
@@ -0,0 +1,74 @@
package identity_test
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
"git.wntrmute.dev/kyle/crossbar/internal/identity"
)
// The middleware sits in front of the proxy: it names the route the same way the proxy does
// (X-Crossbar-Route header, else first path segment) and refuses callers a route does not list.
func TestMiddleware(t *testing.T) {
peers := map[string][]string{"locked": {"talos"}, "open": nil}
inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) })
h := identity.Middleware(identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"}),
func(route string) ([]string, bool) { p, ok := peers[route]; return p, ok }, inner)
for _, tc := range []struct {
name, path, hdr, addr string
want int
}{
{"open route, anyone", "/open/v1/models", "", "203.0.113.1:5", 204},
{"locked, right peer", "/locked/v1/models", "", "100.64.0.5:5", 204},
{"locked, wrong peer", "/locked/v1/models", "", "100.64.0.9:5", 403},
{"locked, not a peer", "/locked/v1/models", "", "203.0.113.1:5", 403},
{"locked via header", "/v1/models", "locked", "100.64.0.9:5", 403},
{"header wins over path", "/open/v1/models", "locked", "203.0.113.1:5", 403},
{"unknown route passes through to the proxy's own 404", "/nope/v1/models", "", "203.0.113.1:5", 204},
{"admin path is never gated here", "/_crossbar/hosts", "", "203.0.113.1:5", 204},
} {
t.Run(tc.name, func(t *testing.T) {
req := httptest.NewRequest(http.MethodGet, tc.path, nil)
req.RemoteAddr = tc.addr
if tc.hdr != "" {
req.Header.Set("X-Crossbar-Route", tc.hdr)
}
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != tc.want {
t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String())
}
if rec.Code == 403 && (!strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") || !strings.Contains(rec.Body.String(), `"forbidden route"`)) {
t.Errorf("403 must be JSON {\"error\":\"forbidden route\"}: %q", rec.Body.String())
}
})
}
}
// HeaderResolver is the test/smoke identity source: it trusts X-Crossbar-Peer. It exists so the
// smoke run can exercise the gate without a tailnet; config must call it out as insecure.
func TestHeaderResolver(t *testing.T) {
inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) })
h := identity.Middleware(identity.NewHeaderChecker(), func(route string) ([]string, bool) { return []string{"talos"}, true }, inner)
req := httptest.NewRequest(http.MethodGet, "/r/v1/models", nil)
req.Header.Set("X-Crossbar-Peer", "talos")
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != 204 {
t.Errorf("header peer talos: %d", rec.Code)
}
req.Header.Set("X-Crossbar-Peer", "titan")
rec = httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != 403 {
t.Errorf("header peer titan: %d, want 403", rec.Code)
}
req.Header.Del("X-Crossbar-Peer")
rec = httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != 403 {
t.Errorf("no header: %d, want 403", rec.Code)
}
}
@@ -0,0 +1,24 @@
{
"Node": {
"ID": 1,
"StableID": "nEXAMPLE",
"Name": "titan.example.ts.net.",
"User": 2,
"Addresses": [
"100.64.0.9/32",
"fd7a:115c:a1e0::9/128"
],
"HomeDERP": 2,
"Created": "2026-01-01T00:00:00Z",
"Cap": 138,
"Online": true,
"ComputedName": "titan",
"ComputedNameWithHost": "titan"
},
"UserProfile": {
"ID": 2,
"LoginName": "user@example.com",
"DisplayName": "Example User"
},
"CapMap": null
}
@@ -0,0 +1,142 @@
package proxy_test
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"testing"
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
)
// ctxUpstream is a fake router that reports a context size in /props and echoes completions.
func ctxUpstream(t *testing.T, name string, nCtx, slots int) *upstream {
u := &upstream{name: name}
mux := http.NewServeMux()
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) })
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"shared"}]}`) })
mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) {
fmt.Fprintf(w, `{"default_generation_settings":{"n_ctx":%d},"total_slots":%d}`, nCtx, slots)
})
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
u.hits.Add(1)
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"choices":[{"message":{"role":"assistant","content":"ok"}}],"usage":{"prompt_tokens":1,"completion_tokens":1}}`)
})
u.srv = httptest.NewServer(mux)
t.Cleanup(u.srv.Close)
return u
}
const ctxHosts = `
listen = "127.0.0.1:1"
queue_max = 2
[hosts.small]
base_url = %q
weight = 10.0
models = { "shared" = { parallel = 2 } }
[hosts.big]
base_url = %q
weight = 1.0
models = { "shared" = { parallel = 1 } }
[routes.r]
hosts = ["small", "big"]
default_model = "shared"
`
// bodyOfTokens builds a chat body whose byte size implies roughly n tokens under the guard's
// estimate (bytes/4 × 1.2): n tokens ≈ 3.33 n bytes ≈ 2n/3 five-byte words.
func bodyOfTokens(n int) string {
text := strings.Repeat("word ", n*2/3)
return fmt.Sprintf(`{"model":"shared","stream":false,"messages":[{"role":"user","content":"%s"}]}`, text)
}
func TestOversizedPromptMovesToAHostWhereItFits(t *testing.T) {
small := ctxUpstream(t, "small", 8192, 2) // 4096 per slot
big := ctxUpstream(t, "big", 131072, 1) // 131072 per slot
r := newRig(t, ctxHosts, small, big)
// A small prompt starts on `small` (weight 10).
resp := r.post("/r/v1/chat/completions", bodyOfTokens(100))
drain(resp)
if resp.Header.Get(proxy.HostHeader) != "small" {
t.Fatalf("small prompt went to %q, want small", resp.Header.Get(proxy.HostHeader))
}
// A new conversation with ~10k tokens does not fit small's 4096-token slot: it must be
// placed on big, with the reason visible in a header.
resp = r.post("/r/v1/chat/completions", bodyOfTokens(10000))
drain(resp)
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" {
t.Fatalf("oversized prompt: %d from %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader))
}
if got := resp.Header.Get(proxy.CtxHeader); !strings.HasPrefix(got, "moved") {
t.Errorf("%s = %q, want moved:… ", proxy.CtxHeader, got)
}
}
func TestOversizedPromptWithNoFitIs400(t *testing.T) {
small := ctxUpstream(t, "small", 8192, 2)
tiny := ctxUpstream(t, "big", 4096, 2) // also too small
r := newRig(t, ctxHosts, small, tiny)
resp := r.post("/r/v1/chat/completions", bodyOfTokens(10000))
body := drain(resp)
if resp.StatusCode != http.StatusBadRequest {
t.Fatalf("status %d body %s, want 400", resp.StatusCode, body)
}
var e map[string]any
if err := json.Unmarshal([]byte(body), &e); err != nil || e["error"] != "prompt too large" {
t.Fatalf("body = %s, want error 'prompt too large'", body)
}
if est, _ := e["estimate"].(float64); est < 8000 || est > 13000 {
t.Errorf("estimate = %v, want roughly 10000 tokens", e["estimate"])
}
if max, _ := e["max"].(float64); max != 4096 {
t.Errorf("max = %v, want the largest per-slot context among the route's hosts (4096)", e["max"])
}
if small.hits.Load()+tiny.hits.Load() != 0 {
t.Errorf("a refused prompt must not reach any upstream")
}
}
func TestUnknownContextNeverBlocks(t *testing.T) {
// /props missing on both hosts: NCtx 0 means "unknown", and the guard must stay out of the way.
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
r := newRig(t, twoHosts, alpha, beta)
resp := r.post("/r/v1/chat/completions", bodyOfTokens(50000))
drain(resp)
if resp.StatusCode != 200 || resp.Header.Get(proxy.CtxHeader) != "" {
t.Errorf("unknown context: %d %q, want 200 and no ctx header", resp.StatusCode, resp.Header.Get(proxy.CtxHeader))
}
}
func TestStickyLeaseSurvivesGrowthUntilItDoesNotFit(t *testing.T) {
small := ctxUpstream(t, "small", 8192, 2)
big := ctxUpstream(t, "big", 131072, 1)
r := newRig(t, ctxHosts, small, big)
body := bodyOfTokens(100)
resp := r.post("/r/v1/chat/completions", body)
drain(resp)
if resp.Header.Get(proxy.HostHeader) != "small" {
t.Fatal("setup: first turn must be on small")
}
// Same conversation (same first user message), later turn well under 4096: stays.
longer := strings.Replace(body, `"content":"`, `"content":"`+strings.Repeat("x ", 500), 1)
resp = r.post("/r/v1/chat/completions", longer)
drain(resp)
if resp.Header.Get(proxy.HostHeader) != "small" || resp.Header.Get(proxy.LeaseHeader) != "reused" {
t.Errorf("turn 2: %q %q, want small reused", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader))
}
// A turn that outgrows the slot moves the lease — once — and the move is recorded as an event.
huge := strings.Replace(body, `"content":"`, `"content":"`+strings.Repeat("x ", 30000), 1)
resp = r.post("/r/v1/chat/completions", huge)
drain(resp)
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" {
t.Fatalf("outgrown turn: %d %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader))
}
resp = r.post("/r/v1/chat/completions", huge)
drain(resp)
if resp.Header.Get(proxy.HostHeader) != "big" || resp.Header.Get(proxy.LeaseHeader) != "reused" {
t.Errorf("after the move the lease is on big: %q %q", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader))
}
}
@@ -0,0 +1,114 @@
package wake_test
import (
"bytes"
"context"
"net"
"testing"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/wake"
)
func listen(t *testing.T) (*net.UDPConn, string) {
conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { conn.Close() })
return conn, conn.LocalAddr().String()
}
func TestMagicPacket(t *testing.T) {
pkt, err := wake.MagicPacket("aa:bb:cc:dd:ee:ff")
if err != nil {
t.Fatal(err)
}
if len(pkt) != 102 || !bytes.Equal(pkt[:6], bytes.Repeat([]byte{0xff}, 6)) {
t.Fatalf("packet = % x", pkt)
}
mac := []byte{0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff}
for i := 0; i < 16; i++ {
if !bytes.Equal(pkt[6+6*i:12+6*i], mac) {
t.Fatalf("repetition %d wrong: % x", i, pkt[6+6*i:12+6*i])
}
}
for _, bad := range []string{"", "aa:bb", "zz:bb:cc:dd:ee:ff", "aabbccddeeff00"} {
if _, err := wake.MagicPacket(bad); err == nil {
t.Errorf("MagicPacket(%q) must fail", bad)
}
}
if p2, _ := wake.MagicPacket("AA-BB-CC-DD-EE-FF"); !bytes.Equal(p2, pkt) {
t.Errorf("dash-separated upper-case MAC must give the same packet")
}
}
func TestSendReachesTheBroadcastAddress(t *testing.T) {
conn, addr := listen(t)
if err := wake.Send("aa:bb:cc:dd:ee:ff", addr); err != nil {
t.Fatal(err)
}
buf := make([]byte, 200)
_ = conn.SetReadDeadline(time.Now().Add(time.Second))
n, _, err := conn.ReadFromUDP(buf)
if err != nil || n != 102 {
t.Fatalf("received %d bytes, err %v", n, err)
}
if err := wake.Send("aa:bb:cc:dd:ee:ff", "256.1.1.1:9"); err == nil {
t.Error("an unresolvable broadcast address must be an error")
}
}
// fakeHealth flips to healthy after `after` calls to Healthy.
type fakeHealth struct{ calls, after int }
func (f *fakeHealth) Healthy(name string) bool { f.calls++; return f.calls > f.after }
func TestWakerSendsOncePerWindowAndWaitsForHealth(t *testing.T) {
conn, addr := listen(t)
h := &fakeHealth{after: 3}
w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 2 * time.Second}}, h)
w.PollEvery(20 * time.Millisecond) // test hook: how often Wake re-checks health
start := time.Now()
ok := w.Wake(context.Background(), "titan")
if !ok {
t.Fatal("Wake must return true once the host reports healthy")
}
if time.Since(start) > time.Second {
t.Errorf("Wake waited %v for a host that came up after 3 checks", time.Since(start))
}
_ = conn.SetReadDeadline(time.Now().Add(200 * time.Millisecond))
buf := make([]byte, 200)
if n, _, err := conn.ReadFromUDP(buf); err != nil || n != 102 {
t.Fatalf("no magic packet received: %d %v", n, err)
}
// A second Wake inside the same window does not send again (the host is booting).
_ = w.Wake(context.Background(), "titan")
_ = conn.SetReadDeadline(time.Now().Add(150 * time.Millisecond))
if n, _, err := conn.ReadFromUDP(buf); err == nil {
t.Errorf("a second packet (%d bytes) was sent inside the wait window", n)
}
if w.Wake(context.Background(), "nobody") {
t.Errorf("unknown host: Wake must return false")
}
}
func TestWakeGivesUpAfterWait(t *testing.T) {
_, addr := listen(t)
h := &fakeHealth{after: 1 << 30}
w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 300 * time.Millisecond}}, h)
w.PollEvery(20 * time.Millisecond)
start := time.Now()
if w.Wake(context.Background(), "titan") {
t.Fatal("Wake must return false when the host never comes up")
}
if d := time.Since(start); d < 250*time.Millisecond || d > 900*time.Millisecond {
t.Errorf("Wake returned after %v, want about the 300ms wait", d)
}
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
start = time.Now()
if w.Wake(ctx, "titan") || time.Since(start) > 200*time.Millisecond {
t.Errorf("a cancelled context must end the wait early (took %v)", time.Since(start))
}
}
+58
View File
@@ -0,0 +1,58 @@
#!/bin/sh
# 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 -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 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. 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 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. 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. 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
h=$(curl -s -o "$tmp/toolarge.json" -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d "$(printf '{"model":"ornith-1.5-35b-a3b","messages":[{"role":"user","content":"%s"}]}' "$(head -c 300000 /dev/zero | tr '\0' 'x')")" "$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. 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. v1 regression: streaming still incremental, usage and metrics present
start=$(date +%s%N)
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"
sleep 1
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)"
+3
View File
@@ -141,6 +141,9 @@ func (hx *handler) usageGet(w http.ResponseWriter, r *http.Request) {
writeError(w, http.StatusInternalServerError, "usage: "+err.Error()) writeError(w, http.StatusInternalServerError, "usage: "+err.Error())
return return
} }
if rows == nil {
rows = []store.UsageRow{}
}
if r.Header.Get("Accept") == "text/plain" { if r.Header.Get("Accept") == "text/plain" {
writeUsageTable(w, rows) writeUsageTable(w, rows)
return return
+31
View File
@@ -0,0 +1,31 @@
package admin_test
import (
"encoding/json"
"strings"
"testing"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
// An empty usage table is an empty JSON array, not null: clients iterate it.
func TestUsageEmptyIsAnArray(t *testing.T) {
r := newRig(t)
for _, q := range []string{"/_crossbar/usage", "/_crossbar/usage?by=host", "/_crossbar/usage?by=model&since=1h"} {
rec := r.do(t, "GET", q, "")
if rec.Code != 200 {
t.Fatalf("%s: %d", q, rec.Code)
}
if strings.TrimSpace(rec.Body.String()) != "[]" {
t.Errorf("%s: body %q, want []", q, rec.Body.String())
}
var rows []store.UsageRow
if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil || rows == nil || len(rows) != 0 {
t.Errorf("%s: decoded %v %v, want an empty non-nil slice", q, rows, err)
}
}
rec := r.do(t, "GET", "/_crossbar/usage?by=route", "", "Accept", "text/plain")
if rec.Code != 200 || !strings.Contains(rec.Body.String(), "key") {
t.Errorf("text form with no rows must still print the header: %d %q", rec.Code, rec.Body.String())
}
}
+109
View File
@@ -0,0 +1,109 @@
package proxy_test
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
// A client that goes away mid-stream is still a request that happened: it held a slot, it cost
// prefill, and it belongs in the accounting. The row records status 499 and a non-empty err.
func TestClientCancelMidStreamIsRecorded(t *testing.T) {
slow := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/health":
fmt.Fprint(w, `{"status":"ok"}`)
case "/v1/models":
fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"}]}`)
default:
w.Header().Set("Content-Type", "text/event-stream")
w.WriteHeader(200)
fmt.Fprint(w, "data: {\"choices\":[{\"delta\":{\"content\":\"first\"}}]}\n\n")
w.(http.Flusher).Flush()
select {
case <-r.Context().Done():
case <-time.After(3 * time.Second):
}
}
}))
t.Cleanup(slow.Close)
beta := newUpstream(t, "beta")
r := newRig(t, twoHosts, &upstream{name: "alpha", srv: slow}, beta)
ctx, cancel := context.WithCancel(context.Background())
body := `{"model":"alpha-only","stream":true,"messages":[{"role":"user","content":"cancel me"}]}`
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(body))
req.Header.Set("Content-Type", "application/json")
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
buf := make([]byte, 64)
if _, err := resp.Body.Read(buf); err != nil {
t.Fatalf("first chunk: %v", err)
}
cancel()
resp.Body.Close()
deadline := time.Now().Add(3 * time.Second)
var counts []store.StatusCount
for time.Now().Before(deadline) {
counts, _ = r.store.StatusCounts(time.Time{})
if len(counts) > 0 {
break
}
time.Sleep(25 * time.Millisecond)
}
if len(counts) != 1 || counts[0].Status != 499 || counts[0].Route != "r" || counts[0].Count != 1 {
t.Fatalf("status counts after a cancelled stream = %+v, want one row: route r, status 499", counts)
}
rows, _ := r.store.Usage(time.Time{}, store.ByRoute)
if len(rows) != 1 || rows[0].Requests != 1 || rows[0].Errors != 1 {
t.Errorf("usage = %+v, want 1 request counted as an error", rows)
}
}
// The same when the client gives up while waiting in the queue: a 499 row, no slot leaked.
func TestClientCancelWhileQueuedIsRecorded(t *testing.T) {
alpha := newUpstream(t, "alpha")
alpha.delay = 800 * time.Millisecond
r := newRig(t, `
listen = "127.0.0.1:1"
queue_max = 2
[hosts.alpha]
base_url = %q
models = { "shared" = { parallel = 1 } }
[routes.r]
hosts = ["alpha"]
default_model = "shared"
`, alpha)
go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot
time.Sleep(100 * time.Millisecond)
ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1)))
req.Header.Set("Content-Type", "application/json")
if _, err := http.DefaultClient.Do(req); err == nil {
t.Fatal("the queued request should have been cancelled by its context")
}
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
counts, _ := r.store.StatusCounts(time.Time{})
for _, c := range counts {
if c.Status == 499 {
if r.lim.Queued("alpha", "shared") != 0 {
t.Errorf("queued = %d after the waiter cancelled", r.lim.Queued("alpha", "shared"))
}
return
}
}
time.Sleep(25 * time.Millisecond)
}
t.Fatal("no 499 row recorded for the request cancelled while queued")
}
+47 -17
View File
@@ -30,26 +30,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
rev := &forwardState{started: started} rev := &forwardState{started: started}
rp := newReverseProxy(p.health, host, leaseState, target, rest, rev) rp := newReverseProxy(p.health, host, leaseState, target, rest, rev)
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK} rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
// ServeHTTP unwinds with http.ErrAbortHandler when a client leaves mid-stream; recover so the
// row the request earned is still written, then re-panic so the server keeps its semantics.
defer func() {
if pv := recover(); pv != nil {
total := time.Since(started).Milliseconds()
req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total)
if perr, ok := pv.(error); ok && errors.Is(perr, http.ErrAbortHandler) {
req.Status = 499
req.Err = "client cancelled"
} else {
req.Err = "upstream error"
}
p.writeRecord(req)
panic(pv)
}
}()
rp.ServeHTTP(rec, r) rp.ServeHTTP(rec, r)
total := time.Since(started) total := time.Since(started)
req := store.Request{ req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total.Milliseconds())
Route: route, if r.Context().Err() != nil {
FP: fp, req.Status = 499
Model: model, req.Err = "client cancelled"
Host: host,
Started: started,
QueuedMs: waited.Milliseconds(),
TTFBMs: ttfbMs(rev),
TotalMs: total.Milliseconds(),
Status: rec.status,
Streamed: rev.streamed,
}
if rev.tee != nil {
prompt, cached, completion := rev.tee.tokens()
req.PromptTokens = int64(prompt)
req.CachedTokens = int64(cached)
req.CompletionTokens = int64(completion)
} }
p.writeRecord(req) p.writeRecord(req)
@@ -70,6 +75,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
) )
} }
// forwardRow builds the accounting row from the state a forward gathered: what the tee scanned and
// what the recorder captured. totalMs is measured from start to the caller's exit, so the forward
// path and the recovery path above build identical rows.
func forwardRow(route, fp, model, host string, start time.Time, waited time.Duration, rev *forwardState, status int, totalMs int64) store.Request {
req := store.Request{
Route: route,
FP: fp,
Model: model,
Host: host,
Started: start,
QueuedMs: waited.Milliseconds(),
TTFBMs: ttfbMs(rev),
TotalMs: totalMs,
Status: status,
Streamed: rev.streamed,
}
if rev.tee != nil {
prompt, cached, completion := rev.tee.tokens()
req.PromptTokens = int64(prompt)
req.CachedTokens = int64(cached)
req.CompletionTokens = int64(completion)
}
return req
}
// leaseState is "reused" when the lease already held the conversation, else "new". // leaseState is "reused" when the lease already held the conversation, else "new".
func leaseState(reused bool) string { func leaseState(reused bool) string {
if reused { if reused {
+10
View File
@@ -249,6 +249,16 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
return return
} }
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499) p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
p.writeRecord(store.Request{
Route: route,
FP: fp,
Model: model,
Host: host,
Started: started,
TotalMs: time.Since(started).Milliseconds(),
Status: 499,
Err: "client cancelled while queued",
})
return return
} }
defer release() defer release()