From 74b8af4ce0a693f2beea04b43e765e731717ba8f Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 08:37:38 -0700 Subject: [PATCH] AGENTS.md: a refused tool call is not a reason to end the turn; v2 plan (props, ctx guard, wake, identity, wiring) as acceptance tests; v1.1 run note v2 given tests compiled against a panic-only skeleton (go vet clean); no reference implementation. Co-Authored-By: Claude Fable 5.1 --- AGENTS.md | 5 + docs/plans/v1.1/README.md | 7 + docs/plans/v2/01-props.md | 64 +++++++ docs/plans/v2/02-ctxguard.md | 67 ++++++++ docs/plans/v2/03-wake.md | 65 +++++++ docs/plans/v2/04-identity.md | 92 ++++++++++ docs/plans/v2/05-wiring-smoke.md | 58 +++++++ docs/plans/v2/README.md | 54 ++++++ docs/plans/v2/_files/cmd/fakeupstream/main.go | 161 ++++++++++++++++++ docs/plans/v2/_files/example.toml | 32 ++++ .../_files/internal/config/config_v2_test.go | 75 ++++++++ .../v2/_files/internal/health/props_test.go | 75 ++++++++ .../_files/internal/identity/identity_test.go | 90 ++++++++++ .../internal/identity/middleware_test.go | 74 ++++++++ .../internal/identity/testdata/whois.json | 24 +++ .../v2/_files/internal/proxy/ctxguard_test.go | 142 +++++++++++++++ .../v2/_files/internal/wake/wake_test.go | 114 +++++++++++++ docs/plans/v2/_files/tools/smoke.sh | 58 +++++++ 18 files changed, 1257 insertions(+) create mode 100644 docs/plans/v2/01-props.md create mode 100644 docs/plans/v2/02-ctxguard.md create mode 100644 docs/plans/v2/03-wake.md create mode 100644 docs/plans/v2/04-identity.md create mode 100644 docs/plans/v2/05-wiring-smoke.md create mode 100644 docs/plans/v2/README.md create mode 100644 docs/plans/v2/_files/cmd/fakeupstream/main.go create mode 100644 docs/plans/v2/_files/example.toml create mode 100644 docs/plans/v2/_files/internal/config/config_v2_test.go create mode 100644 docs/plans/v2/_files/internal/health/props_test.go create mode 100644 docs/plans/v2/_files/internal/identity/identity_test.go create mode 100644 docs/plans/v2/_files/internal/identity/middleware_test.go create mode 100644 docs/plans/v2/_files/internal/identity/testdata/whois.json create mode 100644 docs/plans/v2/_files/internal/proxy/ctxguard_test.go create mode 100644 docs/plans/v2/_files/internal/wake/wake_test.go create mode 100755 docs/plans/v2/_files/tools/smoke.sh diff --git a/AGENTS.md b/AGENTS.md index 607d897..0c0ef18 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -54,6 +54,11 @@ These come from defects found in review; the evidence is in `docs/implementer-lo - When a rule says "every" or "everywhere", finish by listing each place it applies and checking them one by one. The task shows one place; the rule covers all of them. - Never end a turn by describing what you are about to do. Do it, then report. +- Work only inside this repository. Scratch programs under `/tmp` or anywhere else are refused + by the sandbox, and **a refused tool call is not a reason to end the turn**: write the + experiment as a `_test.go` file inside the repository (delete it before committing), or reason + it out. Two sessions have ended with a plan and no tool call right after a refusal; that + leaves the owner with no commit and no `stopped` row, the worst outcome. ## The gate diff --git a/docs/plans/v1.1/README.md b/docs/plans/v1.1/README.md index 8ceb439..f1aafe2 100644 --- a/docs/plans/v1.1/README.md +++ b/docs/plans/v1.1/README.md @@ -25,3 +25,10 @@ Branch `v1.1`. One task, one fresh OpenCode session, one commit. 2. `cmp` both copied tests; `git diff master..v1.1 --stat -- PLAN.md AGENTS.md docs/plans` empty. 3. `make gate`, `make smoke`. 4. Probe: cut a stream with `curl -m 0.4 -N …` against the smoke rig and confirm one `status="499"` line in `/_crossbar/metrics`. + +## Changes during the run + +- 2026-09-25, task 01, first session: ended after ~8 min with no commit and no row, right after + the sandbox refused a `/tmp` scratch program (the I9 pattern, third time tonight). The rule + against ending a turn on a refusal lived only in v1's task 05; it is now in `AGENTS.md`, so every + task carries it. Resumed from the working tree. diff --git a/docs/plans/v2/01-props.md b/docs/plans/v2/01-props.md new file mode 100644 index 0000000..72a7f93 --- /dev/null +++ b/docs/plans/v2/01-props.md @@ -0,0 +1,64 @@ +# v2 task 01: learn each host's context size from `/props` + +**Branch:** `v2` (create it from `master`: `git switch master && git switch -c v2`; `git status --short` must be empty first, otherwise stop) +**Commit subject:** `Health: learn n_ctx and total_slots from /props` + +## Goal + +The poller already asks each host `/health` and `/v1/models`. It now also reads `/props` and +remembers the context size and slot count, so the proxy (task 02) can tell whether a prompt fits. +A missing or malformed `/props` is **not** a health failure: context is then simply unknown (0). + +## Context + +`llama-server` answers `GET /props` with a JSON object containing +`default_generation_settings.n_ctx` (the total context the server was started with) and +`total_slots` (how many parallel slots share it). With unified KV, one slot can use up to +`n_ctx / total_slots` tokens. Some builds omit fields; old ones 404. Read at most `MaxModelsBody` +bytes as for the other endpoints. + +## Files + +- Copy: `internal/health/props_test.go` +- Modify: `internal/health/health.go`, `internal/admin/admin.go` (or wherever `HostView` is built), `docs/implementer-log.md` + +## Interfaces + +`internal/health`, additions: + +```go +type Status struct { + // …existing fields… + NCtx int `json:"n_ctx"` // total context from /props; 0 = unknown + Slots int `json:"slots"` // total_slots from /props; 0 = unknown +} +// PerSlotCtx is the context one request may use: NCtx / Slots, or NCtx when Slots is 0. +func (s Status) PerSlotCtx() int +``` + +Rules the tests check: + +1. A poll is `/health`, `/v1/models` (as before), then `GET /props`. If that request fails, + returns non-200, is not JSON, or lacks the fields, set `NCtx = 0`, `Slots = 0` and **do not + count the poll as failed**. Otherwise `NCtx = default_generation_settings.n_ctx`, + `Slots = total_slots` (negative values → 0). +2. `PerSlotCtx()` is integer division; `NCtx` when `Slots == 0`; 0 when `NCtx == 0`. +3. `admin.HostView` gains `NCtx int \`json:"n_ctx"\`` and `Slots int \`json:"slots"\`` copied + from the status (the v2 smoke reads `"n_ctx":8192` from `/_crossbar/hosts`). `admin_test.go` + must keep passing unchanged. + +## Steps + +- [ ] **1.** `git switch master && git switch -c v2`; `cp docs/plans/v2/_files/internal/health/props_test.go internal/health/`. +- [ ] **2. See it fail** (compile: `NCtx` undefined). **3. Write the code.** `gofmt -w internal/`. +- [ ] **4.** `go test -race -count=1 ./internal/health/ ./internal/admin/` → both `ok`. +- [ ] **5.** `make gate` → `gate: ok`. **6.** Row `v2/01-props`; commit. + +```sh +git add internal/health internal/admin docs/implementer-log.md +git commit +``` + +## Done when + +- Both packages pass; gate ok; `props_test.go` byte-identical to `_files/`. diff --git a/docs/plans/v2/02-ctxguard.md b/docs/plans/v2/02-ctxguard.md new file mode 100644 index 0000000..270ab03 --- /dev/null +++ b/docs/plans/v2/02-ctxguard.md @@ -0,0 +1,67 @@ +# v2 task 02: the context-size guard + +**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Proxy: move or refuse prompts that do not fit the leased host's context` + +## Goal + +The single worst failure a client sees today is the upstream's "Context size has been exceeded" +after a long wait. With `PerSlotCtx` known (task 01), crossbar can estimate a prompt's size from +its body and act before forwarding: move the conversation to a host where it fits, or answer 400 +with the estimate and the largest slot available. `PLAN.md` §4b. + +## Files + +- Copy: `internal/proxy/ctxguard_test.go` +- Modify: `internal/proxy/proxy.go` and/or `forward.go` (new file `ctxguard.go` if that keeps files under 400 lines), `internal/lease/lease.go` (one addition, below), `docs/implementer-log.md` + +## Interfaces + +```go +// internal/proxy +const CtxHeader = "X-Crossbar-Ctx" // set only when the guard moved a conversation: "moved:>>" e.g. "moved:small>big" + +// internal/lease — one addition, the only change allowed there: +// Move re-leases k onto host (deleting any existing lease for k), records a LeaseEvent with +// Reason "ctx", FromHost the previous host ("" if none), ToHost host. Returns an error only from +// the persister. +func (t *Table) Move(k Key, host string, now time.Time) error +``` + +Rules the tests check (`ServeHTTP`, between acquiring the lease and taking a slot): + +1. `estimate := int(float64(len(body)) / 4 * 1.2)` tokens (body = the bytes already peeked; GET/HEAD → 0). +2. `limit := PerSlotCtx` of the leased host's `health.Status`. If `limit == 0` (unknown) or + `estimate <= limit`: no action, no header. +3. Otherwise find, among the route's candidate hosts (in order), the healthy, non-draining hosts + whose `PerSlotCtx() >= estimate` — prefer one that lists the model as loaded, else one that + can serve it (`cfg.Serves`). If one exists: `leases.Move(key, host, now)`, set + `CtxHeader` to `moved:>`, and continue with the new host (**this request and the + following turns**: the lease moved). +4. If none exists: **400** `{"error":"prompt too large","estimate":E,"max":M}` where `M` is the + largest `PerSlotCtx()` among the route's healthy hosts (0 if all unknown — but then rule 2 + already let the request through). Record an accounting row with status 400 and + `Err: "prompt too large"`; do not mark anything down; do not forward. +5. The estimate is never logged with the body; the log line gains `ctx_est=E` only. + +Add `ReasonCtx = "ctx"` to `internal/store` constants (one-line change, allowed). + +## Steps + +- [ ] **1.** `git switch v2`; `cp docs/plans/v2/_files/internal/proxy/ctxguard_test.go internal/proxy/`. +- [ ] **2. See it fail** (compile: `CtxHeader`). **3. Write the code.** `gofmt -w internal/`. +- [ ] **4.** `go test -race -count=2 ./internal/proxy/ ./internal/lease/` → `ok`. +- [ ] **5.** `make gate` → `gate: ok`. **6.** Row `v2/02-ctxguard`; commit. + +```sh +git add internal/proxy internal/lease internal/store docs/implementer-log.md +git commit +``` + +## Done when + +- Tests pass with `-race -count=2`; gate ok; the copied test is byte-identical. + +## Stop and report if + +- `TestStickyLeaseSurvivesGrowthUntilItDoesNotFit` fails on the *second* request after the move (the lease did not actually move): quote the lease table. diff --git a/docs/plans/v2/03-wake.md b/docs/plans/v2/03-wake.md new file mode 100644 index 0000000..2c7fa19 --- /dev/null +++ b/docs/plans/v2/03-wake.md @@ -0,0 +1,65 @@ +# v2 task 03: wake-on-LAN + +**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the wake package: magic packets and a waiter` + +## Goal + +A pure package. `wake.MagicPacket` builds the 102-byte wake-on-LAN frame, `wake.Send` puts it on +the wire as UDP, and `wake.Waker` wakes a named host at most once per wait window and waits for +the health table to report it healthy. Task 05 uses it when a route has no healthy host. + +## Context + +A magic packet is six `0xff` bytes followed by the target MAC sixteen times, sent as a UDP +datagram to the LAN broadcast address (port 9 by convention). The sleeping Mac (titan) has +wake-on-magic-packet enabled; it takes 20–40 s to be reachable. Waking twice inside that window +is harmless but pointless, so the waker remembers when it last sent. + +## Files + +- Copy: `internal/wake/wake_test.go` +- Create: `internal/wake/wake.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +```go +package wake + +type Target struct { + MAC, Broadcast string // "aa:bb:cc:dd:ee:ff" (also "-" separated, any case); "host:port" + Wait time.Duration +} +type Health interface{ Healthy(name string) bool } + +func MagicPacket(mac string) ([]byte, error) // net.ParseMAC; must be 6 bytes; 102-byte frame +func Send(mac, broadcast string) error // one UDP datagram via net.DialUDP("udp4", …); errors from parse/resolve/write +type Waker struct { /* private: targets, health, mutex, last-sent per host, poll interval (default 1s) */ } +func New(targets map[string]Target, h Health) *Waker +func (w *Waker) PollEvery(d time.Duration) // test hook; production keeps the 1 s default +// Wake returns true as soon as h.Healthy(host) is true, false if host is unknown, if Wait passes, +// or if ctx ends first. It sends the packet only if none was sent for host in the last Wait. +func (w *Waker) Wake(ctx context.Context, host string) bool +``` + +Rules the tests check: packet layout; separators; errors for bad MACs and unresolvable +addresses; one packet per window; return within about `Wait` when the host never comes up; +early return on a cancelled context; `false` for an unknown host without sending anything. +Never panic; safe for concurrent `Wake` calls on different hosts. + +## Steps + +- [ ] **1.** `git switch v2`; `mkdir -p internal/wake`; copy the test. +- [ ] **2. See it fail** (compile). **3. Write `wake.go`.** `gofmt -w internal/wake/`. +- [ ] **4.** `go test -race -count=3 ./internal/wake/` → `ok` (timing tests; three runs). +- [ ] **5.** `make gate`. **6.** Row `v2/03-wake`; commit. + +```sh +git add internal/wake docs/implementer-log.md +git commit +``` + +## Done when + +- `-race -count=3` passes; gate ok; the copied test is byte-identical. diff --git a/docs/plans/v2/04-identity.md b/docs/plans/v2/04-identity.md new file mode 100644 index 0000000..2a406b9 --- /dev/null +++ b/docs/plans/v2/04-identity.md @@ -0,0 +1,92 @@ +# v2 task 04: tailnet identity and the route gate; config additions + +**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add identity: whois resolver, checker, header mode, middleware; config for wake, peers, identity` + +## Goal + +Some routes should be usable only from particular tailnet nodes ("`hermes-talos` only from +talos"). `identity` resolves a caller's address to a tailnet node name — in production through +`tailscale whois --json `, in tests through a fake, in the smoke run through a header — and a +middleware in front of the proxy refuses other callers with 403. Config gains the keys the rest +of v2 needs. + +## Files + +- Copy: `internal/identity/identity_test.go`, `internal/identity/middleware_test.go`, `internal/identity/testdata/whois.json`, `internal/config/config_v2_test.go` +- Create: `internal/identity/identity.go`, `internal/identity/middleware.go` +- Modify: `internal/config/config.go`, `docs/implementer-log.md` + +## Interfaces + +```go +package identity + +var ( + ErrNotAPeer = errors.New("identity: not a tailnet peer") + ErrForbidden = errors.New("identity: forbidden route") +) +type ID struct{ Node, Login string } +type Resolver interface { Identity(ctx context.Context, ip string) (ID, error) } + +// ParseWhois reads `tailscale whois --json` output: Node = Node.ComputedName (else Node.Name +// without its trailing dot and domain), Login = UserProfile.LoginName. Empty node → error. +func ParseWhois(raw []byte) (ID, error) + +// TailscaleResolver runs `tailscale whois --json ` (exec, 3 s timeout) and parses it; a +// non-zero exit is ErrNotAPeer; a missing binary is an error that the Checker treats as "deny". +type TailscaleResolver struct{ Bin string } // Bin default "tailscale" +func (TailscaleResolver) Identity(ctx context.Context, ip string) (ID, error) + +type Checker struct { /* private: resolver, cache map[ip]ID with a 5-minute TTL, mutex */ } +func NewChecker(r Resolver) *Checker +// NewHeaderChecker trusts the X-Crossbar-Peer request header as the node name. TEST/SMOKE ONLY. +func NewHeaderChecker() *Checker +// Allow: nil when peers is empty (open route); otherwise the caller's node (from remoteAddr's +// IP, or the header in header mode) must be in peers, else ErrForbidden. Any resolver error, +// unparsable address or loopback → ErrForbidden. +func (c *Checker) Allow(ctx context.Context, peers []string, remoteAddr string) error + +// Middleware names the route like the proxy (X-Crossbar-Route header, else first path segment), +// asks peersFor(route), and answers 403 {"error":"forbidden route"} when Allow refuses. Paths +// under /_crossbar/ and routes peersFor does not know pass straight through. +func Middleware(c *Checker, peersFor func(route string) ([]string, bool), next http.Handler) http.Handler +``` + +Header mode: `Allow` needs the request to read the header, but its signature takes an address. +Make the header checker's resolver read from a `context.Context` value that `Middleware` sets +(`identity.WithHeaderPeer(ctx, r.Header.Get("X-Crossbar-Peer"))`); the tests only observe the +behaviour. Cache: per address, 5 minutes, for both hit and `ErrNotAPeer`. + +`internal/config` gains: + +```go +Identity string `toml:"identity"` // "off" (default) | "tailscale" | "header"; anything else → *Error field "identity" +// on Host: +Wake *Wake `toml:"wake"` // nil when absent +type Wake struct { MAC string `toml:"mac"`; Broadcast string `toml:"broadcast"`; Wait Duration `toml:"wait"` } +// on Route: +Peers []string `toml:"peers"` +``` + +Validation (after the existing host/route checks): `wake.mac` must parse (`net.ParseMAC`, 6 +bytes) → field `hosts..wake.mac`; `wake.broadcast` non-empty `host:port` → `hosts..wake.broadcast`; +`wake.wait` default 45 s, less than 5 s → `hosts..wake.wait`. `routes..peers` non-empty +while `identity == "off"` → `routes..peers` ("peers need identity = tailscale or header"); +`peers = []` (present but empty) with identity on → same field ("empty peers list"). + +## Steps + +- [ ] **1.** `git switch v2`; `mkdir -p internal/identity/testdata`; copy the four given files. +- [ ] **2. See them fail** (compile). **3. Write the code.** `gofmt -w internal/`. +- [ ] **4.** `go test -race -count=1 ./internal/identity/ ./internal/config/` → `ok` (v0/v1 config tests included). +- [ ] **5.** `make gate`. **6.** Row `v2/04-identity`; commit. + +```sh +git add internal/identity internal/config docs/implementer-log.md +git commit +``` + +## Done when + +- Both packages pass; gate ok; all four copied files byte-identical. diff --git a/docs/plans/v2/05-wiring-smoke.md b/docs/plans/v2/05-wiring-smoke.md new file mode 100644 index 0000000..9574bdd --- /dev/null +++ b/docs/plans/v2/05-wiring-smoke.md @@ -0,0 +1,58 @@ +# v2 task 05: wiring, the smoke run, README + +**Branch:** `v2` (run `git switch v2`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Wire wake and identity into crossbar; v2 smoke and README` + +## Goal + +Put the pieces together: the proxy wakes a sleeping host when a route has no healthy host left, +`main` builds the waker from config and wraps the proxy in the identity middleware when +`identity` is on, the given fake upstream can be woken, and `tools/smoke.sh` proves the whole of +v2 over real HTTP. + +## Files + +- Copy (**replaces** v1's): `cmd/fakeupstream/main.go`, `tools/smoke.sh`, `example.toml` +- Modify: `internal/proxy/proxy.go` (or a new file), `cmd/crossbar/main.go`, `README.md`, `docs/implementer-log.md` + +## Rules + +1. `proxy.Handler` gains `func (p *Handler) SetWaker(w Waker)` where + `type Waker interface{ Wake(ctx context.Context, host string) bool }` (defined in `proxy`). + When `leases.Acquire` returns `ErrNoHost` and a waker is set: for each candidate host of the + route, in order, that has a wake target (ask `cfg.Hosts[h].Wake != nil`), call + `Wake(r.Context(), h)`; on `true`, retry `Acquire` once; on `false` for every candidate, 503 + `{"error":"no healthy host","woke":[""]}`. The context-guard's "no host fits" + path (task 02) also tries waking a host whose `PerSlotCtx` is unknown or large enough, before + answering 400. +2. `cmd/crossbar`: build `wake.New(targets, hosts)` from every host with `Wake != nil` + (`hosts` is the `proxy.HostView`, which has `Healthy`), call `p.SetWaker(w)`. When + `cfg.Identity != "off"`: `checker := identity.NewChecker(identity.TailscaleResolver{})` or + `identity.NewHeaderChecker()`; wrap the proxy handler: + `mux.Handle("/", identity.Middleware(checker, func(route string) ([]string, bool) { rt, ok := cfg.Routes[route]; return rt.Peers, ok }, p))`. + Log at start which mode is active; with `"header"` log a warning that it is insecure. +3. Copy the three given files; `make build`; `make smoke` → `smoke: ok (…)`. The smoke's check 2 + waits up to 40 s for the wake; the fake wakes in ~1 s. +4. README: sections stay; add under `## Operate` the wake behaviour and the identity modes with + the `peers` example; under `## Configure` the three new keys; `## What v2 does not do`: + `/slots`, request coalescing, TLS — `PLAN.md`. + +## Steps + +- [ ] **1.** `git switch v2`; copy the three given files. +- [ ] **2. Write the code** (proxy waker path, `main`). `gofmt -w .` +- [ ] **3.** `go test -race -count=1 ./...` → all `ok`. **4.** `make smoke` → `smoke: ok`. +- [ ] **5.** README. **6.** `make gate`. **7.** Row `v2/05-wiring-smoke`; commit. + +```sh +git add internal/proxy cmd/crossbar cmd/fakeupstream tools/smoke.sh example.toml README.md docs/implementer-log.md +git commit +``` + +## Done when + +- `make smoke` prints `smoke: ok (…)`; gate ok; the three copied files byte-identical. + +## Stop and report if + +- `make smoke` fails twice in the same way; quote the failing check and the crossbar log. diff --git a/docs/plans/v2/README.md b/docs/plans/v2/README.md new file mode 100644 index 0000000..71f8834 --- /dev/null +++ b/docs/plans/v2/README.md @@ -0,0 +1,54 @@ +# v2 implementation plan: learned context, the context guard, wake-on-LAN, identity + +> **For the implementing model:** do not work from this file. The owner gives you one task file at +> a time (`01-…` to `05-…`). This file is the index for the owner and the reviewer. + +**Goal:** `PLAN.md` §4b and §10 v2. The poller learns each host's context size from `/props`; a +prompt that cannot fit the leased host's per-slot context moves to one where it fits or is +refused with a clear 400; a route whose hosts are all down can wake a sleeping host by +wake-on-LAN and wait for it; a route can be restricted to named tailnet peers. + +**Architecture:** `health.Status` gains `NCtx`/`Slots` (task 01); the proxy gains the guard +(task 02); two new small packages, `wake` (magic packets + a waiter, task 03) and `identity` +(whois resolver, checker, middleware, task 04); config gains `identity`, `[hosts.x.wake]`, +`routes.x.peers`; `main` wires the waker and the middleware (task 05). + +**How this plan was made:** acceptance tests first, from `PLAN.md`; no reference implementation. +Every given test compiled against a panic-only skeleton of the names in the tasks (`go vet` +clean). The given tests were walked against the task rules and against the other given files +(helpers, line limits, `main.go` call sites) before handover — the v1 findings list is the +reason. + +**Tech stack:** as v1; no new module. `identity`'s production resolver shells out to +`tailscale whois --json`, which exists on every fleet host. + +## Global constraints + +- Everything in `AGENTS.md`. Branch `v2`. One task, one fresh OpenCode session, one commit. +- Bodies never logged. Type assertions two-valued. Files under 400 lines. +- Given files are copied and never edited; some **replace** earlier ones (the task says so). + +## Tasks + +| # | File | Delivers | Tests that define it | +|---|---|---|---| +| 01 | `01-props.md` | `Status.NCtx`, `Status.Slots`, `PerSlotCtx()`; `/props` in the poll; hosts view shows them | `health/props_test.go` | +| 02 | `02-ctxguard.md` | prompt-size estimate; move or 400; `X-Crossbar-Ctx` | `proxy/ctxguard_test.go` | +| 03 | `03-wake.md` | `internal/wake`: magic packet, `Send`, `Waker` | `wake/wake_test.go` | +| 04 | `04-identity.md` | `internal/identity`: whois parse, checker, header mode, middleware; config `identity`/`peers`/`wake` | `identity/*_test.go`, `config/config_v2_test.go` | +| 05 | `05-wiring-smoke.md` | proxy wakes on no-host; `main` wires waker + middleware; given fakeupstream/smoke/example; README | `make smoke` | + +## For the owner + +`tools/run-plan.sh docs/plans/v2` from a clean checkout on `master`. + +## For the reviewer: after task 05 + +1. Five task commits with the trailer; given files byte-identical; protected files untouched. +2. `make gate`, `make smoke`. +3. Probe: a `/props` that returns 200 with a huge body (bounded read); a MAC with an unusual + separator in config; `identity = "tailscale"` on a host where `tailscale` is not on PATH + (must log and refuse the gated routes, never allow); two routes, one gated one open, from + the same peer; a wake target whose broadcast address is unroutable (503 within `wait`, no + hang); the guard with a body of exactly `MaxBody`. +4. Findings under "Reviews" in `docs/implementer-log.md`, by fault. diff --git a/docs/plans/v2/_files/cmd/fakeupstream/main.go b/docs/plans/v2/_files/cmd/fakeupstream/main.go new file mode 100644 index 0000000..0cb6b44 --- /dev/null +++ b/docs/plans/v2/_files/cmd/fakeupstream/main.go @@ -0,0 +1,161 @@ +// fakeupstream stands in for a llama-server router in tests and the smoke run. Do not edit. +// +// fakeupstream -listen 127.0.0.1:18081 -name alpha -models a,b -down-file /tmp/alpha.down -slow 0 +// +// /health answers 503 while the down file exists, 200 otherwise. /v1/models lists -models. +// /props answers a small JSON object. /v1/chat/completions echoes: a streamed answer of five +// SSE chunks 200 ms apart when the body has "stream": true, then a final chunk carrying +// "usage" and llama-server style "timings", then [DONE]; one JSON answer with usage and +// timings otherwise. -slow adds that many milliseconds before answering (for queue tests). +// Every response carries X-Upstream: . /props reports -n-ctx and -slots. With -wol-listen, +// a valid wake-on-LAN magic packet for -wol-mac received on that UDP address removes the down +// file, so the fake "boots" when woken. +package main + +import ( + "encoding/json" + "flag" + "fmt" + "io" + "log" + "net" + "net/http" + "os" + "strings" + "time" +) + +func main() { + listen := flag.String("listen", "127.0.0.1:18081", "address to listen on") + name := flag.String("name", "fake", "name reported in X-Upstream and answers") + models := flag.String("models", "m", "comma-separated model ids for /v1/models") + downFile := flag.String("down-file", "", "while this file exists, /health answers 503") + slow := flag.Int("slow", 0, "milliseconds to wait before answering a completion") + nCtx := flag.Int("n-ctx", 8192, "n_ctx reported by /props") + slots := flag.Int("slots", 2, "total_slots reported by /props") + wolListen := flag.String("wol-listen", "", "UDP address to listen on for a wake-on-LAN magic packet") + wolMAC := flag.String("wol-mac", "aa:bb:cc:dd:ee:01", "MAC the magic packet must carry") + flag.Parse() + if *wolListen != "" && *downFile != "" { + go wakeOnPacket(*wolListen, *wolMAC, *downFile) + } + + ids := strings.Split(*models, ",") + mux := http.NewServeMux() + stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) } + usage := map[string]any{"prompt_tokens": 100, "completion_tokens": 10, "total_tokens": 110} + timings := map[string]any{"prompt_n": 100, "cache_n": 90, "predicted_n": 10, "predicted_ms": 50.0} + + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { + stamp(w) + if *downFile != "" { + if _, err := os.Stat(*downFile); err == nil { + http.Error(w, `{"error":{"message":"Loading model"}}`, http.StatusServiceUnavailable) + return + } + } + writeJSON(w, map[string]string{"status": "ok"}) + }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { + stamp(w) + data := []map[string]any{} + for _, id := range ids { + data = append(data, map[string]any{"id": id, "object": "model", "owned_by": *name}) + } + writeJSON(w, map[string]any{"object": "list", "data": data}) + }) + mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { + stamp(w) + writeJSON(w, map[string]any{"default_generation_settings": map[string]any{"n_ctx": *nCtx}, "total_slots": *slots, "model_path": *name}) + }) + mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) { + stamp(w) + body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + var req struct { + Model string `json:"model"` + Stream bool `json:"stream"` + } + _ = json.Unmarshal(body, &req) + time.Sleep(time.Duration(*slow) * time.Millisecond) + if !req.Stream { + writeJSON(w, map[string]any{ + "id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model, + "choices": []map[string]any{{"index": 0, "message": map[string]string{"role": "assistant", "content": "hello from " + *name}, "finish_reason": "stop"}}, + "usage": usage, "timings": timings, + }) + return + } + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.WriteHeader(http.StatusOK) + fl, _ := w.(http.Flusher) + flush := func() { + if fl != nil { + fl.Flush() + } + } + for i := 1; i <= 5; i++ { + chunk := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model, + "choices": []map[string]any{{"index": 0, "delta": map[string]string{"content": fmt.Sprintf("%s chunk %d ", *name, i)}}}} + b, _ := json.Marshal(chunk) + fmt.Fprintf(w, "data: %s\n\n", b) + flush() + time.Sleep(200 * time.Millisecond) + } + final := map[string]any{"id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model, + "choices": []map[string]any{}, "usage": usage, "timings": timings} + b, _ := json.Marshal(final) + fmt.Fprintf(w, "data: %s\n\n", b) + flush() + fmt.Fprint(w, "data: [DONE]\n\n") + }) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + stamp(w) + http.Error(w, `{"error":"not found"}`, http.StatusNotFound) + }) + + log.Printf("fakeupstream %s listening on %s models=%v slow=%dms", *name, *listen, ids, *slow) + srv := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second} + log.Fatal(srv.ListenAndServe()) +} + +func writeJSON(w http.ResponseWriter, v any) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(v) +} + +// wakeOnPacket removes downFile when a magic packet for mac arrives: 6×0xff then the MAC 16 times. +func wakeOnPacket(addr, mac, downFile string) { + hw, err := net.ParseMAC(mac) + if err != nil { + log.Fatalf("wol-mac: %v", err) + } + pc, err := net.ListenPacket("udp4", addr) + if err != nil { + log.Fatalf("wol-listen: %v", err) + } + log.Printf("fakeupstream listening for wake-on-LAN on %s (mac %s)", addr, hw) + buf := make([]byte, 256) + for { + n, _, err := pc.ReadFrom(buf) + if err != nil { + return + } + if n != 102 { + continue + } + ok := true + for i := 0; i < 6; i++ { + ok = ok && buf[i] == 0xff + } + for i := 0; i < 16 && ok; i++ { + for j := 0; j < 6; j++ { + ok = ok && buf[6+6*i+j] == hw[j] + } + } + if ok { + log.Printf("magic packet received: waking (removing %s)", downFile) + _ = os.Remove(downFile) + } + } +} diff --git a/docs/plans/v2/_files/example.toml b/docs/plans/v2/_files/example.toml new file mode 100644 index 0000000..03f0372 --- /dev/null +++ b/docs/plans/v2/_files/example.toml @@ -0,0 +1,32 @@ +# 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 +identity = "off" # "tailscale" gates routes with `peers` by `tailscale whois`; "header" trusts X-Crossbar-Peer (TEST ONLY) + +[hosts.alpha] +base_url = "http://127.0.0.1:18081" # e.g. http://straylight.:11434 +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 } } +[hosts.beta.wake] # v2: wake a sleeping host when nothing else can take a new lease +mac = "aa:bb:cc:dd:ee:02" +broadcast = "127.0.0.1:19082" # the LAN broadcast address, port 9, in production +wait = "20s" + +# v1: a route is a set of candidate hosts; each conversation gets a sticky lease on the host with +# the most free slots × weight at the time it starts. Pins and drains come from the admin API. +[routes.opencode-a] +hosts = ["alpha", "beta"] +default_model = "ornith-1.5-35b-a3b" + +[routes.hermes-x] +hosts = ["beta", "alpha"] +# peers = ["talos"] # v2: with identity = "tailscale", only these tailnet nodes may use the route diff --git a/docs/plans/v2/_files/internal/config/config_v2_test.go b/docs/plans/v2/_files/internal/config/config_v2_test.go new file mode 100644 index 0000000..0c4af39 --- /dev/null +++ b/docs/plans/v2/_files/internal/config/config_v2_test.go @@ -0,0 +1,75 @@ +package config_test + +import ( + "strings" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +const v2Base = ` +listen = "127.0.0.1:1" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[hosts.b] +base_url = "http://b:1" +models = { "m" = { } } +[hosts.b.wake] +mac = "aa:bb:cc:dd:ee:ff" +broadcast = "192.168.1.255:9" +wait = "45s" +[routes.r] +hosts = ["a", "b"] +peers = ["talos", "imladris"] +` + +func TestV2Defaults(t *testing.T) { + c, err := config.Parse(strings.NewReader(v2Base)) + if err != nil { + t.Fatal(err) + } + if c.Identity != "off" { + t.Errorf("identity default = %q, want off", c.Identity) + } + if c.Hosts["a"].Wake != nil { + t.Errorf("host without [wake] must have nil Wake") + } + w := c.Hosts["b"].Wake + if w == nil || w.MAC != "aa:bb:cc:dd:ee:ff" || w.Broadcast != "192.168.1.255:9" || w.Wait.Duration != 45*time.Second { + t.Errorf("wake = %+v", w) + } + if p := c.Routes["r"].Peers; len(p) != 2 || p[0] != "talos" { + t.Errorf("peers = %v", p) + } +} + +func TestV2Validation(t *testing.T) { + good := v2Base + for _, tc := range []struct{ name, text, field string }{ + {"bad identity", "identity = \"maybe\"\n" + good, "identity"}, + {"peers without identity", "identity = \"off\"\n" + good, "routes.r.peers"}, + {"bad mac", strings.Replace(good, `mac = "aa:bb:cc:dd:ee:ff"`, `mac = "nope"`, 1), "hosts.b.wake.mac"}, + {"no broadcast", strings.Replace(good, `broadcast = "192.168.1.255:9"`, `broadcast = ""`, 1), "hosts.b.wake.broadcast"}, + {"wait too short", strings.Replace(good, `wait = "45s"`, `wait = "2s"`, 1), "hosts.b.wake.wait"}, + {"peers on unknown route field", "identity = \"tailscale\"\n" + strings.Replace(good, `peers = ["talos", "imladris"]`, `peers = []`, 1), "routes.r.peers"}, + } { + _, err := config.Parse(strings.NewReader(tc.text)) + e, ok := config.IsError(err) + if !ok || e.Field != tc.field { + t.Errorf("%s: %v, want *Error on %s", tc.name, err, tc.field) + } + } + // identity = "header" is the test/smoke mode; "tailscale" the real one; both accept peers. + for _, mode := range []string{"header", "tailscale"} { + if _, err := config.Parse(strings.NewReader("identity = \"" + mode + "\"\n" + good)); err != nil { + t.Errorf("identity=%s with peers: %v", mode, err) + } + } + // wait defaults to 45s when the [wake] table omits it + c, err := config.Parse(strings.NewReader("identity = \"header\"\n" + strings.Replace(good, "wait = \"45s\"\n", "", 1))) + if err != nil || c.Hosts["b"].Wake == nil || c.Hosts["b"].Wake.Wait.Duration != 45*time.Second { + t.Errorf("wake.wait default: %v %+v", err, c.Hosts["b"].Wake) + } +} diff --git a/docs/plans/v2/_files/internal/health/props_test.go b/docs/plans/v2/_files/internal/health/props_test.go new file mode 100644 index 0000000..ed77b95 --- /dev/null +++ b/docs/plans/v2/_files/internal/health/props_test.go @@ -0,0 +1,75 @@ +package health_test + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/health" +) + +// propsFake answers /health, /v1/models and a configurable /props. +func propsFake(t *testing.T, props string, status int) *httptest.Server { + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"m"}]}`) }) + mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(status) + fmt.Fprint(w, props) + }) + srv := httptest.NewServer(mux) + t.Cleanup(srv.Close) + return srv +} + +func TestPropsLearned(t *testing.T) { + srv := propsFake(t, `{"default_generation_settings":{"n_ctx":131072,"params":{}},"total_slots":4,"model_path":"/x/m.gguf","chat_template":"..."}`, 200) + tbl := health.New(map[string]string{"a": srv.URL}, time.Hour, nil) + tbl.PollOnce(context.Background()) + s, _ := tbl.Get("a") + if !s.Healthy || s.NCtx != 131072 || s.Slots != 4 { + t.Fatalf("status = %+v, want healthy with NCtx 131072 and Slots 4", s) + } + if got := s.PerSlotCtx(); got != 32768 { + t.Errorf("PerSlotCtx = %d, want 131072/4", got) + } +} + +func TestPropsAbsentOrBrokenIsNotAFailure(t *testing.T) { + for name, tc := range map[string]struct { + props string + status int + }{ + "404": {`not found`, 404}, + "not json": {``, 200}, + "no fields": {`{"model_path":"/x"}`, 200}, + "zero ctx": {`{"default_generation_settings":{"n_ctx":0},"total_slots":0}`, 200}, + } { + t.Run(name, func(t *testing.T) { + srv := propsFake(t, tc.props, tc.status) + tbl := health.New(map[string]string{"a": srv.URL}, time.Hour, nil) + tbl.PollOnce(context.Background()) + s, _ := tbl.Get("a") + if !s.Healthy { + t.Fatalf("a bad /props must not make the host unhealthy: %+v", s) + } + if s.NCtx != 0 || s.Slots != 0 || s.PerSlotCtx() != 0 { + t.Errorf("unknown context must read as 0: %+v", s) + } + }) + } +} + +func TestPerSlotCtxWithUnknownSlots(t *testing.T) { + s := health.Status{NCtx: 8192, Slots: 0} + if s.PerSlotCtx() != 8192 { + t.Errorf("with Slots unknown the whole context is the per-slot value; got %d", s.PerSlotCtx()) + } + s = health.Status{NCtx: 8192, Slots: 3} + if s.PerSlotCtx() != 2730 { + t.Errorf("integer division: got %d, want 2730", s.PerSlotCtx()) + } +} diff --git a/docs/plans/v2/_files/internal/identity/identity_test.go b/docs/plans/v2/_files/internal/identity/identity_test.go new file mode 100644 index 0000000..b75e2de --- /dev/null +++ b/docs/plans/v2/_files/internal/identity/identity_test.go @@ -0,0 +1,90 @@ +package identity_test + +import ( + "context" + "errors" + "os" + "path/filepath" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/identity" +) + +func TestParseWhois(t *testing.T) { + raw, err := os.ReadFile(filepath.Join("testdata", "whois.json")) + if err != nil { + t.Fatal(err) + } + id, err := identity.ParseWhois(raw) + if err != nil { + t.Fatal(err) + } + if id.Node != "titan" || id.Login == "" { + t.Errorf("parsed %+v, want Node titan and a login", id) + } + if _, err := identity.ParseWhois([]byte(`{"Node":{}}`)); err == nil { + t.Error("a whois answer without a node name must be an error") + } + if _, err := identity.ParseWhois([]byte(`nope`)); err == nil { + t.Error("non-JSON must be an error") + } +} + +// fakeResolver answers from a map; "" means not a tailnet peer. +type fakeResolver map[string]string + +func (f fakeResolver) Identity(ctx context.Context, ip string) (identity.ID, error) { + n, ok := f[ip] + if !ok { + return identity.ID{}, identity.ErrNotAPeer + } + return identity.ID{Node: n, Login: n + "@example"}, nil +} + +func TestChecker(t *testing.T) { + c := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"}) + for _, tc := range []struct { + name string + peers []string + addr string + want error + }{ + {"open route", nil, "203.0.113.7:1", nil}, + {"allowed peer", []string{"talos", "titan"}, "100.64.0.5:44444", nil}, + {"other peer", []string{"talos"}, "100.64.0.9:1", identity.ErrForbidden}, + {"not a peer", []string{"talos"}, "203.0.113.7:1", identity.ErrForbidden}, + {"loopback", []string{"talos"}, "127.0.0.1:1", identity.ErrForbidden}, + {"garbage addr", []string{"talos"}, "nonsense", identity.ErrForbidden}, + } { + t.Run(tc.name, func(t *testing.T) { + got := c.Allow(context.Background(), tc.peers, tc.addr) + if !errors.Is(got, tc.want) && !(got == nil && tc.want == nil) { + t.Errorf("Allow(%v, %q) = %v, want %v", tc.peers, tc.addr, got, tc.want) + } + }) + } +} + +func TestCheckerCachesPerAddress(t *testing.T) { + calls := 0 + r := countingResolver{f: fakeResolver{"100.64.0.5": "talos"}, calls: &calls} + c := identity.NewChecker(r) + for i := 0; i < 5; i++ { + if err := c.Allow(context.Background(), []string{"talos"}, "100.64.0.5:1"); err != nil { + t.Fatal(err) + } + } + if calls != 1 { + t.Errorf("resolver called %d times for one address, want 1 (cache)", calls) + } +} + +type countingResolver struct { + f fakeResolver + calls *int +} + +func (c countingResolver) Identity(ctx context.Context, ip string) (identity.ID, error) { + *c.calls++ + return c.f.Identity(ctx, ip) +} diff --git a/docs/plans/v2/_files/internal/identity/middleware_test.go b/docs/plans/v2/_files/internal/identity/middleware_test.go new file mode 100644 index 0000000..19bf778 --- /dev/null +++ b/docs/plans/v2/_files/internal/identity/middleware_test.go @@ -0,0 +1,74 @@ +package identity_test + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/identity" +) + +// The middleware sits in front of the proxy: it names the route the same way the proxy does +// (X-Crossbar-Route header, else first path segment) and refuses callers a route does not list. +func TestMiddleware(t *testing.T) { + peers := map[string][]string{"locked": {"talos"}, "open": nil} + inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) }) + h := identity.Middleware(identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"}), + func(route string) ([]string, bool) { p, ok := peers[route]; return p, ok }, inner) + for _, tc := range []struct { + name, path, hdr, addr string + want int + }{ + {"open route, anyone", "/open/v1/models", "", "203.0.113.1:5", 204}, + {"locked, right peer", "/locked/v1/models", "", "100.64.0.5:5", 204}, + {"locked, wrong peer", "/locked/v1/models", "", "100.64.0.9:5", 403}, + {"locked, not a peer", "/locked/v1/models", "", "203.0.113.1:5", 403}, + {"locked via header", "/v1/models", "locked", "100.64.0.9:5", 403}, + {"header wins over path", "/open/v1/models", "locked", "203.0.113.1:5", 403}, + {"unknown route passes through to the proxy's own 404", "/nope/v1/models", "", "203.0.113.1:5", 204}, + {"admin path is never gated here", "/_crossbar/hosts", "", "203.0.113.1:5", 204}, + } { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, tc.path, nil) + req.RemoteAddr = tc.addr + if tc.hdr != "" { + req.Header.Set("X-Crossbar-Route", tc.hdr) + } + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + if rec.Code != tc.want { + t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String()) + } + if rec.Code == 403 && (!strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") || !strings.Contains(rec.Body.String(), `"forbidden route"`)) { + t.Errorf("403 must be JSON {\"error\":\"forbidden route\"}: %q", rec.Body.String()) + } + }) + } +} + +// HeaderResolver is the test/smoke identity source: it trusts X-Crossbar-Peer. It exists so the +// smoke run can exercise the gate without a tailnet; config must call it out as insecure. +func TestHeaderResolver(t *testing.T) { + inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) }) + h := identity.Middleware(identity.NewHeaderChecker(), func(route string) ([]string, bool) { return []string{"talos"}, true }, inner) + req := httptest.NewRequest(http.MethodGet, "/r/v1/models", nil) + req.Header.Set("X-Crossbar-Peer", "talos") + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + if rec.Code != 204 { + t.Errorf("header peer talos: %d", rec.Code) + } + req.Header.Set("X-Crossbar-Peer", "titan") + rec = httptest.NewRecorder() + h.ServeHTTP(rec, req) + if rec.Code != 403 { + t.Errorf("header peer titan: %d, want 403", rec.Code) + } + req.Header.Del("X-Crossbar-Peer") + rec = httptest.NewRecorder() + h.ServeHTTP(rec, req) + if rec.Code != 403 { + t.Errorf("no header: %d, want 403", rec.Code) + } +} diff --git a/docs/plans/v2/_files/internal/identity/testdata/whois.json b/docs/plans/v2/_files/internal/identity/testdata/whois.json new file mode 100644 index 0000000..877b5a8 --- /dev/null +++ b/docs/plans/v2/_files/internal/identity/testdata/whois.json @@ -0,0 +1,24 @@ +{ + "Node": { + "ID": 1, + "StableID": "nEXAMPLE", + "Name": "titan.example.ts.net.", + "User": 2, + "Addresses": [ + "100.64.0.9/32", + "fd7a:115c:a1e0::9/128" + ], + "HomeDERP": 2, + "Created": "2026-01-01T00:00:00Z", + "Cap": 138, + "Online": true, + "ComputedName": "titan", + "ComputedNameWithHost": "titan" + }, + "UserProfile": { + "ID": 2, + "LoginName": "user@example.com", + "DisplayName": "Example User" + }, + "CapMap": null +} \ No newline at end of file diff --git a/docs/plans/v2/_files/internal/proxy/ctxguard_test.go b/docs/plans/v2/_files/internal/proxy/ctxguard_test.go new file mode 100644 index 0000000..65235f0 --- /dev/null +++ b/docs/plans/v2/_files/internal/proxy/ctxguard_test.go @@ -0,0 +1,142 @@ +package proxy_test + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "git.wntrmute.dev/kyle/crossbar/internal/proxy" +) + +// ctxUpstream is a fake router that reports a context size in /props and echoes completions. +func ctxUpstream(t *testing.T, name string, nCtx, slots int) *upstream { + u := &upstream{name: name} + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"shared"}]}`) }) + mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprintf(w, `{"default_generation_settings":{"n_ctx":%d},"total_slots":%d}`, nCtx, slots) + }) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + u.hits.Add(1) + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"choices":[{"message":{"role":"assistant","content":"ok"}}],"usage":{"prompt_tokens":1,"completion_tokens":1}}`) + }) + u.srv = httptest.NewServer(mux) + t.Cleanup(u.srv.Close) + return u +} + +const ctxHosts = ` +listen = "127.0.0.1:1" +queue_max = 2 +[hosts.small] +base_url = %q +weight = 10.0 +models = { "shared" = { parallel = 2 } } +[hosts.big] +base_url = %q +weight = 1.0 +models = { "shared" = { parallel = 1 } } +[routes.r] +hosts = ["small", "big"] +default_model = "shared" +` + +// bodyOfTokens builds a chat body whose byte size implies roughly n tokens under the guard's +// estimate (bytes/4 × 1.2): n tokens ≈ 3.33 n bytes ≈ 2n/3 five-byte words. +func bodyOfTokens(n int) string { + text := strings.Repeat("word ", n*2/3) + return fmt.Sprintf(`{"model":"shared","stream":false,"messages":[{"role":"user","content":"%s"}]}`, text) +} + +func TestOversizedPromptMovesToAHostWhereItFits(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) // 4096 per slot + big := ctxUpstream(t, "big", 131072, 1) // 131072 per slot + r := newRig(t, ctxHosts, small, big) + // A small prompt starts on `small` (weight 10). + resp := r.post("/r/v1/chat/completions", bodyOfTokens(100)) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" { + t.Fatalf("small prompt went to %q, want small", resp.Header.Get(proxy.HostHeader)) + } + // A new conversation with ~10k tokens does not fit small's 4096-token slot: it must be + // placed on big, with the reason visible in a header. + resp = r.post("/r/v1/chat/completions", bodyOfTokens(10000)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" { + t.Fatalf("oversized prompt: %d from %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader)) + } + if got := resp.Header.Get(proxy.CtxHeader); !strings.HasPrefix(got, "moved") { + t.Errorf("%s = %q, want moved:… ", proxy.CtxHeader, got) + } +} + +func TestOversizedPromptWithNoFitIs400(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) + tiny := ctxUpstream(t, "big", 4096, 2) // also too small + r := newRig(t, ctxHosts, small, tiny) + resp := r.post("/r/v1/chat/completions", bodyOfTokens(10000)) + body := drain(resp) + if resp.StatusCode != http.StatusBadRequest { + t.Fatalf("status %d body %s, want 400", resp.StatusCode, body) + } + var e map[string]any + if err := json.Unmarshal([]byte(body), &e); err != nil || e["error"] != "prompt too large" { + t.Fatalf("body = %s, want error 'prompt too large'", body) + } + if est, _ := e["estimate"].(float64); est < 8000 || est > 13000 { + t.Errorf("estimate = %v, want roughly 10000 tokens", e["estimate"]) + } + if max, _ := e["max"].(float64); max != 4096 { + t.Errorf("max = %v, want the largest per-slot context among the route's hosts (4096)", e["max"]) + } + if small.hits.Load()+tiny.hits.Load() != 0 { + t.Errorf("a refused prompt must not reach any upstream") + } +} + +func TestUnknownContextNeverBlocks(t *testing.T) { + // /props missing on both hosts: NCtx 0 means "unknown", and the guard must stay out of the way. + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + r := newRig(t, twoHosts, alpha, beta) + resp := r.post("/r/v1/chat/completions", bodyOfTokens(50000)) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.CtxHeader) != "" { + t.Errorf("unknown context: %d %q, want 200 and no ctx header", resp.StatusCode, resp.Header.Get(proxy.CtxHeader)) + } +} + +func TestStickyLeaseSurvivesGrowthUntilItDoesNotFit(t *testing.T) { + small := ctxUpstream(t, "small", 8192, 2) + big := ctxUpstream(t, "big", 131072, 1) + r := newRig(t, ctxHosts, small, big) + body := bodyOfTokens(100) + resp := r.post("/r/v1/chat/completions", body) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" { + t.Fatal("setup: first turn must be on small") + } + // Same conversation (same first user message), later turn well under 4096: stays. + longer := strings.Replace(body, `"content":"`, `"content":"`+strings.Repeat("x ", 500), 1) + resp = r.post("/r/v1/chat/completions", longer) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "small" || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Errorf("turn 2: %q %q, want small reused", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } + // A turn that outgrows the slot moves the lease — once — and the move is recorded as an event. + huge := strings.Replace(body, `"content":"`, `"content":"`+strings.Repeat("x ", 30000), 1) + resp = r.post("/r/v1/chat/completions", huge) + drain(resp) + if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" { + t.Fatalf("outgrown turn: %d %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader)) + } + resp = r.post("/r/v1/chat/completions", huge) + drain(resp) + if resp.Header.Get(proxy.HostHeader) != "big" || resp.Header.Get(proxy.LeaseHeader) != "reused" { + t.Errorf("after the move the lease is on big: %q %q", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader)) + } +} diff --git a/docs/plans/v2/_files/internal/wake/wake_test.go b/docs/plans/v2/_files/internal/wake/wake_test.go new file mode 100644 index 0000000..53e6179 --- /dev/null +++ b/docs/plans/v2/_files/internal/wake/wake_test.go @@ -0,0 +1,114 @@ +package wake_test + +import ( + "bytes" + "context" + "net" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/wake" +) + +func listen(t *testing.T) (*net.UDPConn, string) { + conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { conn.Close() }) + return conn, conn.LocalAddr().String() +} + +func TestMagicPacket(t *testing.T) { + pkt, err := wake.MagicPacket("aa:bb:cc:dd:ee:ff") + if err != nil { + t.Fatal(err) + } + if len(pkt) != 102 || !bytes.Equal(pkt[:6], bytes.Repeat([]byte{0xff}, 6)) { + t.Fatalf("packet = % x", pkt) + } + mac := []byte{0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff} + for i := 0; i < 16; i++ { + if !bytes.Equal(pkt[6+6*i:12+6*i], mac) { + t.Fatalf("repetition %d wrong: % x", i, pkt[6+6*i:12+6*i]) + } + } + for _, bad := range []string{"", "aa:bb", "zz:bb:cc:dd:ee:ff", "aabbccddeeff00"} { + if _, err := wake.MagicPacket(bad); err == nil { + t.Errorf("MagicPacket(%q) must fail", bad) + } + } + if p2, _ := wake.MagicPacket("AA-BB-CC-DD-EE-FF"); !bytes.Equal(p2, pkt) { + t.Errorf("dash-separated upper-case MAC must give the same packet") + } +} + +func TestSendReachesTheBroadcastAddress(t *testing.T) { + conn, addr := listen(t) + if err := wake.Send("aa:bb:cc:dd:ee:ff", addr); err != nil { + t.Fatal(err) + } + buf := make([]byte, 200) + _ = conn.SetReadDeadline(time.Now().Add(time.Second)) + n, _, err := conn.ReadFromUDP(buf) + if err != nil || n != 102 { + t.Fatalf("received %d bytes, err %v", n, err) + } + if err := wake.Send("aa:bb:cc:dd:ee:ff", "256.1.1.1:9"); err == nil { + t.Error("an unresolvable broadcast address must be an error") + } +} + +// fakeHealth flips to healthy after `after` calls to Healthy. +type fakeHealth struct{ calls, after int } + +func (f *fakeHealth) Healthy(name string) bool { f.calls++; return f.calls > f.after } + +func TestWakerSendsOncePerWindowAndWaitsForHealth(t *testing.T) { + conn, addr := listen(t) + h := &fakeHealth{after: 3} + w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 2 * time.Second}}, h) + w.PollEvery(20 * time.Millisecond) // test hook: how often Wake re-checks health + start := time.Now() + ok := w.Wake(context.Background(), "titan") + if !ok { + t.Fatal("Wake must return true once the host reports healthy") + } + if time.Since(start) > time.Second { + t.Errorf("Wake waited %v for a host that came up after 3 checks", time.Since(start)) + } + _ = conn.SetReadDeadline(time.Now().Add(200 * time.Millisecond)) + buf := make([]byte, 200) + if n, _, err := conn.ReadFromUDP(buf); err != nil || n != 102 { + t.Fatalf("no magic packet received: %d %v", n, err) + } + // A second Wake inside the same window does not send again (the host is booting). + _ = w.Wake(context.Background(), "titan") + _ = conn.SetReadDeadline(time.Now().Add(150 * time.Millisecond)) + if n, _, err := conn.ReadFromUDP(buf); err == nil { + t.Errorf("a second packet (%d bytes) was sent inside the wait window", n) + } + if w.Wake(context.Background(), "nobody") { + t.Errorf("unknown host: Wake must return false") + } +} + +func TestWakeGivesUpAfterWait(t *testing.T) { + _, addr := listen(t) + h := &fakeHealth{after: 1 << 30} + w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 300 * time.Millisecond}}, h) + w.PollEvery(20 * time.Millisecond) + start := time.Now() + if w.Wake(context.Background(), "titan") { + t.Fatal("Wake must return false when the host never comes up") + } + if d := time.Since(start); d < 250*time.Millisecond || d > 900*time.Millisecond { + t.Errorf("Wake returned after %v, want about the 300ms wait", d) + } + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + start = time.Now() + if w.Wake(ctx, "titan") || time.Since(start) > 200*time.Millisecond { + t.Errorf("a cancelled context must end the wait early (took %v)", time.Since(start)) + } +} diff --git a/docs/plans/v2/_files/tools/smoke.sh b/docs/plans/v2/_files/tools/smoke.sh new file mode 100755 index 0000000..965deb9 --- /dev/null +++ b/docs/plans/v2/_files/tools/smoke.sh @@ -0,0 +1,58 @@ +#!/bin/sh +# Smoke run (v2): everything v1 checked, plus the context guard, wake-on-LAN and identity gating. +# Prints "smoke: ok" or fails with the crossbar log. +set -eu +cd "$(dirname "$0")/.." +tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM +pids="" +sed -e "s#^db .*#db = \"$tmp/crossbar.db\"#" -e 's#^identity .*#identity = "header"#' -e 's#^\# peers = \["talos"\]#peers = ["talos"]#' example.toml > "$tmp/crossbar.toml" +# alpha: small context (4096 per slot = 8192/2); beta: large, sleeps until woken +bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" -slow 600 -n-ctx 8192 -slots 2 >"$tmp/alpha.log" 2>&1 & pids="$pids $!" +bin/fakeupstream -listen 127.0.0.1:18082 -name beta -models ornith-1.5-35b-a3b -down-file "$tmp/beta.down" -n-ctx 131072 -slots 2 -wol-listen 127.0.0.1:19082 -wol-mac aa:bb:cc:dd:ee:02 >"$tmp/beta.log" 2>&1 & pids="$pids $!" +touch "$tmp/beta.down" # beta starts "asleep" +bin/crossbar -config "$tmp/crossbar.toml" >"$tmp/crossbar.log" 2>&1 & pids="$pids $!" +sleep 2.5 # two polls: alpha healthy, beta down +fail() { echo "smoke: FAIL: $*" >&2; echo "--- crossbar.log"; cat "$tmp/crossbar.log"; exit 1; } +base=http://127.0.0.1:17777 +conv() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"system","content":"smoke"},{"role":"user","content":"conversation %s"}]}' "$1"; } +big() { printf '{"model":"ornith-1.5-35b-a3b","stream":false,"messages":[{"role":"user","content":"%s"}]}' "$(head -c 40000 /dev/zero | tr '\0' 'x')"; } +hdrs() { curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host} %header{X-Crossbar-Lease}' "$@"; } + +# 1. v1 behaviour: with beta asleep, opencode-a goes to alpha +h=$(hdrs -X POST -H 'Content-Type: application/json' -d "$(conv A)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 alpha new" ] || fail "with beta asleep conversation A should be '200 alpha new', got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{[^}]*"n_ctx":8192' || fail "hosts view does not show alpha n_ctx 8192: $(curl -s $base/_crossbar/hosts)" + +# 2. context guard: a ~12k-token prompt does not fit alpha's 4096-token slot; beta is asleep and +# wakeable, so crossbar must send the magic packet, wait for beta, and place the prompt there. +start=$(date +%s) +h=$(hdrs -m 40 -X POST -H 'Content-Type: application/json' -d "$(big)" "$base/opencode-a/v1/chat/completions") +[ "$h" = "200 beta new" ] || fail "oversized prompt should wake beta and land there, got '$h' after $(( $(date +%s) - start ))s" +grep -q "magic packet received" "$tmp/beta.log" || fail "beta never saw a magic packet" +curl -s "$base/_crossbar/hosts" | grep -q '"beta":{"healthy":true' || fail "beta not healthy after wake" + +# 3. with beta awake, a prompt that fits nowhere is a 400 (both slots too small? no — beta fits): +# check the guard's refusal with a prompt beyond beta's 65536-per-slot too +h=$(curl -s -o "$tmp/toolarge.json" -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d "$(printf '{"model":"ornith-1.5-35b-a3b","messages":[{"role":"user","content":"%s"}]}' "$(head -c 300000 /dev/zero | tr '\0' 'x')")" "$base/opencode-a/v1/chat/completions") +[ "$h" = "400" ] && grep -q '"prompt too large"' "$tmp/toolarge.json" || fail "300 KB prompt should be 400 prompt too large, got $h $(cat "$tmp/toolarge.json")" + +# 4. identity: hermes-x is locked to peer talos (header mode) +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x without a peer header should be 403, got $h" +h=$(curl -s -o /dev/null -w '%{http_code}' -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: titan' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +[ "$h" = "403" ] || fail "hermes-x as titan should be 403, got $h" +h=$(hdrs -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d "$(conv B)" "$base/hermes-x/v1/chat/completions") +case "$h" in 200*) ;; *) fail "hermes-x as talos should be 200, got '$h'";; esac +h=$(curl -s -o /dev/null -w '%{http_code}' "$base/_crossbar/hosts"); [ "$h" = "200" ] || fail "admin must not be gated, got $h" + +# 5. v1 regression: streaming still incremental, usage and metrics present +start=$(date +%s%N) +curl -sN -X POST -H 'Content-Type: application/json' -H 'X-Crossbar-Peer: talos' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[{"role":"user","content":"stream me"}]}' \ + "$base/hermes-x/v1/chat/completions" | while IFS= read -r line; do [ -n "$line" ] || continue; now=$(date +%s%N); echo "$(( (now - start) / 1000000 )) $line"; done > "$tmp/stream.txt" +firstms=$(head -1 "$tmp/stream.txt" | cut -d' ' -f1); lastms=$(tail -1 "$tmp/stream.txt" | cut -d' ' -f1) +[ -n "$firstms" ] && [ "$((lastms - firstms))" -ge 600 ] || fail "stream arrived in one burst" +sleep 1 +curl -s "$base/_crossbar/usage?by=host" | grep -q '"cached_tokens":[1-9]' || fail "usage has no cached tokens" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_requests_total{route="hermes-x",host="beta",status="403"}' && fail "403s are refused before a lease and must not be counted as requests" +curl -s "$base/_crossbar/metrics" | grep -q 'crossbar_host_healthy{host="beta"} 1' || fail "metrics missing beta health" +echo "smoke: ok (stream spread $((lastms - firstms)) ms)"