From c457046f8baa0575cb6f2ebb1371b8668efa7a6f Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 03:09:50 -0700 Subject: [PATCH] =?UTF-8?q?v1=20plan:=20leases,=20limiter,=20chooser,=20fi?= =?UTF-8?q?ngerprint,=20SQLite=20store,=20accounting,=20admin=20=E2=80=94?= =?UTF-8?q?=20acceptance=20tests=20first,=20no=20reference?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every given test compiled against a panic-only interface skeleton (go vet clean); nothing was implemented. modernc.org/sqlite v1.59.0 vetted in a scratch module (WAL works); go.sum given. Co-Authored-By: Claude Fable 5.1 --- docs/plans/v1/01-store.md | 154 +++++++ docs/plans/v1/02-fingerprint-config.md | 86 ++++ docs/plans/v1/03-limiter-choose.md | 94 +++++ docs/plans/v1/04-lease.md | 117 +++++ docs/plans/v1/05-proxy.md | 122 ++++++ docs/plans/v1/06-admin-main.md | 121 ++++++ docs/plans/v1/07-smoke-readme.md | 53 +++ docs/plans/v1/README.md | 61 +++ docs/plans/v1/_files/cmd/fakeupstream/main.go | 115 +++++ docs/plans/v1/_files/example.toml | 26 ++ docs/plans/v1/_files/go.sum | 52 +++ .../v1/_files/internal/admin/admin_test.go | 293 +++++++++++++ .../v1/_files/internal/choose/choose_test.go | 83 ++++ .../_files/internal/config/config_v1_test.go | 81 ++++ .../internal/fingerprint/fingerprint_test.go | 70 +++ .../v1/_files/internal/lease/lease_test.go | 291 +++++++++++++ .../_files/internal/limiter/limiter_test.go | 178 ++++++++ .../v1/_files/internal/proxy/proxy_test.go | 398 ++++++++++++++++++ .../v1/_files/internal/store/store_test.go | 185 ++++++++ docs/plans/v1/_files/tools/smoke.sh | 80 ++++ 20 files changed, 2660 insertions(+) create mode 100644 docs/plans/v1/01-store.md create mode 100644 docs/plans/v1/02-fingerprint-config.md create mode 100644 docs/plans/v1/03-limiter-choose.md create mode 100644 docs/plans/v1/04-lease.md create mode 100644 docs/plans/v1/05-proxy.md create mode 100644 docs/plans/v1/06-admin-main.md create mode 100644 docs/plans/v1/07-smoke-readme.md create mode 100644 docs/plans/v1/README.md create mode 100644 docs/plans/v1/_files/cmd/fakeupstream/main.go create mode 100644 docs/plans/v1/_files/example.toml create mode 100644 docs/plans/v1/_files/go.sum create mode 100644 docs/plans/v1/_files/internal/admin/admin_test.go create mode 100644 docs/plans/v1/_files/internal/choose/choose_test.go create mode 100644 docs/plans/v1/_files/internal/config/config_v1_test.go create mode 100644 docs/plans/v1/_files/internal/fingerprint/fingerprint_test.go create mode 100644 docs/plans/v1/_files/internal/lease/lease_test.go create mode 100644 docs/plans/v1/_files/internal/limiter/limiter_test.go create mode 100644 docs/plans/v1/_files/internal/proxy/proxy_test.go create mode 100644 docs/plans/v1/_files/internal/store/store_test.go create mode 100755 docs/plans/v1/_files/tools/smoke.sh diff --git a/docs/plans/v1/01-store.md b/docs/plans/v1/01-store.md new file mode 100644 index 0000000..a72103b --- /dev/null +++ b/docs/plans/v1/01-store.md @@ -0,0 +1,154 @@ +# v1 task 01: the SQLite store + +**Branch:** `v1` (create it from `master`: `git switch master && git switch -c v1`; `git status --short` must be empty first, otherwise stop) +**Commit subject:** `Add the SQLite store for leases and accounting` + +## Goal + +One SQLite file holds crossbar's durable state (the lease table) and its accounting log (one row +per proxied request, lease events, poller observations), with rollup queries that answer +per-route / per-model / per-host usage. This is `PLAN.md` §7a. + +## Context + +The driver is `modernc.org/sqlite` (pure Go, no cgo — the arm64 static build stays a plain +`go build`), registered under the `database/sql` name `"sqlite"`. Open with WAL and a busy +timeout: `sql.Open("sqlite", "file:"+path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")`. +Volume is a few rows per request, so nothing here is performance-sensitive; correctness of the +sums is what matters. Times are stored as Unix milliseconds (`INTEGER`) and returned as +`time.Time` in UTC. + +## Files + +- Copy (never edit afterwards): `go.sum` (replaces; adds the driver's lines), `internal/store/store_test.go` +- Create: `internal/store/store.go` (and `schema.go` if you want the SQL separate; both under 400 lines) +- Modify: `go.mod` (add `modernc.org/sqlite v1.59.0` to `require`), `docs/implementer-log.md` + +## Interfaces + +Produces, in `internal/store`, package `store`: + +```go +type State string +const ( Active State = "active"; Pinned State = "pinned" ) + +const ( + ReasonNew = "new"; ReasonUnhealthy = "unhealthy"; ReasonIdle = "idle" + ReasonPin = "pin"; ReasonRelease = "release"; ReasonDrain = "drain" +) + +type By string +const ( ByRoute By = "route"; ByModel By = "model"; ByHost By = "host" ) + +type Lease struct { + Route, FP, Model, Host string + State State + Created, LastUsed time.Time +} +type LeaseEvent struct { + TS time.Time + Route, Model, FromHost, ToHost string + Reason string +} +type Request struct { + Route, FP, Model, Host string + Started time.Time + QueuedMs, TTFBMs, TotalMs int64 + Status int + Streamed bool + PromptTokens, CachedTokens, CompletionTokens int64 + Err string +} +type HostHealth struct { + TS time.Time + Host string + Healthy bool + Loaded []string // stored as a JSON array in one TEXT column +} +type UsageRow struct { + Key string `json:"key"` + Requests int64 `json:"requests"` + Errors int64 `json:"errors"` // rows with Status >= 400 + BusyMs int64 `json:"busy_ms"` // sum(TotalMs) + QueuedMs int64 `json:"queued_ms"` + PromptTokens int64 `json:"prompt_tokens"` + CachedTokens int64 `json:"cached_tokens"` + CompletionTokens int64 `json:"completion_tokens"` +} +func (u UsageRow) CacheHitRatio() float64 // CachedTokens / PromptTokens; 0 when PromptTokens == 0 + +type Store struct { /* private: *sql.DB */ } +func Open(path string) (*Store, error) // creates tables if missing; fails if the directory does not exist +func (s *Store) Close() error +func (s *Store) JournalMode() string // "wal" +func (s *Store) SaveLease(l Lease) error // INSERT OR REPLACE on (route, fp, model) +func (s *Store) DeleteLease(route, fp, model string) error +func (s *Store) ListLeases() ([]Lease, error) +func (s *Store) RecordEvent(e LeaseEvent) error +func (s *Store) RecordRequest(r Request) error +func (s *Store) RecordHostHealth(h HostHealth) error +func (s *Store) Usage(since time.Time, by By) ([]UsageRow, error) // rows with Started >= since, grouped by `by`; a zero `since` means all time +func (s *Store) Events(since time.Time, limit int) ([]LeaseEvent, error) // oldest first +func (s *Store) Prune(now time.Time, retention time.Duration) (int64, error) +``` + +Rules the tests check: + +1. **Schema** (create with `IF NOT EXISTS`, so `Open` twice on one file works): + `leases(route, fp, model, host, state, created, last_used, PRIMARY KEY(route, fp, model))`, + `lease_events(ts, route, model, from_host, to_host, reason)`, + `requests(id INTEGER PRIMARY KEY, route, fp, model, host, started, queued_ms, ttfb_ms, total_ms, status, streamed, prompt_tokens, cached_tokens, completion_tokens, err)`, + `host_health(ts, host, healthy, loaded_models)`, + `requests_daily(day, route, model, host, requests, errors, busy_ms, queued_ms, prompt_tokens, cached_tokens, completion_tokens, PRIMARY KEY(day, route, model, host))`. +2. **`Usage`** sums `requests` rows with `started >= since` **plus** `requests_daily` rows whose + `day >= since` (day = UTC midnight of `started`), grouped by the `by` column. `Errors` counts + `status >= 400`. A zero `since` (`time.Time{}`) means everything. Result order: by `Key`. +3. **`Prune(now, retention)`** moves every `requests` row with `started < now - retention` into + `requests_daily` (adding into the existing day row if there is one), deletes them, and returns + the number deleted. In one transaction. +4. Nothing here panics; every `sql` error is returned wrapped (`fmt.Errorf("store: …: %w", err)`). +5. `Loaded` in `HostHealth` is written as a JSON array; `nil` is written as `[]`. + +## Steps + +- [ ] **1. Branch and copy.** + +```sh +git switch master && git switch -c v1 +cp docs/plans/v1/_files/go.sum go.sum +mkdir -p internal/store && cp docs/plans/v1/_files/internal/store/store_test.go internal/store/ +``` + +- [ ] **2. Add the dependency.** In `go.mod`, the `require` becomes a block with both modules: + +``` +require ( + github.com/BurntSushi/toml v1.6.0 + modernc.org/sqlite v1.59.0 +) +``` + +Then `go mod download modernc.org/sqlite` (network, once) and `go mod verify`. Expected: +`all modules verified`. If `go mod tidy` wants to change `go.sum` or add `// indirect` lines to +`go.mod`, let it, and stage the result; `go.sum` must end up a superset of the given file. + +- [ ] **3. See the test fail.** `go test ./internal/store/`. Expected: it does not compile. +- [ ] **4. Write `internal/store/store.go`.** `gofmt -w internal/store/`. +- [ ] **5. See the test pass.** `go test -race -count=1 ./internal/store/`. Expected: `ok`. The + first compile of the driver takes a minute or two. +- [ ] **6. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **7. Log and commit.** Row `v1/01-store`. + +```sh +git add go.mod go.sum internal/store docs/implementer-log.md +git commit +``` + +## Done when + +- `go test -race -count=1 ./internal/store/` is `ok`; `make gate` prints `gate: ok`. +- `cmp internal/store/store_test.go docs/plans/v1/_files/internal/store/store_test.go` prints nothing. + +## Stop and report if + +- The module cannot be downloaded, or `go vet` rejects the driver on this Go version. diff --git a/docs/plans/v1/02-fingerprint-config.md b/docs/plans/v1/02-fingerprint-config.md new file mode 100644 index 0000000..926460d --- /dev/null +++ b/docs/plans/v1/02-fingerprint-config.md @@ -0,0 +1,86 @@ +# v1 task 02: the conversation fingerprint; config additions + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the conversation fingerprint and the v1 config keys` + +## Goal + +Two small things. `fingerprint.Of` turns a chat-completions body into a stable key for the +conversation it belongs to (`PLAN.md` §4a), and `config` learns `db`, `lease_idle` and +`retention`, with durations that accept a `d` suffix. + +## Context + +OpenCode and Hermes send no session id. A conversation's system prompt and its **first user +message** do not change from turn to turn, so hashing those two identifies the conversation +without client support. Only the first 4 KiB of each is hashed, so a huge first message does +not make every turn slow, and a body with no user message has no fingerprint (`""`): such +requests fall back to the route-level lease (task 04). + +## Files + +- Copy: `internal/fingerprint/fingerprint_test.go`, `internal/config/config_v1_test.go` +- Create: `internal/fingerprint/fingerprint.go` +- Modify: `internal/config/config.go`, `docs/implementer-log.md` + +## Interfaces + +`internal/fingerprint`, package `fingerprint`: + +```go +// Of returns the lowercase hex SHA-256 of the system prompt and the first user message of a +// chat-completions body (first 4 KiB of each, joined with "\n"), or "" when the body is not a +// JSON object with a "messages" array containing a user message. +func Of(body []byte) string +``` + +Rules the tests check: + +1. Decode `{"messages":[{"role":…,"content":…}, …]}`. `content` is either a string or an array + of parts `[{"type":"text","text":"…"}, …]`; for an array, join the `text` of the text parts + with `""` (other part types are ignored). Use `json.Unmarshal` into a struct with + `Content json.RawMessage`, then decide. +2. System prompt = content of the **first** message with `role == "system"` (or `""` if none). + First user message = content of the **first** message with `role == "user"`; **no user + message → return `""`**. Not JSON, or no `messages` → `""`. +3. Truncate each of the two strings to its first 4096 **bytes**, hash `system + "\n" + user` + with `crypto/sha256`, return `hex.EncodeToString`. + +`internal/config` gains, in `Config`: + +```go +DB string `toml:"db"` // default "crossbar.db"; empty string is an error (field "db") +LeaseIdle Duration `toml:"lease_idle"` // default 30m; less than 1m is an error (field "lease_idle") +Retention Duration `toml:"retention"` // default 180d; less than 1d is an error (field "retention") +``` + +and `Duration.UnmarshalText` accepts an integer followed by `d` (`"7d"` = 7 × 24 h) **in addition +to** `time.ParseDuration` syntax. Exactly: if the text matches `^[0-9]+d$`, multiply; otherwise +`time.ParseDuration`. `"1.5d"`, `"d"`, `"1d2h"` are errors. Validation order of the new fields: +after `queue_max`, before `hosts`. + +## Steps + +- [ ] **1. Copy.** + +```sh +git switch v1 +mkdir -p internal/fingerprint +cp docs/plans/v1/_files/internal/fingerprint/fingerprint_test.go internal/fingerprint/ +cp docs/plans/v1/_files/internal/config/config_v1_test.go internal/config/ +``` + +- [ ] **2. See them fail.** `go test ./internal/fingerprint/ ./internal/config/`. Expected: compile errors. +- [ ] **3. Write `fingerprint.go`; extend `config.go`.** `gofmt -w internal/`. +- [ ] **4. See them pass.** `go test -race -count=1 ./internal/fingerprint/ ./internal/config/`. Expected: both `ok` (the v0 config tests must still pass). +- [ ] **5. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **6. Log and commit.** Row `v1/02-fingerprint-config`. + +```sh +git add internal/fingerprint internal/config docs/implementer-log.md +git commit +``` + +## Done when + +- Both packages pass with `-race`; `make gate` prints `gate: ok`; both copied files are byte-identical to `_files/`. diff --git a/docs/plans/v1/03-limiter-choose.md b/docs/plans/v1/03-limiter-choose.md new file mode 100644 index 0000000..6a01f6f --- /dev/null +++ b/docs/plans/v1/03-limiter-choose.md @@ -0,0 +1,94 @@ +# v1 task 03: the per-(host, model) limiter and the chooser + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the per-host-model limiter and the host chooser` + +## Goal + +Two pure packages. `limiter` hands out at most `parallel` concurrent slots per (host, model) and +lets at most `queue_max` requests wait in line, first come first served. `choose` picks the host +for a **new** lease: most free slots × weight, ties to the shortest queue, then list order. +This is `PLAN.md` §4 (`choose`) and §6 (concurrency). + +## Context + +A `llama-server` router with `parallel = 4` and unified KV falls over when a fifth request +arrives ("Context size has been exceeded"). The limiter is where that is prevented: the fifth +request waits, the ninth (with `queue_max = 4`) is refused at once so the client can retry +elsewhere. Waiting is cancellable (the client may go away) and must leak neither a slot nor a +queue place. + +## Files + +- Copy: `internal/limiter/limiter_test.go`, `internal/choose/choose_test.go` +- Create: `internal/limiter/limiter.go`, `internal/choose/choose.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +`internal/limiter`, package `limiter`: + +```go +var ErrQueueFull = errors.New("queue full") + +type Limiter struct { /* private: mutex, map[(host, model)]*pair */ } +func New() *Limiter +func (l *Limiter) Configure(host, model string, parallel, queueMax int) +// Acquire returns when a slot is held. release gives it back (idempotent: a second call is a +// no-op). waited is how long the caller sat in the queue. Errors: ErrQueueFull immediately +// when queueMax waiters are already queued; ctx.Err() if ctx ends while waiting. +func (l *Limiter) Acquire(ctx context.Context, host, model string) (release func(), waited time.Duration, err error) +func (l *Limiter) InFlight(host, model string) int +func (l *Limiter) Queued(host, model string) int +func (l *Limiter) FreeSlots(host string) int // sum over the host's configured models of parallel - inflight (never below 0); 0 for an unknown host +``` + +Rules the tests check: + +1. An unconfigured (host, model) behaves as `parallel = 1, queueMax = 0`. +2. FIFO: waiters get slots in arrival order. Suggested shape: a mutex, `inflight`, and a slice + of waiter channels; `release` pops the head waiter (if any) and hands the slot over without + ever decrementing `inflight`, else decrements. A waiter whose ctx ends removes itself from + the queue under the mutex; if the slot was handed to it in the same instant, it releases it. +3. `release` idempotent via `sync.Once`. +4. Never panic; all methods safe for concurrent use. + +`internal/choose`, package `choose`: + +```go +type Info struct { + Healthy, Draining, Loaded, CanServe bool // Loaded: model is resident; CanServe: config lists the model + Free, Queued int + Weight float64 +} +// Best returns the candidate with the highest Free*Weight among those that are healthy, not +// draining and known (info ok) — preferring hosts with Loaded over merely CanServe; ties go to +// the lowest Queued, then to candidate order. A host with Free == 0 is still eligible (it will +// queue). ok is false when nothing is eligible. +func Best(candidates []string, info func(host string) (Info, bool)) (string, bool) +``` + +Rules: two passes — first over eligible candidates with `Loaded`, then, if none, over eligible +candidates with `CanServe`. Score `float64(Free) * Weight`. Compare with `>`; on equality prefer +lower `Queued`; on equality keep the earlier candidate. + +## Steps + +- [ ] **1. Copy.** `git switch v1`; `mkdir -p internal/limiter internal/choose`; copy both tests from `docs/plans/v1/_files/internal/…`. +- [ ] **2. See them fail** (compile). **3. Write both packages.** `gofmt -w internal/`. +- [ ] **4. See them pass.** `go test -race -count=1 ./internal/limiter/ ./internal/choose/`. The + limiter tests are timing-based with generous margins; run them three times: `-count=3`. +- [ ] **5. Run the gate.** `make gate`. **6. Log and commit.** Row `v1/03-limiter-choose`. + +```sh +git add internal/limiter internal/choose docs/implementer-log.md +git commit +``` + +## Done when + +- Both packages pass `-race -count=3`; `make gate` prints `gate: ok`; copied files byte-identical. + +## Stop and report if + +- A limiter test fails only sometimes: report which and how often; do not loosen it. diff --git a/docs/plans/v1/04-lease.md b/docs/plans/v1/04-lease.md new file mode 100644 index 0000000..6a6bf5a --- /dev/null +++ b/docs/plans/v1/04-lease.md @@ -0,0 +1,117 @@ +# v1 task 04: the lease table + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the sticky lease table` + +## Goal + +The heart of crossbar: `lease.Table` remembers which host each conversation is on and keeps it +there. A lease moves only when its host is unhealthy, when it has been idle longer than +`lease_idle`, or when an operator releases or pins the route. It is written through to the store +on every change and loaded back at start, so a restart does not reshuffle sessions. +`PLAN.md` §5, §4 step 3. + +## Context + +Why sticky: a `llama-server` prompt cache is per process; moving a 200k-token conversation to +another host costs minutes of re-prefill. A "better" host appearing is never a reason to move. +A pin ("project A goes to titan right now") is an operator decision and outranks everything, +including health: a pinned host that is down yields an error, not a silent move. + +## Files + +- Copy: `internal/lease/lease_test.go` +- Create: `internal/lease/lease.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +`internal/lease`, package `lease`: + +```go +var ( + ErrNoHost = errors.New("lease: no usable host") + ErrPinnedDown = errors.New("lease: pinned host is not healthy") + ErrUnknownHost = errors.New("lease: unknown host") +) + +type Key struct{ Route, FP, Model string } +type Lease struct { + Key + Host string + State store.State + Created, LastUsed time.Time +} +type Persister interface { // *store.Store satisfies it + SaveLease(store.Lease) error + DeleteLease(route, fp, model string) error + ListLeases() ([]store.Lease, error) + RecordEvent(store.LeaseEvent) error +} +type Hosts interface { + Healthy(name string) bool + Draining(name string) bool +} +type Chooser interface { + Choose(candidates []string, model string) (string, bool) +} +type Table struct { /* private: mutex, leases map[Key]*Lease, pins map[route]host, p, hosts, choose, idle */ } + +func New(p Persister, h Hosts, c Chooser, idle time.Duration) (*Table, error) // loads p.ListLeases(): rows with FP=="" && Model=="" && State==Pinned are pins +func (t *Table) Acquire(k Key, candidates []string, now time.Time) (host string, reused bool, err error) +func (t *Table) ExpireIdle(now time.Time) int // removes leases with now - LastUsed > idle (not pins); events ReasonIdle; returns how many +func (t *Table) Pin(route, host string, now time.Time) error // host must be in the candidates of at least one existing lease of the route, or in a candidate list seen for that route; else ErrUnknownHost +func (t *Table) Unpin(route string) +func (t *Table) Pinned(route string) string // "" if not pinned +func (t *Table) Release(route string) int // drops all leases (not the pin) of the route; events ReasonRelease; returns how many +func (t *Table) Snapshot() []Lease // copies, sorted by Route, FP, Model; pins excluded +``` + +`Acquire`, in this order: + +1. **Pinned route** (`pins[k.Route]` set): if `hosts.Healthy(pin)` → host = pin; if no lease for + `k` exists, create one (state `Active`, event `ReasonPin` only if this is the first lease + created under this pin for this key… keep it simple: event `ReasonNew` with `ToHost = pin`); + return `(pin, existed, nil)`. If the pin is not healthy → `ErrPinnedDown`. +2. **Existing lease for `k`**: if its host is healthy → touch `LastUsed = now`, save, return + `(host, true, nil)`. (A draining host still serves its existing leases.) If not healthy → + record `ReasonUnhealthy` (`FromHost` = old host) and fall through to choose, excluding that + host. +3. **Inherit**: if `k.FP != ""` and a lease for `Key{k.Route, "", k.Model}` exists on a healthy + host, create `k`'s lease on that host, save, return `(host, true, nil)`. +4. **Choose**: candidates minus unhealthy minus draining → `choose.Choose(filtered, k.Model)`. + None → `ErrNoHost` (nothing is created). Else create the lease (`Created = LastUsed = now`), + save, record `ReasonNew` (or the `ReasonUnhealthy` event from step 2 instead, with `ToHost` + filled), return `(host, false, nil)`. + +`Pin` records `ReasonPin` (`ToHost` = host) and saves a pin row +(`store.Lease{Route: route, FP: "", Model: "", Host: host, State: store.Pinned}`); existing +leases of the route on other hosts are **deleted** (event `ReasonPin` per lease) so the next turn +lands on the pin. `Unpin` deletes the pin row and records `ReasonRelease`; existing leases stay. +`Pin` to a host that no candidate list for that route has ever contained → `ErrUnknownHost` +(keep a `map[route]map[host]bool` of candidates seen in `Acquire`; at load time, hosts of stored +leases count as seen). + +All persister errors are returned; nothing is left half-changed in memory when a save fails +(apply to memory after the save succeeds). + +## Steps + +- [ ] **1. Copy.** `git switch v1`; `mkdir -p internal/lease`; `cp docs/plans/v1/_files/internal/lease/lease_test.go internal/lease/`. + Read `TestPinAndUnpin` and `TestFingerprintInheritsRouteLease` twice: they are the rules above as stories. +- [ ] **2. See it fail** (compile). **3. Write `lease.go`.** `gofmt -w internal/lease/`. +- [ ] **4. See it pass.** `go test -race -count=1 ./internal/lease/`. +- [ ] **5. Run the gate.** `make gate`. **6. Log and commit.** Row `v1/04-lease`. + +```sh +git add internal/lease docs/implementer-log.md +git commit +``` + +## Done when + +- `go test -race -count=1 ./internal/lease/` is `ok`; `make gate` prints `gate: ok`; the copied test is byte-identical. + +## Stop and report if + +- A test expects an event sequence you cannot produce under the rules above: quote the test and the rule that conflict. diff --git a/docs/plans/v1/05-proxy.md b/docs/plans/v1/05-proxy.md new file mode 100644 index 0000000..f98d913 --- /dev/null +++ b/docs/plans/v1/05-proxy.md @@ -0,0 +1,122 @@ +# v1 task 05: the proxy uses leases, the limiter, the SSE tee and the store + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Route by lease, queue per host and model, record every request` + +## Goal + +`internal/proxy` becomes the v1 proxy: route from path **or** header, fingerprint the body, +acquire a lease, take a limiter slot (queue or 503), forward with streaming, tee the response +to read `usage`/`timings`, mark hosts down on failure, and record one `store.Request` per +request. `PLAN.md` §4, §6, §7a. + +## Context + +The v0 proxy stays the skeleton of this one: `SplitRoute`, the ordered error answers, the model +peek, `httputil.ReverseProxy` with `FlushInterval: -1`, the status recorder with a checked +`Flush`, the log line. What changes is who picks the host and what happens around the forward. +Two new response headers make the behaviour observable: `X-Crossbar-Host` (existing) and +`X-Crossbar-Lease: new|reused`. The given test drives the whole handler over real HTTP with a +real `store`, `lease.Table`, `limiter` and `health.Table`. + +## Files + +- Copy (**replaces** v0's file): `internal/proxy/proxy_test.go` +- Modify: `internal/proxy/proxy.go` (split into more files if it passes 400 lines: `proxy.go`, `tee.go`, `hosts.go`), `docs/implementer-log.md` +- Keep: `internal/proxy/recorder_test.go` from v0.1 — it must still pass. Its `proxy.New(cfg, h, nil)` + call no longer compiles, so **this is the one given test you edit**: change that call to + `proxy.New(cfg, h, nil, nil, nil, nil)` and nothing else; `New` must accept nils for + `leases`, `lim`, `rec` and then behave like v0 (first healthy host, no queue, no recording). + Say so in Deviations. + +## Interfaces + +`internal/proxy`, package `proxy` (v0 names kept; additions): + +```go +const ( + MaxBody = 16 << 20 + HostHeader = "X-Crossbar-Host" + LeaseHeader = "X-Crossbar-Lease" // "new" or "reused" + RouteHeader = "X-Crossbar-Route" // client may name the route here instead of the path +) +type Health interface { Get(name string) (health.Status, bool); MarkDown(name, reason string) } +type Recorder interface { RecordRequest(store.Request) error } // *store.Store satisfies it + +// Hosts adapts the health table and config for the lease table, and holds the drain set. +type Hosts struct { /* private */ } +func HostView(h *health.Table, cfg *config.Config) *Hosts +func (h *Hosts) Healthy(name string) bool +func (h *Hosts) Draining(name string) bool +func (h *Hosts) SetDraining(name string, on bool) + +// Chooser adapts config, health and limiter to lease.Chooser using choose.Best: +// Info{Healthy, Draining: false (the lease table already filtered), Loaded: model in Loaded, +// CanServe: cfg.Serves, Free: lim.FreeSlots(host), Queued: lim.Queued(host, model), Weight}. +func Chooser(cfg *config.Config, h *health.Table, l *limiter.Limiter) lease.Chooser + +func New(cfg *config.Config, h Health, leases *lease.Table, lim *limiter.Limiter, rec Recorder, log *slog.Logger) *Handler +func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) +func SplitRoute(path string) (route, rest string, ok bool) +``` + +`ServeHTTP`, in order (every error answer is JSON `{"error":"…"}` as in v0): + +1. **Route.** `hdr := r.Header.Get(RouteHeader)`. If `hdr != ""`: the path is used **whole** as + `rest` (it must then start with `/v1/` or be `/health` or `/props`); if the path *also* starts + with a known route name and it differs from `hdr` → **400** `conflicting route`. If `hdr == ""`: + `SplitRoute` as in v0 (400 `missing route`). Unknown route (either source) → 404 + `unknown route`. Disallowed `rest` → 404 `not found`. +2. **Peek** (v0 rule): body up to `MaxBody` → 413; `model` from the body or the route default; + `fp := fingerprint.Of(body)` (GET/HEAD → `""`). +3. **Lease.** `host, reused, err := leases.Acquire(lease.Key{route, fp, model}, rt.Hosts, time.Now())`. + `ErrNoHost` → **503** `no healthy host`; `ErrPinnedDown` → **503** `pinned host down`. +4. **Slot.** `release, waited, err := lim.Acquire(r.Context(), host, model)`. `ErrQueueFull` → **503** + `queue full`; ctx error → **499**-style: just return (the client left; log status 499). + `defer release()`. +5. **Forward** as in v0 (`Rewrite`, `FlushInterval: -1`, `ModifyResponse` sets `HostHeader` and + `LeaseHeader`, `ErrorHandler` marks down + 502 with host). **Tee:** in `ModifyResponse`, wrap + `resp.Body` in a reader that passes every byte through unchanged and, when + `Content-Type` starts with `text/event-stream`, scans complete `data: ` lines for a JSON object + with `usage` and/or `timings`, remembering the **last** one seen; for non-streamed JSON + answers, remember the whole body's `usage`/`timings` (bounded: keep at most 1 MiB for the + parse; beyond that, record no tokens). The scanner must not hold data back: `Read` returns + what the upstream returned. +6. **Record**, after the upstream body is closed (the tee's `Close`, or the error handler): + `store.Request{Route, FP: fp, Model, Host, Started, QueuedMs: waited, TTFBMs (first byte of + the response head), TotalMs, Status, Streamed, PromptTokens: usage.prompt_tokens (or + timings.prompt_n), CachedTokens: timings.cache_n, CompletionTokens: usage.completion_tokens + (or timings.predicted_n), Err}`. Also record the 503/502 cases (Status set, no tokens). Do it + from the request goroutine after `rp.ServeHTTP` returns, so tests that read the store right + after the response see the row; if the tee cannot tell that the body closed, record what + you have when `ServeHTTP` returns. A `rec` error is logged, never returned to the client. +7. When `leases == nil` (v0.1 compatibility path used by `recorder_test.go`): choose with the + v0 `Choose` rule, skip the limiter and the store, still set `HostHeader`. +8. Log line as v0, adding `lease=new|reused`, `queued_ms`, `fp` (**first 8 hex chars only**). + Never the body. + +## Steps + +- [ ] **1. Copy (replace).** `git switch v1`; `cp docs/plans/v1/_files/internal/proxy/proxy_test.go internal/proxy/proxy_test.go`. + Edit the one call in `internal/proxy/recorder_test.go` as described above. +- [ ] **2. See it fail** (compile). **3. Write the code.** `gofmt -w internal/proxy/`. +- [ ] **4. See it pass.** `go test -race -count=1 ./internal/proxy/`. `TestDifferentConversationsSpreadByFreeSlots` + and `TestQueueFullIs503` are timing-based with generous margins; run `-count=3`. +- [ ] **5. Run the gate.** `make gate`. `cmd/crossbar` will not compile until task 06 — if `go vet ./...` + fails only in `cmd/crossbar/main.go` because of the new `New` signature, change that one call + to pass `nil, nil, nil` for the new arguments (task 06 wires it properly) and say so in Deviations. +- [ ] **6. Log and commit.** Row `v1/05-proxy`. + +```sh +git add internal/proxy cmd/crossbar docs/implementer-log.md +git commit +``` + +## Done when + +- `go test -race -count=3 ./internal/proxy/` is `ok`; `make gate` prints `gate: ok`; + `cmp internal/proxy/proxy_test.go docs/plans/v1/_files/internal/proxy/proxy_test.go` prints nothing. + +## Stop and report if + +- `TestAccountingRowsFromUsageAndTimings` fails on the token sums while the stream test passes: quote the recorded row. diff --git a/docs/plans/v1/06-admin-main.md b/docs/plans/v1/06-admin-main.md new file mode 100644 index 0000000..044ce7e --- /dev/null +++ b/docs/plans/v1/06-admin-main.md @@ -0,0 +1,121 @@ +# v1 task 06: admin v1 and the wiring + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Admin: leases, pin, release, drain, usage, metrics; wire the store into crossbar` + +## Goal + +Operators see and steer the system: the hosts view gains slots and drain state, the routes view +shows leases and pins, `POST` endpoints pin/release a route and drain a host, `/usage` answers +the accounting questions, `/metrics` exposes them to Prometheus. `cmd/crossbar` opens the store, +builds the lease table and limiter, runs idle expiry and pruning, and shuts down cleanly. +`PLAN.md` §7, §7a. + +## Files + +- Copy (**replaces** v0's): `internal/admin/admin_test.go`, `example.toml` +- Modify: `internal/admin/admin.go` (split if over 400 lines), `cmd/crossbar/main.go`, `docs/implementer-log.md` + +## Interfaces + +`internal/admin`, package `admin`: + +```go +type Hosts interface { All() map[string]health.Status } +type Drainer interface { Draining(name string) bool; SetDraining(name string, on bool) } // *proxy.Hosts satisfies it + +type HostView struct { + Healthy bool `json:"healthy"` + Loaded []string `json:"loaded"` // never null + LastOK string `json:"last_ok"` // RFC 3339 UTC or "" + LastErr string `json:"last_err"` + FreeSlots int `json:"free_slots"` // lim.FreeSlots(host) + InFlight int `json:"in_flight"` // sum over the host's configured models + Queued int `json:"queued"` // same + Draining bool `json:"draining"` +} +type LeaseView struct { FP, Model, Host, State, Created, LastUsed string } // json tags: fp, model, host, state, created, last_used (RFC 3339 UTC) +type RouteView struct { + Hosts []string `json:"hosts"` + DefaultModel string `json:"default_model"` + Pinned string `json:"pinned"` // "" when not pinned + Leases []LeaseView `json:"leases"` // never null +} +func Handler(cfg *config.Config, h Hosts, lt *lease.Table, lim *limiter.Limiter, st *store.Store, d Drainer) http.Handler +``` + +Endpoints (all JSON unless said; errors `{"error":"…"}`; wrong method → 405 with `Allow`): + +- `GET /_crossbar/hosts` → `map[string]HostView`. +- `GET /_crossbar/routes` → `map[string]RouteView` from config + `lt.Snapshot()` + `lt.Pinned`. +- `POST /_crossbar/routes/{route}` body `{"host":"…","pin":true}` → `lt.Pin(route, host, now)`; + `{"release":true}` → `lt.Release(route)` **and** `lt.Unpin(route)`. Unknown route → 404; + `lt.Pin` returning `ErrUnknownHost` → 404; not JSON, neither form, or both forms at once, + or `pin` without `host` → 400. Answer `{"ok":true}` (plus `"released": n` for a release). + `GET` on this path → 405. +- `POST /_crossbar/hosts/{host}` body `{"drain":true|false}` → `d.SetDraining`; unknown host + (not in config) → 404; bad body → 400. `{"ok":true}`. +- `GET /_crossbar/usage?since=…&by=route|model|host` → `[]store.UsageRow` (JSON array; sorted by + key). `by` defaults to `route`; `since` is either RFC 3339 or a duration like `24h`/`7d` + (meaning `now - d`); absent = all time; anything else → 400. With `Accept: text/plain`, a + fixed-width table with a header line containing `key requests errors busy_ms queued_ms + prompt cached completion cache_hit` and one line per row (`cache_hit` as `0.80`). +- `GET /_crossbar/metrics` → `text/plain; version=0.0.4`, computed on request. Request counters + need (route, host, status), which `Usage` (one key) cannot give, so add **one** method to + `internal/store` — the only change to that package allowed in this task: + ```go + type StatusCount struct { Route, Host string; Status int; Count int64 } + func (s *Store) StatusCounts(since time.Time) ([]StatusCount, error) // from `requests` only; rolled-up days are not in it, say so in a comment + ``` + Then emit, in this order: + ``` + # TYPE crossbar_requests_total counter + crossbar_requests_total{route="…",host="…",status="…"} N (one line per StatusCount) + # TYPE crossbar_prompt_tokens_total counter + crossbar_prompt_tokens_total{route="…"} N (Usage(zero, ByRoute)) + # TYPE crossbar_cached_tokens_total counter + # TYPE crossbar_completion_tokens_total counter + # TYPE crossbar_queue_wait_ms_total counter crossbar_queue_wait_ms_total{route="…"} N + # TYPE crossbar_host_healthy gauge crossbar_host_healthy{host="…"} 0|1 + # TYPE crossbar_host_free_slots gauge + # TYPE crossbar_host_in_flight gauge + # TYPE crossbar_host_queued gauge + ``` + Label values escaped (`"` and `\`), lines sorted, no trailing spaces. + +`cmd/crossbar/main.go`: + +- After `config.Load`: `store.Open(cfg.DB)` (error → exit 1 `crossbar: …`); `defer st.Close()`. +- `hosts := proxy.HostView(table, cfg)`; `lim := limiter.New()` configured for every host/model + from config with `cfg.QueueMax`; `leases, err := lease.New(st, hosts, proxy.Chooser(cfg, table, lim), cfg.LeaseIdle.Duration)`. +- `proxy.New(cfg, table, leases, lim, st, log)`; `admin.Handler(cfg, table, leases, lim, st, hosts)`. +- Background loops until ctx is done: every minute `leases.ExpireIdle(time.Now())`; every hour + `st.Prune(time.Now(), cfg.Retention.Duration)`; after every poll round the health table's + observations are recorded with `st.RecordHostHealth` — do this from a goroutine that every + `poll_interval` reads `table.All()` and writes one row per host. +- Shutdown as v0, then `st.Close()`. + +## Steps + +- [ ] **1. Copy (replace).** `git switch v1`; copy `admin_test.go` and `example.toml` from `docs/plans/v1/_files/`. +- [ ] **2. See it fail** (compile). **3. Write the code** (admin, the store's `StatusCounts`, main). `gofmt -w .` +- [ ] **4. See it pass.** `go test -race -count=1 ./...`. +- [ ] **5. Build and run for three seconds.** + +```sh +make build +timeout --preserve-status --signal=TERM 3 bin/crossbar -config example.toml; echo "exit=$?" +ls -la crossbar.db* && rm -f crossbar.db crossbar.db-wal crossbar.db-shm +``` + +Expected: `listening`, `shutting down`, `exit=0`; the SQLite file was created (then removed). +- [ ] **6. Run the gate.** `make gate`. **7. Log and commit.** Row `v1/06-admin-main`. + +```sh +git add internal/admin internal/store cmd/crossbar example.toml docs/implementer-log.md +git commit +``` + +## Done when + +- All tests pass with `-race`; the three-second run exits 0 and created the db; `make gate` prints `gate: ok`; both copied files byte-identical. diff --git a/docs/plans/v1/07-smoke-readme.md b/docs/plans/v1/07-smoke-readme.md new file mode 100644 index 0000000..5b03179 --- /dev/null +++ b/docs/plans/v1/07-smoke-readme.md @@ -0,0 +1,53 @@ +# v1 task 07: the smoke run and the README + +**Branch:** `v1` (run `git switch v1`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Smoke run for v1; README for leases, admin and accounting` + +## Goal + +Prove v1 end to end with the given `tools/smoke.sh` — leases, header route, pin, queue, drain, +failover and recovery, streaming with the usage chunk intact, usage and metrics — and bring the +README up to date. + +## Files + +- Copy (**replaces** v0's): `tools/smoke.sh`, `cmd/fakeupstream/main.go` +- Modify: `README.md`, `docs/implementer-log.md` + +## Steps + +- [ ] **1. Copy and run.** + +```sh +git switch v1 +cp docs/plans/v1/_files/tools/smoke.sh tools/ +cp docs/plans/v1/_files/cmd/fakeupstream/main.go cmd/fakeupstream/ +make smoke +``` + +Expected last line: `smoke: ok (stream spread N ms)`, N ≥ 600. The script prints which numbered +check failed and crossbar's log. A failure is a defect in tasks 01–06 **or in the script**: fix +code only when a rule from an earlier task was broken; if the script's expectation contradicts a +task rule, stop and report which. + +- [ ] **2. Update `README.md`.** Keep the seven v0 sections; change: the intro (leases, not "first + healthy host"); `## Configure` gets `db`, `lease_idle`, `retention`, `queue_max` in the table and + the new `example.toml`; `## Point clients at it` adds the `X-Crossbar-Route` header + alternative; `## Inspect` becomes `## Operate` and documents all six endpoints with one example + each (`GET hosts`, `GET routes`, `POST routes/{route}` pin and release, `POST hosts/{host}` + drain, `GET usage` JSON and text, `GET metrics`), taken from the smoke run; `## What v1 does not + do`: context-size guard, wake-on-LAN, Tailscale identity, `/slots` — see `PLAN.md` v2. +- [ ] **3. Run the gate.** `make gate`. **4. Log and commit.** Row `v1/07-smoke-readme`, with the smoke line in Notes. + +```sh +git add tools/smoke.sh cmd/fakeupstream/main.go README.md docs/implementer-log.md +git commit +``` + +## Done when + +- `make smoke` prints `smoke: ok (…)`; `make gate` prints `gate: ok`; both copied files byte-identical. + +## Stop and report if + +- `make smoke` fails twice in the same way. diff --git a/docs/plans/v1/README.md b/docs/plans/v1/README.md new file mode 100644 index 0000000..0365496 --- /dev/null +++ b/docs/plans/v1/README.md @@ -0,0 +1,61 @@ +# v1 implementation plan: leases, queueing, accounting + +> **For the implementing model:** do not work from this file. The owner gives you one task file at +> a time (`01-…` to `07-…`). This file is the index for the owner and the reviewer. + +**Goal:** `PLAN.md` §4–§7a. Every conversation gets a sticky lease on one host (chosen by free +slots × weight when it starts), requests queue per (host, model) instead of overflowing a +router, an operator can pin a route or drain a host, and SQLite keeps the leases and an accounting +log that answers "which session used which host and model, for how long, at what cache-hit rate". + +**Architecture:** five new packages — `store` (SQLite, `modernc.org/sqlite`), `fingerprint` +(conversation key), `choose` (the scoring rule), `limiter` (per-(host, model) slots + bounded +FIFO), `lease` (the sticky table, persisted through `store`) — and v1 versions of `proxy`, +`admin`, `config` and `cmd/crossbar`. The proxy tees streamed responses through an SSE scanner +to read the final `usage`/`timings` chunk; it never buffers or alters the stream. + +**How this plan was made:** acceptance tests first, from `PLAN.md`; no reference implementation. +Every given test file was compiled against a panic-only skeleton of the interfaces named in the +tasks (`go vet ./...` clean), and nothing else was run. If a test turns out to be wrong, that is +the owner's finding: stop and report as `AGENTS.md` says; do not edit it. + +**Tech stack:** Go 1.26, stdlib, `github.com/BurntSushi/toml` v1.6.0, `modernc.org/sqlite` +v1.59.0 (pure Go; `go.sum` given). No other module. + +## Global constraints + +- Everything in `AGENTS.md`. Branch `v1`. One task, one fresh OpenCode session, one commit. +- Bodies are never logged or stored. Rows carry names, counts and timings only. +- Given files (tests, `example.toml`, `cmd/fakeupstream/main.go`, `tools/smoke.sh`, `go.sum`) + are copied and never edited. Some **replace** v0 files of the same name; the task says so. + +## Tasks + +| # | File | Delivers | Tests that define it | +|---|---|---|---| +| 01 | `01-store.md` | `internal/store`: SQLite state + accounting | `internal/store/store_test.go` | +| 02 | `02-fingerprint-config.md` | `internal/fingerprint`; `config` gains `db`, `lease_idle`, `retention`, day suffix | `fingerprint_test.go`, `config_v1_test.go` | +| 03 | `03-limiter-choose.md` | `internal/limiter`, `internal/choose` | `limiter_test.go`, `choose_test.go` | +| 04 | `04-lease.md` | `internal/lease`: the sticky table | `lease_test.go` | +| 05 | `05-proxy.md` | proxy v1: leases, queue, SSE tee, accounting, header route | `proxy_test.go` (replaces v0's) | +| 06 | `06-admin-main.md` | admin v1 (pin/release/drain/usage/metrics), `cmd/crossbar` wiring | `admin_test.go` (replaces), start/stop check | +| 07 | `07-smoke-readme.md` | `tools/smoke.sh` v1 run, README update | `make smoke` | + +## For the owner: running a task + +```sh +tools/run-plan.sh docs/plans/v1 # from a clean checkout on master +``` + +## For the reviewer: after task 07 + +1. `git log --oneline master..v1`: seven commits with the trailer. +2. Copied files byte-identical to `_files/`; `git diff ..v1 --stat -- PLAN.md AGENTS.md docs/plans` empty. +3. `make gate`, `make smoke`. +4. Read every source file against its task. Probe outside the tests: a lease whose host is + drained *and* unhealthy; `lease_idle` expiry while a request is in flight; a stream cut by the + client mid-way (the accounting row must still be written, with the status it had); a body + with `"messages"` that is not an array; two crossbars on the same `db` file; `Prune` while + requests are being recorded. +5. Every `.(` type assertion in `internal/` is the two-value form or on a value we constructed. +6. Findings under "Reviews" in `docs/implementer-log.md`, by fault (model / task / test). diff --git a/docs/plans/v1/_files/cmd/fakeupstream/main.go b/docs/plans/v1/_files/cmd/fakeupstream/main.go new file mode 100644 index 0000000..c222d38 --- /dev/null +++ b/docs/plans/v1/_files/cmd/fakeupstream/main.go @@ -0,0 +1,115 @@ +// 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: . +package main + +import ( + "encoding/json" + "flag" + "fmt" + "io" + "log" + "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") + flag.Parse() + + 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": 8192}, "total_slots": 2, "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) +} diff --git a/docs/plans/v1/_files/example.toml b/docs/plans/v1/_files/example.toml new file mode 100644 index 0000000..f780b12 --- /dev/null +++ b/docs/plans/v1/_files/example.toml @@ -0,0 +1,26 @@ +# crossbar example configuration (v1). Replace and the addresses with your own. +listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production +db = "crossbar.db" # SQLite: leases + accounting (WAL). /var/lib/crossbar/crossbar.db under systemd +poll_interval = "1s" # 60s in production; 1s makes the smoke run quick +lease_idle = "30m" # a conversation idle this long loses its host +retention = "180d" # per-request rows older than this are rolled up daily +queue_max = 1 # waiting places per (host, model) beyond `parallel`; 503 past that + +[hosts.alpha] +base_url = "http://127.0.0.1:18081" # e.g. http://straylight.:11434 +weight = 1.0 +models = { "ornith-1.5-35b-a3b" = { parallel = 1 }, "small-9b" = { parallel = 6 } } + +[hosts.beta] +base_url = "http://127.0.0.1:18082" # e.g. http://titan.:8081 +weight = 2.0 +models = { "ornith-1.5-35b-a3b" = { parallel = 2 } } + +# 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"] diff --git a/docs/plans/v1/_files/go.sum b/docs/plans/v1/_files/go.sum new file mode 100644 index 0000000..e3728f9 --- /dev/null +++ b/docs/plans/v1/_files/go.sum @@ -0,0 +1,52 @@ +github.com/BurntSushi/toml v1.6.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= +github.com/BurntSushi/toml v1.6.0 h1:dRaEfpa2VI55EwlIW72hMRHdWouJeRF7TPYhI+AUQjk= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= +github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= +github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +modernc.org/ccgo/v4 v4.35.0/go.mod h1:qrVGs9S3Sr2Ztcg9ve+kTAYMp5a3YvWjo+SoN06kJ5I= +modernc.org/ccgo/v4 v4.35.0 h1:F+TUsmw09QxLzmi3aeYYGxjAXarmZaKgj3mKQHNaA8w= +modernc.org/cc/v4 v4.29.2/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI= +modernc.org/cc/v4 v4.29.2 h1:h6+9ciCnPKutf4I03CvheAvDLX7+IHlqR6Iy6J+cgd8= +modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU= +modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v3 v3.1.5/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY= +modernc.org/gc/v3 v3.1.5 h1:21ldfPfRYE31Tb7B3mwAK8gy1AxP4+dKjrOQPfqakoc= +modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI= +modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks= +modernc.org/libc v1.75.7/go.mod h1:bO5o2ztHxBb2rjz0PgdHN0sSMw57CgxGFLZ3Qd/QpVQ= +modernc.org/libc v1.75.7 h1:o3DTP9/0p9pKmY2WCKQaySW6wIiZhNM7wc2lUoyhfew= +modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= +modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= +modernc.org/memory v1.12.1/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/memory v1.12.1 h1:nFMiWrpStgZczNl6XI9GnIk/rWhYIyHGUaR04pGbp9g= +modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sqlite v1.59.0/go.mod h1:+paeT2A3iPRHkQDwG7oA6Tk0zQd5woMEI8q7orfry8k= +modernc.org/sqlite v1.59.0 h1:X1es1GpqBlS/5T+vbM4HLUdaa8OtQx468DF2vrx+38A= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= diff --git a/docs/plans/v1/_files/internal/admin/admin_test.go b/docs/plans/v1/_files/internal/admin/admin_test.go new file mode 100644 index 0000000..f85cc4d --- /dev/null +++ b/docs/plans/v1/_files/internal/admin/admin_test.go @@ -0,0 +1,293 @@ +package admin_test + +// v1 admin: read the tables, pin/release a route, drain a host, usage rollups, metrics. + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/admin" + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +type fakeHosts struct { + st map[string]health.Status + draining map[string]bool +} + +func (f *fakeHosts) All() map[string]health.Status { return f.st } +func (f *fakeHosts) Healthy(n string) bool { return f.st[n].Healthy } +func (f *fakeHosts) Draining(n string) bool { return f.draining[n] } +func (f *fakeHosts) SetDraining(n string, on bool) { f.draining[n] = on } +func (f *fakeHosts) Choose(c []string, model string) (string, bool) { + for _, h := range c { + if f.st[h].Healthy && !f.draining[h] { + return h, true + } + } + return "", false +} + +type rig struct { + h http.Handler + store *store.Store + leases *lease.Table + hosts *fakeHosts +} + +func newRig(t *testing.T) *rig { + cfg, err := config.Parse(strings.NewReader(` +listen = "127.0.0.1:1" +[hosts.alpha] +base_url = "http://alpha:1" +models = { "m" = { parallel = 2 } } +[hosts.beta] +base_url = "http://beta:1" +models = { "m" = { parallel = 4 } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "m" +`)) + if err != nil { + t.Fatal(err) + } + st, err := store.Open(filepath.Join(t.TempDir(), "x.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = st.Close() }) + hosts := &fakeHosts{ + st: map[string]health.Status{ + "alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC)}, + "beta": {Healthy: false, LastErr: "HTTP 503"}, + }, + draining: map[string]bool{}, + } + lt, err := lease.New(st, hosts, hosts, 30*time.Minute) + if err != nil { + t.Fatal(err) + } + lim := limiter.New() + lim.Configure("alpha", "m", 2, 8) + lim.Configure("beta", "m", 4, 8) + return &rig{h: admin.Handler(cfg, hosts, lt, lim, st, hosts), store: st, leases: lt, hosts: hosts} +} + +func (r *rig) do(t *testing.T, method, path, body string, hdr ...string) *httptest.ResponseRecorder { + req := httptest.NewRequest(method, path, strings.NewReader(body)) + if body != "" { + req.Header.Set("Content-Type", "application/json") + } + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } + rec := httptest.NewRecorder() + r.h.ServeHTTP(rec, req) + return rec +} + +func TestHostsShowsSlotsAndDrain(t *testing.T) { + r := newRig(t) + rec := r.do(t, "GET", "/_crossbar/hosts", "") + if rec.Code != 200 { + t.Fatalf("%d %s", rec.Code, rec.Body.String()) + } + var out map[string]admin.HostView + if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { + t.Fatal(err) + } + a := out["alpha"] + if !a.Healthy || a.FreeSlots != 2 || a.InFlight != 0 || a.Queued != 0 || a.Draining || a.LastOK != "2026-09-25T08:00:00Z" { + t.Errorf("alpha = %+v", a) + } + if b := out["beta"]; b.Healthy || b.LastErr != "HTTP 503" || b.FreeSlots != 4 || b.Loaded == nil { + t.Errorf("beta = %+v (loaded must be [] not null)", b) + } +} + +func TestRoutesShowsLeases(t *testing.T) { + r := newRig(t) + now := time.Date(2026, 9, 25, 9, 0, 0, 0, time.UTC) + if _, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "abc", Model: "m"}, []string{"alpha", "beta"}, now); err != nil { + t.Fatal(err) + } + rec := r.do(t, "GET", "/_crossbar/routes", "") + var out map[string]admin.RouteView + if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { + t.Fatalf("%v: %s", err, rec.Body.String()) + } + rv := out["r"] + if len(rv.Hosts) != 2 || rv.DefaultModel != "m" || rv.Pinned != "" { + t.Errorf("route view = %+v", rv) + } + if len(rv.Leases) != 1 || rv.Leases[0].FP != "abc" || rv.Leases[0].Host != "alpha" || rv.Leases[0].State != "active" || rv.Leases[0].LastUsed != "2026-09-25T09:00:00Z" { + t.Errorf("leases = %+v", rv.Leases) + } +} + +func TestPinReleaseDrain(t *testing.T) { + r := newRig(t) + rec := r.do(t, "POST", "/_crossbar/routes/r", `{"host":"beta","pin":true}`) + if rec.Code != 200 { + t.Fatalf("pin: %d %s", rec.Code, rec.Body.String()) + } + if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err == nil || h != "" { + // beta is unhealthy in the rig: a pin to a down host is honoured, not silently moved + t.Errorf("acquire on a route pinned to a down host: %q %v, want ErrPinnedDown", h, err) + } + rec = r.do(t, "GET", "/_crossbar/routes", "") + var out map[string]admin.RouteView + _ = json.Unmarshal(rec.Body.Bytes(), &out) + if out["r"].Pinned != "beta" { + t.Errorf("Pinned = %q after pin", out["r"].Pinned) + } + rec = r.do(t, "POST", "/_crossbar/routes/r", `{"release":true}`) + if rec.Code != 200 { + t.Fatalf("release: %d %s", rec.Code, rec.Body.String()) + } + if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err != nil || h != "alpha" { + t.Errorf("after release: %q %v, want alpha (the only healthy host)", h, err) + } + for _, tc := range []struct { + body string + want int + }{ + {`{"host":"nobody","pin":true}`, 404}, + {`{"pin":true}`, 400}, + {`not json`, 400}, + {`{"release":true,"pin":true,"host":"alpha"}`, 400}, + } { + if rec := r.do(t, "POST", "/_crossbar/routes/r", tc.body); rec.Code != tc.want { + t.Errorf("POST %s: %d, want %d (%s)", tc.body, rec.Code, tc.want, rec.Body.String()) + } + } + if rec := r.do(t, "POST", "/_crossbar/routes/nope", `{"release":true}`); rec.Code != 404 { + t.Errorf("unknown route: %d", rec.Code) + } + + rec = r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":true}`) + if rec.Code != 200 || !r.hosts.Draining("alpha") { + t.Fatalf("drain: %d %s draining=%v", rec.Code, rec.Body.String(), r.hosts.Draining("alpha")) + } + rec = r.do(t, "GET", "/_crossbar/hosts", "") + var hv map[string]admin.HostView + _ = json.Unmarshal(rec.Body.Bytes(), &hv) + if !hv["alpha"].Draining { + t.Errorf("hosts view must show draining") + } + if rec := r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":false}`); rec.Code != 200 || r.hosts.Draining("alpha") { + t.Errorf("undrain: %d draining=%v", rec.Code, r.hosts.Draining("alpha")) + } + if rec := r.do(t, "POST", "/_crossbar/hosts/nobody", `{"drain":true}`); rec.Code != 404 { + t.Errorf("unknown host: %d", rec.Code) + } +} + +func seedUsage(t *testing.T, st *store.Store) { + t0 := time.Now().UTC().Add(-time.Hour) + for i, r := range []store.Request{ + {Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 1000, PromptTokens: 100, CachedTokens: 80, CompletionTokens: 10}, + {Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 500, QueuedMs: 30, PromptTokens: 100, CachedTokens: 100, CompletionTokens: 5}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Status: 503, TotalMs: 1, Err: "queue full"}, + } { + r.Started = t0.Add(time.Duration(i) * time.Minute) + if err := st.RecordRequest(r); err != nil { + t.Fatal(err) + } + } +} + +func TestUsageJSONAndText(t *testing.T) { + r := newRig(t) + seedUsage(t, r.store) + rec := r.do(t, "GET", "/_crossbar/usage?by=route", "") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type")) + } + var rows []store.UsageRow + if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil { + t.Fatalf("%v: %s", err, rec.Body.String()) + } + if len(rows) != 2 { + t.Fatalf("rows = %+v", rows) + } + for _, row := range rows { + if row.Key == "r" && (row.Requests != 2 || row.CachedTokens != 180 || row.QueuedMs != 30) { + t.Errorf("r = %+v", row) + } + if row.Key == "r2" && (row.Requests != 1 || row.Errors != 1) { + t.Errorf("r2 = %+v", row) + } + } + rec = r.do(t, "GET", "/_crossbar/usage?by=host&since=24h", "", "Accept", "text/plain") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") { + t.Fatalf("text: %d %q", rec.Code, rec.Header().Get("Content-Type")) + } + body := rec.Body.String() + if !strings.Contains(body, "alpha") || !strings.Contains(body, "beta") || !strings.Contains(strings.ToLower(body), "cache") { + t.Errorf("text table = %q", body) + } + if rec := r.do(t, "GET", "/_crossbar/usage?by=colour", ""); rec.Code != 400 { + t.Errorf("bad by: %d", rec.Code) + } + if rec := r.do(t, "GET", "/_crossbar/usage?since=yesterday", ""); rec.Code != 400 { + t.Errorf("bad since: %d", rec.Code) + } + rec = r.do(t, "GET", "/_crossbar/usage?since=2026-09-25T00:00:00Z&by=model", "") + if rec.Code != 200 { + t.Errorf("RFC3339 since: %d %s", rec.Code, rec.Body.String()) + } +} + +func TestMetrics(t *testing.T) { + r := newRig(t) + seedUsage(t, r.store) + rec := r.do(t, "GET", "/_crossbar/metrics", "") + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") { + t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type")) + } + body := rec.Body.String() + for _, want := range []string{ + `# TYPE crossbar_requests_total counter`, + `crossbar_requests_total{route="r",host="alpha",status="200"} 2`, + `crossbar_requests_total{route="r2",host="beta",status="503"} 1`, + `crossbar_host_healthy{host="alpha"} 1`, + `crossbar_host_healthy{host="beta"} 0`, + `crossbar_host_free_slots{host="alpha"} 2`, + `crossbar_prompt_tokens_total{route="r"} 200`, + `crossbar_cached_tokens_total{route="r"} 180`, + `crossbar_queue_wait_ms_total{route="r"} 30`, + } { + if !strings.Contains(body, want) { + t.Errorf("metrics missing %q\n%s", want, body) + } + } +} + +func TestMethodsAndUnknown(t *testing.T) { + r := newRig(t) + for _, tc := range []struct { + method, path string + want int + }{ + {http.MethodPost, "/_crossbar/hosts", 405}, + {http.MethodDelete, "/_crossbar/routes", 405}, + {http.MethodGet, "/_crossbar/routes/r", 405}, + {http.MethodGet, "/_crossbar/nope", 404}, + {http.MethodPut, "/_crossbar/usage", 405}, + } { + rec := r.do(t, tc.method, tc.path, "") + if rec.Code != tc.want || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Errorf("%s %s = %d %q, want %d JSON", tc.method, tc.path, rec.Code, rec.Header().Get("Content-Type"), tc.want) + } + } +} diff --git a/docs/plans/v1/_files/internal/choose/choose_test.go b/docs/plans/v1/_files/internal/choose/choose_test.go new file mode 100644 index 0000000..1af398e --- /dev/null +++ b/docs/plans/v1/_files/internal/choose/choose_test.go @@ -0,0 +1,83 @@ +package choose_test + +import ( + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/choose" +) + +func infoFor(m map[string]choose.Info) func(string) (choose.Info, bool) { + return func(name string) (choose.Info, bool) { i, ok := m[name]; return i, ok } +} + +func TestMostFreeSlotsTimesWeightWins(t *testing.T) { + info := infoFor(map[string]choose.Info{ + "alpha": {Healthy: true, Loaded: true, CanServe: true, Free: 3, Weight: 1.0}, + "beta": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 2.0}, // 4 > 3 + "gamma": {Healthy: true, Loaded: true, CanServe: true, Free: 4, Weight: 0.5}, // 2 + }) + got, ok := choose.Best([]string{"alpha", "beta", "gamma"}, info) + if !ok || got != "beta" { + t.Errorf("got %q %v, want beta", got, ok) + } +} + +func TestTieGoesToShortestQueueThenListOrder(t *testing.T) { + info := infoFor(map[string]choose.Info{ + "alpha": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 3}, + "beta": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 1}, + "gamma": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 1}, + }) + if got, _ := choose.Best([]string{"alpha", "beta", "gamma"}, info); got != "beta" { + t.Errorf("tie on score: shortest queue wins, then list order; got %q", got) + } + if got, _ := choose.Best([]string{"gamma", "beta"}, info); got != "gamma" { + t.Errorf("full tie: first in list order wins; got %q", got) + } +} + +func TestLoadedBeatsMerelyCapable(t *testing.T) { + info := infoFor(map[string]choose.Info{ + "alpha": {Healthy: true, Loaded: false, CanServe: true, Free: 8, Weight: 4}, + "beta": {Healthy: true, Loaded: true, CanServe: true, Free: 1, Weight: 1}, + }) + got, ok := choose.Best([]string{"alpha", "beta"}, info) + if !ok || got != "beta" { + t.Errorf("a host that has the model loaded wins over one that would have to load it; got %q", got) + } +} + +func TestFallsBackToCapableHost(t *testing.T) { + info := infoFor(map[string]choose.Info{ + "alpha": {Healthy: true, Loaded: false, CanServe: true, Free: 1, Weight: 1}, + "beta": {Healthy: true, Loaded: false, CanServe: false, Free: 9, Weight: 9}, + }) + got, ok := choose.Best([]string{"beta", "alpha"}, info) + if !ok || got != "alpha" { + t.Errorf("only a host configured to serve the model may load it; got %q %v", got, ok) + } +} + +func TestSkipsUnhealthyDrainingUnknownAndFull(t *testing.T) { + info := infoFor(map[string]choose.Info{ + "down": {Healthy: false, Loaded: true, CanServe: true, Free: 9, Weight: 9}, + "drain": {Healthy: true, Draining: true, Loaded: true, CanServe: true, Free: 9, Weight: 9}, + "full": {Healthy: true, Loaded: true, CanServe: true, Free: 0, Weight: 9, Queued: 0}, + "ok": {Healthy: true, Loaded: true, CanServe: true, Free: 1, Weight: 1}, + }) + got, ok := choose.Best([]string{"down", "drain", "missing", "full", "ok"}, info) + if !ok || got != "ok" { + t.Errorf("got %q %v, want ok", got, ok) + } + // A full host is still better than nothing: it gets the request (it will queue). + got, ok = choose.Best([]string{"down", "full"}, info) + if !ok || got != "full" { + t.Errorf("with only a full host left it must still be chosen; got %q %v", got, ok) + } + if _, ok := choose.Best([]string{"down", "drain", "missing"}, info); ok { + t.Errorf("nothing usable must give ok=false") + } + if _, ok := choose.Best(nil, info); ok { + t.Errorf("empty candidates must give ok=false") + } +} diff --git a/docs/plans/v1/_files/internal/config/config_v1_test.go b/docs/plans/v1/_files/internal/config/config_v1_test.go new file mode 100644 index 0000000..9931c17 --- /dev/null +++ b/docs/plans/v1/_files/internal/config/config_v1_test.go @@ -0,0 +1,81 @@ +package config_test + +import ( + "strings" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const v1Base = ` +listen = "127.0.0.1:1" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.r] +hosts = ["a"] +` + +func TestV1Defaults(t *testing.T) { + c, err := config.Parse(strings.NewReader(v1Base)) + if err != nil { + t.Fatal(err) + } + if c.DB != "crossbar.db" { + t.Errorf("DB default = %q", c.DB) + } + if c.LeaseIdle.Duration != 30*time.Minute { + t.Errorf("LeaseIdle default = %v", c.LeaseIdle.Duration) + } + if c.Retention.Duration != 180*24*time.Hour { + t.Errorf("Retention default = %v", c.Retention.Duration) + } +} + +func TestV1Values(t *testing.T) { + c, err := config.Parse(strings.NewReader(` +db = "/var/lib/crossbar/crossbar.db" +lease_idle = "45m" +retention = "30d" +` + v1Base)) + if err != nil { + t.Fatal(err) + } + if c.DB != "/var/lib/crossbar/crossbar.db" || c.LeaseIdle.Duration != 45*time.Minute || c.Retention.Duration != 30*24*time.Hour { + t.Errorf("got db %q idle %v retention %v", c.DB, c.LeaseIdle.Duration, c.Retention.Duration) + } +} + +func TestDurationAcceptsDays(t *testing.T) { + var d config.Duration + for _, tc := range []struct { + in string + want time.Duration + }{ + {"1d", 24 * time.Hour}, {"7d", 7 * 24 * time.Hour}, {"90m", 90 * time.Minute}, {"2h30m", 150 * time.Minute}, + } { + if err := d.UnmarshalText([]byte(tc.in)); err != nil || d.Duration != tc.want { + t.Errorf("UnmarshalText(%q) = %v %v, want %v", tc.in, d.Duration, err, tc.want) + } + } + for _, bad := range []string{"1.5d", "d", "3 days", "1d2h"} { + if err := d.UnmarshalText([]byte(bad)); err == nil { + t.Errorf("UnmarshalText(%q) must fail", bad) + } + } +} + +func TestV1Validation(t *testing.T) { + for _, tc := range []struct{ name, text, field string }{ + {"empty db", "db = \"\"\n" + v1Base, "db"}, + {"lease_idle too short", "lease_idle = \"10s\"\n" + v1Base, "lease_idle"}, + {"retention too short", "retention = \"12h\"\n" + v1Base, "retention"}, + } { + _, 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) + } + } +} diff --git a/docs/plans/v1/_files/internal/fingerprint/fingerprint_test.go b/docs/plans/v1/_files/internal/fingerprint/fingerprint_test.go new file mode 100644 index 0000000..860d5cb --- /dev/null +++ b/docs/plans/v1/_files/internal/fingerprint/fingerprint_test.go @@ -0,0 +1,70 @@ +package fingerprint_test + +import ( + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/fingerprint" +) + +const conv1 = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Add a config loader."},{"role":"assistant","content":"Sure."},{"role":"user","content":"Now tests."}]}` +const conv1later = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Add a config loader."},{"role":"assistant","content":"Sure."},{"role":"user","content":"Now tests."},{"role":"assistant","content":"Done."},{"role":"user","content":"And docs."}]}` +const conv2 = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Fix the flaky test."}]}` +const conv3 = `{"model":"m","messages":[{"role":"system","content":"You are the project B assistant."},{"role":"user","content":"Add a config loader."}]}` + +func TestSameConversationSameKey(t *testing.T) { + a := fingerprint.Of([]byte(conv1)) + b := fingerprint.Of([]byte(conv1later)) + if a == "" || a != b { + t.Errorf("later turns of one conversation must keep the key: %q vs %q", a, b) + } + if len(a) != 64 || strings.Trim(a, "0123456789abcdef") != "" { + t.Errorf("key must be lowercase hex sha256 (64 chars), got %q", a) + } +} + +func TestDifferentConversationsDifferentKeys(t *testing.T) { + a, b, c := fingerprint.Of([]byte(conv1)), fingerprint.Of([]byte(conv2)), fingerprint.Of([]byte(conv3)) + if a == b { + t.Errorf("different first user message must change the key") + } + if a == c { + t.Errorf("different system prompt must change the key") + } +} + +func TestNoUserMessageIsEmpty(t *testing.T) { + for _, body := range []string{ + `{"model":"m","messages":[{"role":"system","content":"only a system prompt"}]}`, + `{"model":"m","messages":[]}`, + `{"model":"m"}`, + `{"input":"an embeddings request"}`, + `not json at all`, + ``, + } { + if got := fingerprint.Of([]byte(body)); got != "" { + t.Errorf("Of(%q) = %q, want empty", body, got) + } + } +} + +func TestOnlyTheFirstFourKiBCount(t *testing.T) { + long := strings.Repeat("x", 5000) + a := `{"messages":[{"role":"user","content":"` + long + `A"}]}` + b := `{"messages":[{"role":"user","content":"` + long + `B"}]}` + if fingerprint.Of([]byte(a)) != fingerprint.Of([]byte(b)) { + t.Errorf("bytes after the first 4 KiB of a message must not change the key") + } + c := `{"messages":[{"role":"user","content":"A` + long + `"}]}` + if fingerprint.Of([]byte(a)) == fingerprint.Of([]byte(c)) { + t.Errorf("bytes inside the first 4 KiB must change the key") + } +} + +func TestContentPartsAreFlattened(t *testing.T) { + plain := `{"messages":[{"role":"user","content":"hello world"}]}` + parts := `{"messages":[{"role":"user","content":[{"type":"text","text":"hello world"}]}]}` + if fingerprint.Of([]byte(plain)) != fingerprint.Of([]byte(parts)) { + t.Errorf("a content array of text parts must fingerprint like the joined text") + } +} diff --git a/docs/plans/v1/_files/internal/lease/lease_test.go b/docs/plans/v1/_files/internal/lease/lease_test.go new file mode 100644 index 0000000..a1db60e --- /dev/null +++ b/docs/plans/v1/_files/internal/lease/lease_test.go @@ -0,0 +1,291 @@ +package lease_test + +import ( + "errors" + "sync" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +// memPersister is an in-memory Persister that also counts writes. +type memPersister struct { + mu sync.Mutex + leases map[[3]string]store.Lease + events []store.LeaseEvent + saves int +} + +func newPersister() *memPersister { return &memPersister{leases: map[[3]string]store.Lease{}} } + +func (m *memPersister) SaveLease(l store.Lease) error { + m.mu.Lock() + defer m.mu.Unlock() + m.saves++ + m.leases[[3]string{l.Route, l.FP, l.Model}] = l + return nil +} +func (m *memPersister) DeleteLease(route, fp, model string) error { + m.mu.Lock() + defer m.mu.Unlock() + delete(m.leases, [3]string{route, fp, model}) + return nil +} +func (m *memPersister) ListLeases() ([]store.Lease, error) { + m.mu.Lock() + defer m.mu.Unlock() + out := []store.Lease{} + for _, l := range m.leases { + out = append(out, l) + } + return out, nil +} +func (m *memPersister) RecordEvent(e store.LeaseEvent) error { + m.mu.Lock() + defer m.mu.Unlock() + m.events = append(m.events, e) + return nil +} +func (m *memPersister) reasons() []string { + m.mu.Lock() + defer m.mu.Unlock() + var r []string + for _, e := range m.events { + r = append(r, e.Reason) + } + return r +} + +// world is a hand-set view of hosts plus a chooser that returns a fixed answer. +type world struct { + mu sync.Mutex + healthy map[string]bool + draining map[string]bool + pick string + picks []string // candidates seen by Choose, for assertions +} + +func (w *world) Healthy(name string) bool { w.mu.Lock(); defer w.mu.Unlock(); return w.healthy[name] } +func (w *world) Draining(name string) bool { w.mu.Lock(); defer w.mu.Unlock(); return w.draining[name] } +func (w *world) Choose(candidates []string, model string) (string, bool) { + w.mu.Lock() + defer w.mu.Unlock() + w.picks = append([]string{}, candidates...) + for _, c := range candidates { + if c == w.pick { + return c, true + } + } + if len(candidates) > 0 { + return candidates[0], true + } + return "", false +} + +var t0 = time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + +func newTable(t *testing.T, p *memPersister, w *world) *lease.Table { + tbl, err := lease.New(p, w, w, 30*time.Minute) + if err != nil { + t.Fatal(err) + } + return tbl +} + +func TestNewLeaseThenSticky(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "beta"} + tbl := newTable(t, p, w) + k := lease.Key{Route: "r", FP: "conv1", Model: "m"} + host, reused, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0) + if err != nil || host != "beta" || reused { + t.Fatalf("first: %q %v %v", host, reused, err) + } + w.pick = "alpha" // the chooser would now prefer alpha; the lease must hold + for i := 1; i <= 5; i++ { + host, reused, err = tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Duration(i)*time.Minute)) + if err != nil || host != "beta" || !reused { + t.Fatalf("turn %d: %q reused=%v %v, want beta reused", i, host, reused, err) + } + } + if got := p.reasons(); len(got) != 1 || got[0] != store.ReasonNew { + t.Errorf("events = %v, want one 'new'", got) + } + snap := tbl.Snapshot() + if len(snap) != 1 || snap[0].Host != "beta" || !snap[0].LastUsed.Equal(t0.Add(5*time.Minute)) { + t.Errorf("snapshot = %+v", snap) + } + if p.saves < 2 { + t.Errorf("LastUsed must be written through (saves=%d)", p.saves) + } +} + +func TestUnhealthyHostMovesTheLease(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"} + tbl := newTable(t, p, w) + k := lease.Key{Route: "r", FP: "c", Model: "m"} + if host, _, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0); host != "alpha" { + t.Fatalf("first: %q", host) + } + w.mu.Lock() + w.healthy["alpha"] = false + w.pick = "beta" + w.mu.Unlock() + host, reused, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Minute)) + if err != nil || host != "beta" || reused { + t.Fatalf("after alpha down: %q reused=%v %v", host, reused, err) + } + if got := p.reasons(); len(got) != 2 || got[1] != store.ReasonUnhealthy { + t.Errorf("events = %v, want [new unhealthy]", got) + } + if len(w.picks) != 1 || w.picks[0] != "beta" { + t.Errorf("Choose must not see the unhealthy host: %v", w.picks) + } +} + +func TestIdleExpiry(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"} + tbl := newTable(t, p, w) + k := lease.Key{Route: "r", FP: "c", Model: "m"} + tbl.Acquire(k, []string{"alpha", "beta"}, t0) + if n := tbl.ExpireIdle(t0.Add(29 * time.Minute)); n != 0 { + t.Errorf("expired %d before lease_idle", n) + } + if n := tbl.ExpireIdle(t0.Add(31 * time.Minute)); n != 1 { + t.Errorf("expired %d after lease_idle, want 1", n) + } + if got := p.reasons(); got[len(got)-1] != store.ReasonIdle { + t.Errorf("events = %v, want idle last", got) + } + w.pick = "beta" + if host, reused, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(32*time.Minute)); host != "beta" || reused { + t.Errorf("after expiry a new lease is chosen: %q reused=%v", host, reused) + } + if l, _ := p.ListLeases(); len(l) != 1 { + t.Errorf("persister holds %d leases, want 1", len(l)) + } +} + +func TestFingerprintInheritsRouteLease(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "beta"} + tbl := newTable(t, p, w) + // A request without a fingerprint (no user message) leases the route itself… + if host, _, _ := tbl.Acquire(lease.Key{Route: "r", FP: "", Model: "m"}, []string{"alpha", "beta"}, t0); host != "beta" { + t.Fatalf("route lease: %q", host) + } + w.pick = "alpha" + // …and a new conversation on that route starts where the route already is. + host, reused, err := tbl.Acquire(lease.Key{Route: "r", FP: "conv", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second)) + if err != nil || host != "beta" || !reused { + t.Errorf("fingerprint lease must inherit the route's host: %q reused=%v %v", host, reused, err) + } + if len(tbl.Snapshot()) != 2 { + t.Errorf("both the route lease and the conversation lease exist: %+v", tbl.Snapshot()) + } +} + +func TestPinAndUnpin(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"} + tbl := newTable(t, p, w) + k := lease.Key{Route: "r", FP: "c", Model: "m"} + tbl.Acquire(k, []string{"alpha", "beta"}, t0) + if err := tbl.Pin("r", "beta", t0.Add(time.Minute)); err != nil { + t.Fatal(err) + } + host, _, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(2*time.Minute)) + if err != nil || host != "beta" { + t.Fatalf("pinned route must go to beta: %q %v", host, err) + } + host, _, err = tbl.Acquire(lease.Key{Route: "r", FP: "other", Model: "m"}, []string{"alpha", "beta"}, t0.Add(2*time.Minute)) + if err != nil || host != "beta" { + t.Fatalf("new conversations on a pinned route go to the pin too: %q %v", host, err) + } + w.mu.Lock() + w.healthy["beta"] = false + w.mu.Unlock() + if _, _, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(3*time.Minute)); !errors.Is(err, lease.ErrPinnedDown) { + t.Errorf("a pinned host that is down is ErrPinnedDown, never a silent move: %v", err) + } + if err := tbl.Pin("r", "nobody", t0); !errors.Is(err, lease.ErrUnknownHost) { + t.Errorf("pinning to a host not in the candidates of any lease: %v, want ErrUnknownHost", err) + } + tbl.Unpin("r") + w.mu.Lock() + w.healthy["beta"] = true + w.mu.Unlock() + if host, _, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(4*time.Minute)); host != "beta" { + t.Errorf("after unpin the existing lease (on beta) simply continues: %q", host) + } + if got := p.reasons(); got[len(got)-3] != store.ReasonPin || got[len(got)-1] != store.ReasonRelease { + t.Errorf("events = %v, want a pin event and a release event", got) + } +} + +func TestDrainKeepsExistingRefusesNew(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, draining: map[string]bool{}, pick: "alpha"} + tbl := newTable(t, p, w) + k := lease.Key{Route: "r", FP: "c", Model: "m"} + tbl.Acquire(k, []string{"alpha", "beta"}, t0) + w.mu.Lock() + w.draining["alpha"] = true + w.mu.Unlock() + if host, reused, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Minute)); host != "alpha" || !reused { + t.Errorf("an existing lease on a draining host continues: %q reused=%v", host, reused) + } + host, _, err := tbl.Acquire(lease.Key{Route: "r2", FP: "x", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Minute)) + if err != nil || host != "beta" { + t.Errorf("a new lease avoids the draining host: %q %v", host, err) + } + if len(w.picks) != 1 || w.picks[0] != "beta" { + t.Errorf("Choose must not see the draining host: %v", w.picks) + } + if _, _, err := tbl.Acquire(lease.Key{Route: "r3", FP: "y", Model: "m"}, []string{"alpha"}, t0); !errors.Is(err, lease.ErrNoHost) { + t.Errorf("only draining candidates: %v, want ErrNoHost", err) + } +} + +func TestReleaseRoute(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"} + tbl := newTable(t, p, w) + tbl.Acquire(lease.Key{Route: "r", FP: "a", Model: "m"}, []string{"alpha", "beta"}, t0) + tbl.Acquire(lease.Key{Route: "r", FP: "b", Model: "m"}, []string{"alpha", "beta"}, t0) + tbl.Acquire(lease.Key{Route: "other", FP: "c", Model: "m"}, []string{"alpha", "beta"}, t0) + if n := tbl.Release("r"); n != 2 { + t.Errorf("Release removed %d, want 2", n) + } + if n := tbl.Release("r"); n != 0 { + t.Errorf("second Release removed %d", n) + } + if l, _ := p.ListLeases(); len(l) != 1 || l[0].Route != "other" { + t.Errorf("persister after release: %+v", l) + } + w.pick = "beta" + if host, reused, _ := tbl.Acquire(lease.Key{Route: "r", FP: "a", Model: "m"}, []string{"alpha", "beta"}, t0); host != "beta" || reused { + t.Errorf("after release the route is re-chosen: %q reused=%v", host, reused) + } +} + +func TestLoadsFromPersister(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"} + _ = p.SaveLease(store.Lease{Route: "r", FP: "c", Model: "m", Host: "beta", State: store.Active, Created: t0, LastUsed: t0}) + _ = p.SaveLease(store.Lease{Route: "pinned", FP: "", Model: "", Host: "beta", State: store.Pinned, Created: t0, LastUsed: t0}) + tbl := newTable(t, p, w) + if host, reused, _ := tbl.Acquire(lease.Key{Route: "r", FP: "c", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second)); host != "beta" || !reused { + t.Errorf("a restart must not reshuffle: %q reused=%v", host, reused) + } + if host, _, _ := tbl.Acquire(lease.Key{Route: "pinned", FP: "new", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second)); host != "beta" { + t.Errorf("a pin survives a restart: %q", host) + } +} + +func TestNoCandidates(t *testing.T) { + p, w := newPersister(), &world{healthy: map[string]bool{}, pick: ""} + tbl := newTable(t, p, w) + if _, _, err := tbl.Acquire(lease.Key{Route: "r", FP: "c", Model: "m"}, []string{"alpha"}, t0); !errors.Is(err, lease.ErrNoHost) { + t.Errorf("no healthy host: %v, want ErrNoHost", err) + } + if len(tbl.Snapshot()) != 0 { + t.Errorf("a failed acquire must not create a lease") + } +} diff --git a/docs/plans/v1/_files/internal/limiter/limiter_test.go b/docs/plans/v1/_files/internal/limiter/limiter_test.go new file mode 100644 index 0000000..bf90588 --- /dev/null +++ b/docs/plans/v1/_files/internal/limiter/limiter_test.go @@ -0,0 +1,178 @@ +package limiter_test + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/limiter" +) + +func TestParallelAndQueue(t *testing.T) { + l := limiter.New() + l.Configure("alpha", "m", 2, 1) // two slots, one waiting place + ctx := context.Background() + + rel1, w1, err := l.Acquire(ctx, "alpha", "m") + if err != nil || w1 > 50*time.Millisecond { + t.Fatalf("first acquire: err %v waited %v", err, w1) + } + rel2, _, err := l.Acquire(ctx, "alpha", "m") + if err != nil { + t.Fatalf("second acquire: %v", err) + } + if l.InFlight("alpha", "m") != 2 || l.FreeSlots("alpha") != 0 { + t.Errorf("in flight %d free %d, want 2 and 0", l.InFlight("alpha", "m"), l.FreeSlots("alpha")) + } + + // Third waits in the queue. + got3 := make(chan error, 1) + 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") + } + } + got3 <- err + }() + time.Sleep(20 * time.Millisecond) + if l.Queued("alpha", "m") != 1 { + t.Errorf("queued = %d, want 1", l.Queued("alpha", "m")) + } + // Fourth finds the queue full and is refused at once. + start := time.Now() + _, _, err = l.Acquire(ctx, "alpha", "m") + if !errors.Is(err, limiter.ErrQueueFull) { + t.Fatalf("fourth acquire: %v, want ErrQueueFull", err) + } + if time.Since(start) > 50*time.Millisecond { + t.Errorf("a full queue must refuse immediately, took %v", time.Since(start)) + } + time.Sleep(30 * time.Millisecond) + rel1() // frees a slot: the queued third proceeds + select { + case err := <-got3: + if err != nil { + t.Fatalf("third: %v", err) + } + case <-time.After(time.Second): + t.Fatal("queued acquire did not proceed after a release") + } + rel2() + if l.InFlight("alpha", "m") != 0 || l.Queued("alpha", "m") != 0 { + t.Errorf("after releases: inflight %d queued %d", l.InFlight("alpha", "m"), l.Queued("alpha", "m")) + } +} + +func TestReleaseIsIdempotent(t *testing.T) { + l := limiter.New() + l.Configure("h", "m", 1, 0) + rel, _, err := l.Acquire(context.Background(), "h", "m") + if err != nil { + t.Fatal(err) + } + rel() + rel() // a second call must not free a slot that was never taken + if l.InFlight("h", "m") != 0 { + t.Errorf("in flight %d after double release", l.InFlight("h", "m")) + } + if _, _, err := l.Acquire(context.Background(), "h", "m"); err != nil { + t.Errorf("slot must be free again: %v", err) + } +} + +func TestCancelWhileQueuedLeaksNothing(t *testing.T) { + l := limiter.New() + l.Configure("h", "m", 1, 2) + rel, _, err := l.Acquire(context.Background(), "h", "m") + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }() + time.Sleep(20 * time.Millisecond) + cancel() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) { + t.Fatalf("cancelled acquire returned %v", err) + } + case <-time.After(time.Second): + t.Fatal("cancelled acquire did not return") + } + if l.Queued("h", "m") != 0 { + t.Errorf("queued = %d after cancel", l.Queued("h", "m")) + } + rel() + if l.InFlight("h", "m") != 0 { + t.Errorf("in flight %d, the cancelled waiter must not have taken the slot", l.InFlight("h", "m")) + } +} + +func TestQueueIsFIFO(t *testing.T) { + l := limiter.New() + l.Configure("h", "m", 1, 8) + rel, _, err := l.Acquire(context.Background(), "h", "m") + if err != nil { + t.Fatal(err) + } + var mu sync.Mutex + var order []int + var wg sync.WaitGroup + for i := 1; i <= 4; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + r, _, err := l.Acquire(context.Background(), "h", "m") + if err != nil { + t.Errorf("waiter %d: %v", i, err) + return + } + mu.Lock() + order = append(order, i) + mu.Unlock() + time.Sleep(5 * time.Millisecond) + r() + }(i) + time.Sleep(15 * time.Millisecond) // stagger arrivals so the order is defined + } + rel() + wg.Wait() + if len(order) != 4 || order[0] != 1 || order[1] != 2 || order[2] != 3 || order[3] != 4 { + t.Errorf("waiters proceeded in order %v, want [1 2 3 4]", order) + } +} + +func TestUnconfiguredPairIsOneSlotNoQueue(t *testing.T) { + l := limiter.New() + rel, _, err := l.Acquire(context.Background(), "x", "y") + if err != nil { + t.Fatal(err) + } + defer rel() + if _, _, err := l.Acquire(context.Background(), "x", "y"); !errors.Is(err, limiter.ErrQueueFull) { + t.Errorf("second acquire on an unconfigured pair: %v, want ErrQueueFull", err) + } +} + +func TestFreeSlotsSumsModels(t *testing.T) { + l := limiter.New() + l.Configure("h", "a", 4, 0) + l.Configure("h", "b", 2, 0) + if got := l.FreeSlots("h"); got != 6 { + t.Fatalf("free = %d, want 6", got) + } + rel, _, _ := l.Acquire(context.Background(), "h", "a") + defer rel() + if got := l.FreeSlots("h"); got != 5 { + t.Errorf("free = %d, want 5", got) + } + if l.FreeSlots("nobody") != 0 { + t.Errorf("unknown host has no slots") + } +} diff --git a/docs/plans/v1/_files/internal/proxy/proxy_test.go b/docs/plans/v1/_files/internal/proxy/proxy_test.go new file mode 100644 index 0000000..cf33fcc --- /dev/null +++ b/docs/plans/v1/_files/internal/proxy/proxy_test.go @@ -0,0 +1,398 @@ +package proxy_test + +// v1 acceptance tests for the proxy: leases, queueing, accounting, header route override. +// They drive the whole handler over real HTTP against fake upstreams; only what a client or an +// operator can observe is asserted (status codes, headers, the accounting rows, the health table). + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/lease" + "git.wntrmute.dev/kyle/crossbar/internal/limiter" + "git.wntrmute.dev/kyle/crossbar/internal/proxy" + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +// upstream is a llama-server stand-in: streams N chunks with a delay, reports usage/timings in +// the final chunk, counts requests, and can be slowed down or killed. +type upstream struct { + name string + srv *httptest.Server + hits atomic.Int32 + delay time.Duration + mu sync.Mutex + last recorded +} + +type recorded struct{ method, path, host, xff, body string } + +func newUpstream(t *testing.T, name string) *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, `{"object":"list","data":[{"id":"shared"},{"id":"`+name+`-only"}]}`) + }) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + u.hits.Add(1) + b, _ := io.ReadAll(r.Body) + u.mu.Lock() + u.last = recorded{r.Method, r.URL.RequestURI(), r.Host, r.Header.Get("X-Forwarded-For"), string(b)} + u.mu.Unlock() + var req struct { + Stream bool `json:"stream"` + } + _ = json.Unmarshal(b, &req) + w.Header().Set("X-Upstream", name) + time.Sleep(u.delay) + if !req.Stream { + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"choices":[{"message":{"role":"assistant","content":"hi from %s"}}],"usage":{"prompt_tokens":100,"completion_tokens":10,"total_tokens":110},"timings":{"prompt_n":100,"cache_n":90,"predicted_n":10,"predicted_ms":50.0}}`, name) + return + } + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(200) + fl := w.(http.Flusher) + for i := 0; i < 3; i++ { + fmt.Fprintf(w, "data: {\"choices\":[{\"delta\":{\"content\":\"%s %d \"}}]}\n\n", name, i) + fl.Flush() + time.Sleep(10 * time.Millisecond) + } + fmt.Fprint(w, `data: {"choices":[],"usage":{"prompt_tokens":200,"completion_tokens":20,"total_tokens":220},"timings":{"prompt_n":200,"cache_n":150,"predicted_n":20,"predicted_ms":80.0}}`+"\n\n") + fl.Flush() + fmt.Fprint(w, "data: [DONE]\n\n") + }) + u.srv = httptest.NewServer(mux) + t.Cleanup(u.srv.Close) + return u +} + +func (u *upstream) lastReq() recorded { u.mu.Lock(); defer u.mu.Unlock(); return u.last } + +// rig is one crossbar: config, real health table (polled once), real lease table over a real +// SQLite store, real limiter, the proxy handler served by httptest. +type rig struct { + t *testing.T + cfg *config.Config + health *health.Table + store *store.Store + leases *lease.Table + lim *limiter.Limiter + front *httptest.Server +} + +// newRig builds crossbar from a config text where %s placeholders are the upstream base URLs. +func newRig(t *testing.T, cfgText string, ups ...*upstream) *rig { + urls := make([]any, len(ups)) + for i, u := range ups { + urls[i] = u.srv.URL + } + cfg, err := config.Parse(strings.NewReader(fmt.Sprintf(cfgText, urls...))) + if err != nil { + t.Fatal(err) + } + bases := map[string]string{} + for name, h := range cfg.Hosts { + bases[name] = h.BaseURL + } + ht := health.New(bases, time.Hour, nil) + ht.PollOnce(t.Context()) + st, err := store.Open(filepath.Join(t.TempDir(), "crossbar.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = st.Close() }) + lim := limiter.New() + for name, h := range cfg.Hosts { + for model, m := range h.Models { + lim.Configure(name, model, m.Parallel, cfg.QueueMax) + } + } + lt, err := lease.New(st, proxy.HostView(ht, cfg), proxy.Chooser(cfg, ht, lim), cfg.LeaseIdle.Duration) + if err != nil { + t.Fatal(err) + } + p := proxy.New(cfg, ht, lt, lim, st, nil) + front := httptest.NewServer(p) + t.Cleanup(front.Close) + return &rig{t: t, cfg: cfg, health: ht, store: st, leases: lt, lim: lim, front: front} +} + +const twoHosts = ` +listen = "127.0.0.1:1" +queue_max = 1 +lease_idle = "30m" +[hosts.alpha] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 2 }, "alpha-only" = { } } +[hosts.beta] +base_url = %q +weight = 2.0 +models = { "shared" = { parallel = 2 }, "beta-only" = { } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "shared" +[routes.other] +hosts = ["alpha"] +` + +func conversation(id, turn int) string { + msgs := fmt.Sprintf(`{"role":"system","content":"project"},{"role":"user","content":"conversation %d opening"}`, id) + for i := 1; i < turn; i++ { + msgs += fmt.Sprintf(`,{"role":"assistant","content":"ok"},{"role":"user","content":"turn %d"}`, i) + } + return `{"model":"shared","stream":false,"messages":[` + msgs + `]}` +} + +func (r *rig) post(path, body string, hdr ...string) *http.Response { + req, _ := http.NewRequest(http.MethodPost, r.front.URL+path, strings.NewReader(body)) + req.Header.Set("Content-Type", "application/json") + for i := 0; i+1 < len(hdr); i += 2 { + req.Header.Set(hdr[i], hdr[i+1]) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + r.t.Fatal(err) + } + return resp +} + +func drain(resp *http.Response) string { + b, _ := io.ReadAll(resp.Body) + resp.Body.Close() + return string(b) +} + +func TestConversationIsStickyAndLeaseHeaderTellsWhy(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + first := r.post("/r/v1/chat/completions", conversation(1, 1)) + drain(first) + host := first.Header.Get(proxy.HostHeader) + if first.StatusCode != 200 || host != "beta" { // beta: same free slots, double weight + t.Fatalf("first turn: %d from %q, want 200 from beta", first.StatusCode, host) + } + if got := first.Header.Get(proxy.LeaseHeader); got != "new" { + t.Errorf("%s = %q on the first turn, want new", proxy.LeaseHeader, got) + } + // Take alpha's slots away as a "better host" signal: it must not matter, the lease holds. + for turn := 2; turn <= 6; turn++ { + resp := r.post("/r/v1/chat/completions", conversation(1, turn)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != host || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Fatalf("turn %d: host %q lease %q, want %q reused", turn, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader), host) + } + } + if alpha.hits.Load() != 0 || beta.hits.Load() != 6 { + t.Errorf("hits alpha=%d beta=%d, want 0 and 6", alpha.hits.Load(), beta.hits.Load()) + } +} + +func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + beta.delay = 300 * time.Millisecond + r := newRig(t, twoHosts, alpha, beta) + // Two slow conversations occupy beta's two "shared" slots… + var wg sync.WaitGroup + for i := 1; i <= 2; i++ { + wg.Add(1) + go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i) + } + time.Sleep(100 * time.Millisecond) + // …so a third conversation starting now is sent to alpha (beta has 0 free slots, alpha 2). + resp := r.post("/r/v1/chat/completions", conversation(3, 1)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "alpha" { + t.Errorf("third conversation went to %q, want alpha (free slots beat weight)", resp.Header.Get(proxy.HostHeader)) + } + wg.Wait() +} + +func TestQueueFullIs503(t *testing.T) { + alpha := newUpstream(t, "alpha") + alpha.delay = 400 * time.Millisecond + r := newRig(t, ` +listen = "127.0.0.1:1" +queue_max = 1 +[hosts.alpha] +base_url = %q +models = { "shared" = { parallel = 1 } } +[routes.r] +hosts = ["alpha"] +default_model = "shared" +`, alpha) + codes := make(chan int, 3) + for i := 1; i <= 3; i++ { + go func(i int) { + 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 + } + got := map[int]int{} + for i := 0; i < 3; i++ { + got[<-codes]++ + } + if got[200] != 2 || got[503] != 1 { + t.Fatalf("status counts = %v, want two 200 and one 503", got) + } + rows, err := r.store.Usage(time.Time{}, store.ByRoute) + if err != nil || len(rows) != 1 || rows[0].Requests != 3 || rows[0].Errors != 1 { + t.Errorf("usage = %+v %v, want 3 requests, 1 error (the 503 is recorded too)", rows, err) + } + if rows[0].QueuedMs <= 0 { + t.Errorf("the queued request must record its wait: %+v", rows[0]) + } +} + +func TestUnhealthyHostReleasesAndMoves(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // lands on beta + beta.srv.Close() + resp := r.post("/r/v1/chat/completions", conversation(1, 2)) + drain(resp) + if resp.StatusCode != http.StatusBadGateway { + t.Fatalf("first request after beta died: %d, want 502", resp.StatusCode) + } + if s, _ := r.health.Get("beta"); s.Healthy { + t.Fatalf("beta must be marked down after the 502") + } + resp = r.post("/r/v1/chat/completions", conversation(1, 3)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" || resp.Header.Get(proxy.LeaseHeader) != "new" { + t.Errorf("after the move: %d from %q lease %q, want 200 alpha new", resp.StatusCode, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } + ev, _ := r.store.Events(time.Time{}, 10) + var reasons []string + for _, e := range ev { + reasons = append(reasons, e.Reason) + } + if len(reasons) != 2 || reasons[0] != store.ReasonNew || reasons[1] != store.ReasonUnhealthy { + t.Errorf("lease events = %v, want [new unhealthy]", reasons) + } +} + +func TestAccountingRowsFromUsageAndTimings(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // non-streamed + drain(r.post("/r/v1/chat/completions", strings.Replace(conversation(1, 2), `"stream":false`, `"stream":true`, 1))) // streamed + deadline := time.Now().Add(2 * time.Second) + var rows []store.UsageRow + for time.Now().Before(deadline) { + rows, _ = r.store.Usage(time.Time{}, store.ByHost) + if len(rows) == 1 && rows[0].Requests == 2 { + break + } + time.Sleep(20 * time.Millisecond) + } + if len(rows) != 1 || rows[0].Requests != 2 { + t.Fatalf("usage by host = %+v, want one host with 2 requests (rows may be written after the response completes, within 2 s)", rows) + } + u := rows[0] + if u.PromptTokens != 300 || u.CachedTokens != 240 || u.CompletionTokens != 30 { + t.Errorf("tokens = prompt %d cached %d completion %d, want 300/240/30 (100+200, 90+150, 10+20)", u.PromptTokens, u.CachedTokens, u.CompletionTokens) + } + if u.BusyMs <= 0 || u.Errors != 0 { + t.Errorf("busy %d errors %d", u.BusyMs, u.Errors) + } + if got := u.CacheHitRatio(); got < 0.79 || got > 0.81 { + t.Errorf("cache hit ratio = %v, want 0.8", got) + } +} + +func TestStreamIsUnalteredWhileTeed(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.post("/r/v1/chat/completions", strings.Replace(conversation(9, 1), `"stream":false`, `"stream":true`, 1)) + body := drain(resp) + want := 0 + for _, line := range strings.Split(body, "\n") { + if strings.HasPrefix(line, "data: ") { + want++ + } + } + if want != 5 || !strings.HasSuffix(strings.TrimSpace(body), "data: [DONE]") { + t.Errorf("client must receive every SSE line untouched (3 deltas, usage, DONE); got %d data lines:\n%s", want, body) + } +} + +func TestHeaderRouteOverride(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + // The header names the route; the path has none. + resp := r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other") + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" { + t.Errorf("header route 'other' (alpha only): %d from %q", resp.StatusCode, resp.Header.Get(proxy.HostHeader)) + } + if alpha.lastReq().path != "/v1/chat/completions" { + t.Errorf("upstream path = %q", alpha.lastReq().path) + } + // A path route and a header route that disagree: the header is the operator's intent → 400. + resp = r.post("/r/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other") + if drain(resp); resp.StatusCode != 400 { + t.Errorf("conflicting route in path and header: %d, want 400", resp.StatusCode) + } + resp = r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "nope") + if drain(resp); resp.StatusCode != 404 { + t.Errorf("unknown header route: %d, want 404", resp.StatusCode) + } +} + +func TestV0BehaviourStillHolds(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + for _, tc := range []struct { + method, path string + want int + msg string + }{ + {http.MethodGet, "/", 400, "missing route"}, + {http.MethodGet, "/nope/v1/models", 404, "unknown route"}, + {http.MethodGet, "/r/slots", 404, "not found"}, + {http.MethodGet, "/r/_crossbar/hosts", 404, "not found"}, + } { + req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + body := drain(resp) + var e map[string]string + if resp.StatusCode != tc.want || json.Unmarshal([]byte(body), &e) != nil || e["error"] != tc.msg { + t.Errorf("%s: %d %s, want %d %q", tc.path, resp.StatusCode, body, tc.want, tc.msg) + } + } + big := strings.Repeat("x", proxy.MaxBody+1) + resp := r.post("/r/v1/chat/completions", big) + if drain(resp); resp.StatusCode != 413 { + t.Errorf("oversize body: %d, want 413", resp.StatusCode) + } + // GET pass-through with query string, Host and X-Forwarded-For as in v0. + resp, err := http.Get(r.front.URL + "/r/v1/models?x=1") + if err != nil { + t.Fatal(err) + } + drain(resp) + host := resp.Header.Get(proxy.HostHeader) + u := map[string]*upstream{"alpha": alpha, "beta": beta}[host] + if u == nil || u.lastReq().path != "/v1/models?x=1" || u.lastReq().host != strings.TrimPrefix(u.srv.URL, "http://") || u.lastReq().xff == "" { + t.Errorf("GET pass-through: host %q last %+v", host, u.lastReq()) + } +} diff --git a/docs/plans/v1/_files/internal/store/store_test.go b/docs/plans/v1/_files/internal/store/store_test.go new file mode 100644 index 0000000..e7f2c28 --- /dev/null +++ b/docs/plans/v1/_files/internal/store/store_test.go @@ -0,0 +1,185 @@ +package store_test + +import ( + "path/filepath" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/store" +) + +func open(t *testing.T, dir string) *store.Store { + s, err := store.Open(filepath.Join(dir, "crossbar.db")) + if err != nil { + t.Fatalf("Open: %v", err) + } + t.Cleanup(func() { _ = s.Close() }) + return s +} + +func TestOpenIsIdempotentAndWAL(t *testing.T) { + dir := t.TempDir() + s := open(t, dir) + if got := s.JournalMode(); got != "wal" { + t.Errorf("journal_mode = %q, want wal", got) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + open(t, dir) // second open on the same file must not fail on existing tables +} + +func TestLeasesSurviveReopen(t *testing.T) { + dir := t.TempDir() + s := open(t, dir) + now := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + l := store.Lease{Route: "opencode-a", FP: "abc", Model: "m", Host: "alpha", State: store.Active, Created: now, LastUsed: now} + if err := s.SaveLease(l); err != nil { + t.Fatal(err) + } + l2 := l + l2.FP = "def" + l2.Host = "beta" + l2.State = store.Pinned + if err := s.SaveLease(l2); err != nil { + t.Fatal(err) + } + // Saving the same key again replaces, not duplicates. + l.Host = "beta" + l.LastUsed = now.Add(time.Minute) + if err := s.SaveLease(l); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + s = open(t, dir) + got, err := s.ListLeases() + if err != nil { + t.Fatal(err) + } + if len(got) != 2 { + t.Fatalf("ListLeases = %d rows, want 2: %+v", len(got), got) + } + byFP := map[string]store.Lease{} + for _, x := range got { + byFP[x.FP] = x + } + if a := byFP["abc"]; a.Host != "beta" || !a.LastUsed.Equal(now.Add(time.Minute)) || a.State != store.Active { + t.Errorf("abc = %+v", a) + } + if d := byFP["def"]; d.State != store.Pinned || d.Host != "beta" { + t.Errorf("def = %+v", d) + } + if err := s.DeleteLease("opencode-a", "abc", "m"); err != nil { + t.Fatal(err) + } + got, _ = s.ListLeases() + if len(got) != 1 || got[0].FP != "def" { + t.Errorf("after delete: %+v", got) + } +} + +func TestEventsAndRequestsAndUsage(t *testing.T) { + s := open(t, t.TempDir()) + t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + must := func(err error) { + if err != nil { + t.Fatal(err) + } + } + must(s.RecordEvent(store.LeaseEvent{TS: t0, Route: "r1", Model: "m", FromHost: "", ToHost: "alpha", Reason: store.ReasonNew})) + must(s.RecordEvent(store.LeaseEvent{TS: t0.Add(time.Hour), Route: "r1", Model: "m", FromHost: "alpha", ToHost: "beta", Reason: store.ReasonUnhealthy})) + reqs := []store.Request{ + {Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0, QueuedMs: 0, TTFBMs: 100, TotalMs: 1000, Status: 200, Streamed: true, PromptTokens: 1000, CachedTokens: 900, CompletionTokens: 50}, + {Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0.Add(time.Minute), QueuedMs: 40, TTFBMs: 120, TotalMs: 2000, Status: 200, Streamed: true, PromptTokens: 1100, CachedTokens: 1000, CompletionTokens: 60}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(2 * time.Minute), TotalMs: 500, Status: 502, Err: "upstream failed"}, + {Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(-48 * time.Hour), TotalMs: 300, Status: 200, PromptTokens: 10, CompletionTokens: 5}, + } + for _, r := range reqs { + must(s.RecordRequest(r)) + } + must(s.RecordHostHealth(store.HostHealth{TS: t0, Host: "alpha", Healthy: true, Loaded: []string{"m"}})) + + rows, err := s.Usage(t0.Add(-time.Hour), store.ByRoute) + must(err) + if len(rows) != 2 { + t.Fatalf("Usage by route since t0-1h: %d rows, want 2 (r1, r2): %+v", len(rows), rows) + } + byKey := map[string]store.UsageRow{} + for _, r := range rows { + byKey[r.Key] = r + } + r1 := byKey["r1"] + if r1.Requests != 2 || r1.Errors != 0 || r1.BusyMs != 3000 || r1.QueuedMs != 40 { + t.Errorf("r1 = %+v", r1) + } + if r1.PromptTokens != 2100 || r1.CachedTokens != 1900 || r1.CompletionTokens != 110 { + t.Errorf("r1 tokens = %+v", r1) + } + if got := r1.CacheHitRatio(); got < 0.904 || got > 0.905 { + t.Errorf("r1 cache hit ratio = %v, want 1900/2100", got) + } + r2 := byKey["r2"] + if r2.Requests != 1 || r2.Errors != 1 || r2.BusyMs != 500 { + t.Errorf("r2 = %+v (the 48h-old request is outside since)", r2) + } + if r2.CacheHitRatio() != 0 { + t.Errorf("no prompt tokens: ratio must be 0, got %v", r2.CacheHitRatio()) + } + byHost, err := s.Usage(time.Time{}, store.ByHost) + must(err) + if len(byHost) != 2 { + t.Errorf("by host, all time: %+v", byHost) + } + for _, r := range byHost { + if r.Key == "beta" && r.Requests != 2 { + t.Errorf("beta all-time requests = %d, want 2", r.Requests) + } + } + byModel, err := s.Usage(time.Time{}, store.ByModel) + must(err) + if len(byModel) != 1 || byModel[0].Key != "m" || byModel[0].Requests != 4 { + t.Errorf("by model: %+v", byModel) + } + ev, err := s.Events(t0.Add(-time.Minute), 10) + must(err) + if len(ev) != 2 || ev[0].Reason != store.ReasonNew || ev[1].ToHost != "beta" { + t.Errorf("events = %+v", ev) + } +} + +func TestPruneRollsUpOldRequests(t *testing.T) { + s := open(t, t.TempDir()) + t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC) + old := t0.Add(-200 * 24 * time.Hour) + for i := 0; i < 3; i++ { + if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: old.Add(time.Duration(i) * time.Minute), TotalMs: 100, Status: 200, PromptTokens: 10, CachedTokens: 5, CompletionTokens: 1}); err != nil { + t.Fatal(err) + } + } + if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: t0, TotalMs: 100, Status: 200}); err != nil { + t.Fatal(err) + } + n, err := s.Prune(t0, 180*24*time.Hour) + if err != nil { + t.Fatal(err) + } + if n != 3 { + t.Errorf("Prune removed %d rows, want 3", n) + } + rows, _ := s.Usage(time.Time{}, store.ByRoute) + if len(rows) != 1 || rows[0].Requests != 4 || rows[0].PromptTokens != 30 { + t.Errorf("usage must still include pruned traffic through the daily rollup: %+v", rows) + } + live, _ := s.Usage(old.Add(24*time.Hour), store.ByRoute) + if len(live) != 1 || live[0].Requests != 1 { + t.Errorf("recent-only usage = %+v", live) + } +} + +func TestBadPath(t *testing.T) { + if _, err := store.Open(filepath.Join(t.TempDir(), "no", "such", "dir", "x.db")); err == nil { + t.Fatal("Open must fail when the directory does not exist") + } +} diff --git a/docs/plans/v1/_files/tools/smoke.sh b/docs/plans/v1/_files/tools/smoke.sh new file mode 100755 index 0000000..7c2bcb4 --- /dev/null +++ b/docs/plans/v1/_files/tools/smoke.sh @@ -0,0 +1,80 @@ +#!/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. +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 $!" +bin/crossbar -config "$tmp/crossbar.toml" >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" +sleep 1.5 +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"; } +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) +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'" + +# 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" + +# 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)" + +# 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")" + +# 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 +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" +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 +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" +echo "smoke: ok (stream spread $((lastms - firstms)) ms)"