diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..9ef7c4b --- /dev/null +++ b/.gitignore @@ -0,0 +1,2 @@ +.state/ +bin/ diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..4ba75ba --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,72 @@ +# AGENTS.md + +crossbar is an affinity router for the fleet's `llama-server` instances, written in Go. You are +implementing it one task at a time. + +## How you work + +1. The owner gives you one task file, `docs/plans//NN-name.md`. Read it fully. Do that task + and nothing else. Do not start the next task. +2. Do the steps in order. Where a step shows a command and its expected output, run it and compare. +3. Do not read `PLAN.md` or other task files unless the task tells you to. The task file quotes + what you need. +4. If something in the task is impossible, contradictory, or fails twice in the same way, **stop**. + Do not improvise, do not change a test, do not weaken a check. Add your row to + `docs/implementer-log.md` with status `stopped`, say what you tried and what happened, commit + only that file, and tell the owner. +5. Work only from files inside this repository. Never read another checkout (such as + `~/src/crossbar-ref` or `~/src/crossbar-design`), and never search the file system for code. + If you are stuck, stop as in point 4. + +## Files you must never edit + +- `PLAN.md`, `docs/plans/`, `AGENTS.md` +- Anything a task told you to copy from `docs/plans/**/files/`: tests, testdata, `Makefile`, + scripts, `example.toml`, `cmd/fakeupstream`. If a copied test fails, your code is wrong. + +## Code rules + +- Go 1.26 (the version in `go.mod`). Standard library plus the one dependency named in the tasks + (`github.com/BurntSushi/toml`). No other module, ever. +- No source file over 400 lines (`scripts/check-lines.sh`). +- Library code (`internal/...`) never panics on input: no `panic`, no indexing that can go out of + range on data that came from a file, a request or a peer. Check lengths, use `strconv` and + `errors`. `main` packages may exit with a message. +- Errors are values: return them, wrap with `fmt.Errorf("...: %w", err)`, never log-and-continue + in library code. A check that cannot do its job fails; it does not return "ok". +- Never log, print or store a request or response body. Log lines carry names, paths, status + codes and durations only. +- Anything read from the network is bounded (`io.LimitReader`, `http.MaxBytesReader`, timeouts). +- Exported names, field names and JSON/TOML tags are exactly as the task gives them; the tests + compile against them. +- Comments say why, not what. `gofmt` decides layout; run it before the gate. + +## The gate + +`make gate` must print `gate: ok` before a task is done. It runs offline: `gofmt -l`, `go vet`, +`go test -race -count=1 ./...`, and `scripts/check-lines.sh`. Run `gofmt -w` on the files you +touched before the gate. `go test ./internal//` runs one package. + +## Git + +- Work on the branch the task names. One task is one commit. +- Stage only the paths the task lists: `git add ...`. Never `git add -A` or `git add .`. +- Never push, amend, rebase, reset, or switch to another branch. +- Commit message: the subject line the task gives, a blank line, then this trailer: + `Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)` + +## The implementer log + +Before you commit, add one row to the table in `docs/implementer-log.md` and include the file in +the commit. Be honest: the log is how the owner judges the process. + +| Column | What to write | +|---|---| +| Task | The task file name, for example `v0/02-health` | +| Date | Today's date, `YYYY-MM-DD` | +| Status | `done` or `stopped` | +| Gate runs | How many times you ran `make gate` | +| First gate | `pass` or `fail` for the first run | +| Deviations | Anything you did that the task did not say, or `none` | +| Notes | Problems you hit and how you solved them, in one or two sentences | +| Model | Write `?`. The owner fills this in. | diff --git a/docs/implementer-log.md b/docs/implementer-log.md new file mode 100644 index 0000000..90a7119 --- /dev/null +++ b/docs/implementer-log.md @@ -0,0 +1,9 @@ +# Implementer log + +Kept by the implementing model, one row per task. The column meanings are in `AGENTS.md`. The +owner fills in the Model column. The reviewer adds findings under "Reviews" once per plan. + +| Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model | +|---|---|---|---|---|---|---|---| + +## Reviews diff --git a/docs/plans/v0/01-module-gate-config.md b/docs/plans/v0/01-module-gate-config.md new file mode 100644 index 0000000..a08ad55 --- /dev/null +++ b/docs/plans/v0/01-module-gate-config.md @@ -0,0 +1,147 @@ +# v0 task 01: module, gate and the config package + +**Branch:** `v0` (create it: `git switch -c v0`; `git status --short` must be empty first, otherwise stop) +**Commit subject:** `Add the Go module, the gate and the config package` + +## Goal + +Create the Go module and the gate every later task must pass, and write `internal/config`: the +TOML file crossbar starts from, with defaults and validation. At the end `make gate` prints +`gate: ok`. + +## Context + +crossbar is an HTTP reverse proxy in front of several `llama-server` routers ("hosts"). A client's +identity is the first segment of its request path (its "route"); each route has an ordered list +of hosts to try. The config file names the hosts, which models each serves, and the routes. It +must be strict: a misspelt key or a route naming a host that does not exist is an error at start, +not a surprise at 3 a.m. The one dependency is `github.com/BurntSushi/toml` v1.6.0; `go.sum` for it +is given. + +## Files + +- Copy (never edit afterwards): `Makefile`, `scripts/check-lines.sh`, `go.sum`, + `internal/config/config_test.go`, `internal/config/testdata/` (5 files) +- Create: `go.mod`, `internal/config/config.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +Produces, in `internal/config/config.go`, package `config`: + +```go +// Duration is a time.Duration that TOML reads from a string such as "60s" or "30m". +type Duration struct{ time.Duration } +func (d *Duration) UnmarshalText(text []byte) error // time.ParseDuration + +type Model struct { Parallel int `toml:"parallel"` } +type Host struct { + BaseURL string `toml:"base_url"` + Weight float64 `toml:"weight"` + Models map[string]Model `toml:"models"` +} +type Route struct { + Hosts []string `toml:"hosts"` + DefaultModel string `toml:"default_model"` +} +type Config struct { + Listen string `toml:"listen"` + PollInterval Duration `toml:"poll_interval"` + QueueMax int `toml:"queue_max"` + Hosts map[string]Host `toml:"hosts"` + Routes map[string]Route `toml:"routes"` +} + +// Error is a validation error naming the field it is about. +type Error struct { Field, Msg string } +func (e *Error) Error() string // "config: " + Field + ": " + Msg + +const ( + DefaultPollInterval = 60 * time.Second + DefaultQueueMax = 8 + MinPollInterval = time.Second +) + +func Load(path string) (*Config, error) // open, then Parse +func Parse(r io.Reader) (*Config, error) // decode, defaults, validate +func (c *Config) Serves(host, model string) bool // host exists and lists model +func IsError(err error) (*Error, bool) // errors.As on *Error +``` + +Rules the tests check: + +1. **Decoding.** `toml.NewDecoder(r).Decode(&c)`. A TOML syntax error is returned as + `fmt.Errorf("config: %w", err)` — it is *not* an `*Error`. If `md.Undecoded()` is non-empty, + the error is `&Error{Field: , Msg: "unknown key"}`. `Load` + wraps an open failure the same way (`config: …`). +2. **Validation, in this order, first problem wins.** Every problem is an `*Error` with exactly + this `Field`: + - `listen`: required; must be `host:port` (`net.SplitHostPort`); the host part must not be + empty and must not be an unspecified address (`0.0.0.0`, `::`). + - `poll_interval`: `0` (absent) becomes `DefaultPollInterval`; less than `MinPollInterval` is + an error. + - `queue_max`: `0` becomes `DefaultQueueMax`; negative is an error. + - `hosts`: at least one. Then for each host, **in sorted name order**: + - `hosts..base_url`: `url.Parse` must succeed, scheme `http` or `https`, non-empty + host, no query, no fragment. A trailing `/` is trimmed (`strings.TrimRight(url, "/")`) + and the trimmed value stored. + - `hosts..weight`: `0` becomes `1`; negative is an error. + - `hosts..models`: at least one. For each model in sorted order, + `hosts..models..parallel`: `0` becomes `1`; negative is an error. + - `routes`: at least one. For each route in sorted name order: + - `routes.`: the name must match `^[a-z0-9][a-z0-9-]*$`. + - `routes..hosts`: at least one; every entry must be a configured host; no host twice. + - `routes..default_model`: if set, at least one of the route's hosts must list it. +3. Defaults are written back into the returned `Config` (the tests read `Weight == 1`, + `Parallel == 1`, `PollInterval == DefaultPollInterval` after parsing files that omit them). +4. Nothing here panics on a bad file: every map lookup and slice access is on data you checked. + +## Steps + +- [ ] **1. Branch and copy.** + +```sh +git switch -c v0 +mkdir -p scripts internal/config/testdata +cp docs/plans/v0/files/Makefile docs/plans/v0/files/go.sum . +cp docs/plans/v0/files/scripts/check-lines.sh scripts/ +cp docs/plans/v0/files/internal/config/config_test.go internal/config/ +cp docs/plans/v0/files/internal/config/testdata/*.toml internal/config/testdata/ +``` + +Read `Makefile` and `internal/config/config_test.go`. The test names every rule above. + +- [ ] **2. Write `go.mod`** with exactly this content: + +``` +module git.wntrmute.dev/kyle/crossbar + +go 1.26 + +require github.com/BurntSushi/toml v1.6.0 +``` + +Then `go mod download github.com/BurntSushi/toml` (this needs the network once) and +`go mod verify`. Expected: `all modules verified`. + +- [ ] **3. See the test fail.** `go test ./internal/config/`. Expected: it does not compile + (`no Go files` or undefined names). +- [ ] **4. Write `internal/config/config.go`.** Run `gofmt -w internal/config/`. +- [ ] **5. See the test pass.** `go test ./internal/config/`. Expected: `ok`. +- [ ] **6. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **7. Log and commit.** Add your row to `docs/implementer-log.md` (Task `v0/01-module-gate-config`). + +```sh +git add go.mod go.sum Makefile scripts/check-lines.sh internal/config docs/implementer-log.md +git commit +``` + +## Done when + +- `go test ./internal/config/` is `ok`; `make gate` prints `gate: ok`. +- `cmp internal/config/config_test.go docs/plans/v0/files/internal/config/config_test.go` prints nothing. + +## Stop and report if + +- `go mod download` cannot fetch the module (no network): stop, log `stopped`. +- A test expects a `Field` you cannot produce under the rules above: stop; do not edit the test. diff --git a/docs/plans/v0/02-health.md b/docs/plans/v0/02-health.md new file mode 100644 index 0000000..83f0b0d --- /dev/null +++ b/docs/plans/v0/02-health.md @@ -0,0 +1,109 @@ +# v0 task 02: the health poller + +**Branch:** `v0` (run `git switch v0`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the health poller` + +## Goal + +Write `internal/health`: a table that knows, for every host, whether it answered its last poll, +which models it has loaded, and when it last answered. The proxy (task 03) reads this table to +choose a host and tells it when a request to a host fails. + +## Context + +A `llama-server` router answers `GET /health` with 200 when it can serve (503 while loading), +and `GET /v1/models` with `{"object":"list","data":[{"id":"", …}, …]}` listing the models +it has resident. A host that failed must answer **two** polls in a row before it is trusted again +(`RecoveryPolls`), because a router that is loading a model flaps. A host that has never failed is +trusted after its first good poll. **Never call `/slots`**: probing an unloaded model makes the +router load it. + +## Files + +- Copy: `internal/health/health_test.go` +- Create: `internal/health/health.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +Produces, in `internal/health/health.go`, package `health`: + +```go +const RecoveryPolls = 2 // good polls in a row a host needs after a failure +const MaxModelsBody = 1 << 20 // bound on what we read from /v1/models + +type Status struct { + Healthy bool `json:"healthy"` + Loaded []string `json:"loaded"` // sorted, unique model ids from the last good poll + LastOK time.Time `json:"last_ok"` // zero if never + LastErr string `json:"last_err"` // "" after a good poll + Consecutive int `json:"consecutive"` // good polls in a row +} + +type Table struct { /* private: a mutex, name -> base URL, name -> status + "ever failed" flag, client, interval */ } + +func New(hosts map[string]string, interval time.Duration, client *http.Client) *Table +func (t *Table) Run(ctx context.Context) // PollOnce now, then every interval; returns when ctx is done +func (t *Table) PollOnce(ctx context.Context) // polls every host concurrently, returns when all are done +func (t *Table) Get(name string) (Status, bool) // a copy; ok == false for an unknown name +func (t *Table) All() map[string]Status // copies +func (t *Table) MarkDown(name, reason string) // a failure seen by the proxy +``` + +Rules the tests check: + +1. **`New`**: `hosts` maps a host name to its base URL (no trailing slash; the config already + trimmed it). A nil `client` becomes `&http.Client{Timeout: 5 * time.Second}`. Every host + starts with `Healthy false`, `Loaded` an **empty slice, not nil**, `Consecutive 0`. +2. **One poll of one host** is two requests with `ctx`: `GET /health` must return 200, + else the poll fails with `LastErr` starting `health: ` (for example `health: HTTP 503`, or the + client error); then `GET /v1/models` must return 200 and decode as + `{"data":[{"id":"…"}]}`, else the poll fails with `LastErr` starting `models: `. Read at most + `MaxModelsBody` bytes of either body. `Loaded` is the ids, empty ids dropped, duplicates + dropped, sorted. +3. **Success**: `Consecutive++`, `LastOK = time.Now()`, `LastErr = ""`, `Loaded` replaced, and + `Healthy = (never failed) || Consecutive >= RecoveryPolls`. +4. **Failure**, from a poll or from `MarkDown`: remember that the host has failed, `Healthy = + false`, `Consecutive = 0`, `LastErr = reason`. `Loaded` is **left as last seen**. + `MarkDown` uses `reason = "marked down: " + reason`. Unknown names are ignored, never a panic. +5. **A cancelled context is not a failure.** If a request fails and `ctx.Err() != nil`, record + nothing: we were told to stop, that says nothing about the host. +6. **`Run`** polls once immediately, then on a `time.Ticker` every `interval`, and returns when + `ctx` is done (stop the ticker). **`PollOnce`** polls all hosts concurrently (one goroutine + each, a `sync.WaitGroup`) and returns when all have finished. +7. **`Get` and `All` return copies**: changing `Loaded` on what they return must not change the + table (copy the slice). All methods are safe to call from several goroutines: one mutex around + the map, never held while a request is in flight. + +## Steps + +- [ ] **1. Copy.** + +```sh +git switch v0 +cp docs/plans/v0/files/internal/health/health_test.go internal/health/ +``` + +Read the test. `newFake` is the router stand-in; `TestFailureThenRecoveryNeedsTwoPolls` is rule 3 +and 4 in one story. + +- [ ] **2. See the test fail.** `go test ./internal/health/`. Expected: it does not compile. +- [ ] **3. Write `internal/health/health.go`.** `gofmt -w internal/health/`. +- [ ] **4. See the test pass.** `go test -race -count=1 ./internal/health/`. Expected: `ok`. + `TestRunPollsOnStart` polls a fake every 20 ms; if it fails, check rules 5 and 6. +- [ ] **5. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **6. Log and commit.** Row `v0/02-health`. + +```sh +git add internal/health docs/implementer-log.md +git commit +``` + +## Done when + +- `go test -race -count=1 ./internal/health/` is `ok`; `make gate` prints `gate: ok`. +- `cmp internal/health/health_test.go docs/plans/v0/files/internal/health/health_test.go` prints nothing. + +## Stop and report if + +- The race detector reports a race you cannot remove with the one-mutex design above. diff --git a/docs/plans/v0/03-proxy.md b/docs/plans/v0/03-proxy.md new file mode 100644 index 0000000..2932b39 --- /dev/null +++ b/docs/plans/v0/03-proxy.md @@ -0,0 +1,121 @@ +# v0 task 03: the routing reverse proxy + +**Branch:** `v0` (run `git switch v0`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the routing reverse proxy` + +## Goal + +Write `internal/proxy`: an `http.Handler` that takes `/{route}/v1/…`, picks a host from the +route's list using the health table, forwards the request with `httputil.ReverseProxy`, streams +the answer back **as it arrives**, and tells the health table when a host fails. + +## Context + +Clients (OpenCode, Hermes) only know a base URL, so the route is the first path segment: +`http://crossbar:7777/opencode-a/v1/chat/completions`. The upstream must see `/v1/chat/completions` +with the query string kept. Answers are often server-sent-event streams of hundreds of small +chunks over minutes; a proxy that buffers them makes the client look frozen, so **`FlushInterval` +is `-1`** (flush after every write) and anything wrapping the `ResponseWriter` must still +implement `http.Flusher`. Request bodies are looked at once, for a top-level `"model"` field, so +that a request for a model only one host has loaded goes there; the body is then handed to the +upstream unchanged. Bodies are never logged. + +## Files + +- Copy: `internal/proxy/proxy_test.go` +- Create: `internal/proxy/proxy.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +Produces, in `internal/proxy/proxy.go`, package `proxy`: + +```go +const MaxBody = 16 << 20 // largest request body we look at +const HostHeader = "X-Crossbar-Host" // set on every proxied response: the host that answered + +// Health is what the proxy needs from the health table (internal/health satisfies it). +type Health interface { + Get(name string) (health.Status, bool) + MarkDown(name, reason string) +} + +type Handler struct { /* private: *config.Config, Health, *slog.Logger */ } + +func New(cfg *config.Config, h Health, log *slog.Logger) *Handler // nil log -> slog.Default() +func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) + +// SplitRoute takes the first path segment as the route. +// "/a/v1/x" -> ("a", "/v1/x", true) "/a" and "/a/" -> ("a", "/", true) +// "/", "//x", "", "noslash/v1" -> ("", "", false) +func SplitRoute(path string) (route, rest string, ok bool) + +// Choose: the first host in order that is healthy and lists model in Loaded; failing that, the +// first healthy host; ok == false if none. model may be "". +func Choose(hosts []string, model string, h Health) (string, bool) +``` + +Rules the tests check, in the order `ServeHTTP` applies them. Every error answer is JSON +`{"error":""}` with `Content-Type: application/json`: + +1. `SplitRoute(r.URL.Path)` not ok → **400** `missing route`. +2. Route not in `cfg.Routes` → **404** `unknown route`. +3. `rest` must start with `/v1/` or be exactly `/health` or `/props`; else → **404** `not found`. + (`/_crossbar/…` through a route is therefore 404 too.) +4. **Model peek.** For requests other than GET/HEAD with a body: read up to `MaxBody + 1` bytes + (`io.LimitReader`); more than `MaxBody` → **413** `body too large`. Put the bytes back + (`r.Body = io.NopCloser(bytes.NewReader(body))`, `r.ContentLength = len(body)`). Then try to + decode `{"model": "…"}`; a body that is not JSON, or has no model, simply gives `""` — that is + not an error. If the model is `""`, use the route's `DefaultModel`. +5. `Choose(route.Hosts, model, h)` not ok → **503** `no healthy host`. Nothing is marked down by + rules 1–5. +6. **Forward** with a `httputil.ReverseProxy`: + - `Rewrite`: `pr.SetURL(target)` where `target` is the host's `BaseURL` parsed once; + `pr.Out.URL.Path = target.Path + rest`; `pr.Out.URL.RawPath = ""`; `pr.Out.Host = + target.Host`; `pr.SetXForwarded()`. The query string is kept (the tests check + `/v1/models?x=1` arrives as `/v1/models?x=1`). + - `FlushInterval: -1`. + - `ModifyResponse`: set `HostHeader` to the host's name. + - `ErrorHandler`: if `errors.Is(err, context.Canceled)`, do nothing (the client left); + otherwise `h.MarkDown(name, err.Error())` and answer **502** + `{"error":"upstream failed","host":""}`. +7. **One log line per proxied request**, after it finishes, through the logger: + `p.log.Info("request", "route", …, "host", …, "method", …, "path", rest, "status", …, "ms", …)`. + To know the status, wrap the `ResponseWriter` in a small recorder that implements + `WriteHeader` **and `Flush`** (forwarding to the underlying `http.Flusher`). Without `Flush` + the reverse proxy cannot stream and `TestStreamingIsNotBuffered` fails. +8. Never log, print or keep a request or response body. Never panic on a request. + +## Steps + +- [ ] **1. Copy.** + +```sh +git switch v0 +cp docs/plans/v0/files/internal/proxy/proxy_test.go internal/proxy/ +``` + +Read the test. `fakeHealth` stands in for the table; `newUpstream` records what arrived. +`TestStreamingIsNotBuffered` is rule 6/7: the upstream sends one chunk and then *waits until the +test has read it*; a buffering proxy hangs there. + +- [ ] **2. See the test fail.** `go test ./internal/proxy/`. Expected: it does not compile. +- [ ] **3. Write `internal/proxy/proxy.go`.** `gofmt -w internal/proxy/`. +- [ ] **4. See the test pass.** `go test -race -count=1 ./internal/proxy/`. Expected: `ok`. +- [ ] **5. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **6. Log and commit.** Row `v0/03-proxy`. + +```sh +git add internal/proxy docs/implementer-log.md +git commit +``` + +## Done when + +- `go test -race -count=1 ./internal/proxy/` is `ok`; `make gate` prints `gate: ok`. +- `cmp internal/proxy/proxy_test.go docs/plans/v0/files/internal/proxy/proxy_test.go` prints nothing. + +## Stop and report if + +- `TestStreamingIsNotBuffered` still fails with `FlushInterval: -1` and a recorder that + implements `Flush`: stop and describe exactly what you wrote. diff --git a/docs/plans/v0/04-admin-main.md b/docs/plans/v0/04-admin-main.md new file mode 100644 index 0000000..69956c9 --- /dev/null +++ b/docs/plans/v0/04-admin-main.md @@ -0,0 +1,127 @@ +# v0 task 04: admin endpoints, the `crossbar` binary, the fake upstream + +**Branch:** `v0` (run `git switch v0`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the admin endpoints, the crossbar binary and the fake upstream` + +## Goal + +Make crossbar runnable: `internal/admin` shows the health table and the routes as JSON, +`cmd/crossbar` wires config, health, proxy and admin into one HTTP server with graceful shutdown, +and the given `cmd/fakeupstream` stands in for a router so the whole thing can be exercised +without a real model (task 05 does that). + +## Context + +Operators read `/_crossbar/hosts` to see why a request went where it went, so its shape is +fixed: every host, `healthy`, `loaded` (always a JSON array, never `null`), `last_ok` as RFC 3339 +in UTC or `""`, `last_err`. The admin handler is mounted at `/_crossbar/` on the same listener as +the proxy; the proxy already refuses `/{route}/_crossbar/…` (task 03). + +## Files + +- Copy (never edit): `cmd/fakeupstream/main.go`, `example.toml`, `internal/admin/admin_test.go` +- Create: `internal/admin/admin.go`, `cmd/crossbar/main.go` +- Modify: `docs/implementer-log.md` + +## Interfaces + +Produces, in `internal/admin/admin.go`, package `admin`: + +```go +// Hosts is what the admin handler needs from the health table. +type Hosts interface { All() map[string]health.Status } + +type HostView struct { + Healthy bool `json:"healthy"` + Loaded []string `json:"loaded"` // never null: an empty slice when nothing is loaded + LastOK string `json:"last_ok"` // time.RFC3339 in UTC, or "" if never + LastErr string `json:"last_err"` +} +type RouteView struct { + Hosts []string `json:"hosts"` + DefaultModel string `json:"default_model"` +} + +// Handler serves GET /_crossbar/hosts and GET /_crossbar/routes. +func Handler(cfg *config.Config, h Hosts) http.Handler +``` + +Rules the tests check: + +1. `GET /_crossbar/hosts` → 200, `Content-Type: application/json`, a JSON object mapping host + name to `HostView` built from `h.All()`. +2. `GET /_crossbar/routes` → 200, JSON object mapping route name to `RouteView` (copy the hosts + slice; do not hand out the config's). +3. Any other method on those two paths → **405** with an `Allow: GET` header and a JSON + `{"error":"method not allowed"}` body. Any other path under the handler → **404** JSON + `{"error":"not found"}`. (An `http.ServeMux` with the two exact paths plus a `/` fallback does + this.) + +Produces `cmd/crossbar/main.go`, package `main`: + +- Flag `-config` (default `crossbar.toml`). `config.Load`; on error print `crossbar: ` to + stderr and exit 1. +- Logger: `slog.New(slog.NewTextHandler(os.Stderr, nil))`. +- `health.New( BaseURL from cfg.Hosts>, cfg.PollInterval.Duration, nil)`; run it with + `go table.Run(ctx)` where `ctx` comes from + `signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)`. +- `http.ServeMux`: `mux.Handle("/_crossbar/", admin.Handler(cfg, table))`, + `mux.Handle("/", proxy.New(cfg, table, log))`. +- `&http.Server{Addr: cfg.Listen, Handler: mux, ReadHeaderTimeout: 10 * time.Second}`. Log + `"listening"` with `addr` before `ListenAndServe`. On signal, log `"shutting down"` and + `srv.Shutdown` with a 10 s timeout; exit 0. A `ListenAndServe` error other than + `http.ErrServerClosed` exits 1 with the message. + +`cmd/fakeupstream/main.go` is given; read its package comment so you know what it does, and do +not change it. + +## Steps + +- [ ] **1. Copy.** + +```sh +git switch v0 +mkdir -p cmd/fakeupstream cmd/crossbar internal/admin +cp docs/plans/v0/files/cmd/fakeupstream/main.go cmd/fakeupstream/ +cp docs/plans/v0/files/example.toml . +cp docs/plans/v0/files/internal/admin/admin_test.go internal/admin/ +``` + +- [ ] **2. See the test fail.** `go test ./internal/admin/`. Expected: it does not compile. +- [ ] **3. Write `internal/admin/admin.go` and `cmd/crossbar/main.go`.** `gofmt -w .` +- [ ] **4. See the test pass.** `go test -race -count=1 ./internal/admin/`. Expected: `ok`. +- [ ] **5. Build and run for three seconds.** + +```sh +make build +timeout --signal=TERM 3 bin/crossbar -config example.toml; echo "exit=$?" +``` + +Expected on stderr: a line containing `listening` and `addr=127.0.0.1:17777`, then +`shutting down`; then `exit=0`. (`example.toml` names two upstreams that are not running; the +health table simply records them unhealthy — that is fine here.) + +```sh +bin/crossbar -config /nonexistent.toml; echo "exit=$?" +``` + +Expected: `crossbar: config: open /nonexistent.toml: no such file or directory` and `exit=1`. + +- [ ] **6. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **7. Log and commit.** Row `v0/04-admin-main`. `bin/` is build output: do not add it. + +```sh +git add internal/admin cmd/crossbar cmd/fakeupstream example.toml docs/implementer-log.md +git commit +``` + +## Done when + +- The admin test passes, the three-second run exits 0 with both log lines, the missing-config + run exits 1, `make gate` prints `gate: ok`. +- `cmp cmd/fakeupstream/main.go docs/plans/v0/files/cmd/fakeupstream/main.go` and the same for + `example.toml` and `internal/admin/admin_test.go` print nothing. + +## Stop and report if + +- `bin/crossbar` does not exit 0 on SIGTERM within the timeout. diff --git a/docs/plans/v0/05-smoke-readme-deploy.md b/docs/plans/v0/05-smoke-readme-deploy.md new file mode 100644 index 0000000..896e963 --- /dev/null +++ b/docs/plans/v0/05-smoke-readme-deploy.md @@ -0,0 +1,118 @@ +# v0 task 05: the smoke run, README and systemd unit + +**Branch:** `v0` (run `git switch v0`; `git status --short` must be empty, otherwise stop) +**Commit subject:** `Add the smoke run, README and systemd unit` + +## Goal + +Prove the whole binary works over real HTTP against two fake upstreams — routing, failover, +recovery, streaming — with the given `tools/smoke.sh`, then document how to build, configure, +run and point clients at crossbar. + +## Context + +`tools/smoke.sh` starts two `fakeupstream`s on `127.0.0.1:18081/18082` and `bin/crossbar` on +`example.toml` (`poll_interval = "1s"`), then checks with `curl`: `opencode-a` goes to `alpha` +and `hermes-x` to `beta`; an unknown route is 404; after `alpha` is taken down (its `-down-file` +appears) requests fail over to `beta` and `/_crossbar/hosts` shows `alpha` unhealthy; after two +good polls `alpha` is back; a streamed completion of five chunks 200 ms apart reaches the client +spread over at least 600 ms, not in one burst; and crossbar wrote a request log line. It prints +`smoke: ok (stream spread N ms)` or fails with the crossbar log. + +## Files + +- Copy (never edit): `tools/smoke.sh` +- Create: `README.md`, `deploy/crossbar.service` +- Modify: `docs/implementer-log.md` + +## Steps + +- [ ] **1. Copy and run the smoke test.** + +```sh +git switch v0 +mkdir -p tools deploy +cp docs/plans/v0/files/tools/smoke.sh tools/ +make smoke +``` + +Expected last line: `smoke: ok (stream spread N ms)` with N ≥ 600. If it fails, the message says +which check failed and prints crossbar's log; the fault is in code from tasks 02–04 or in this +machine's `curl`. Fix code only if a rule from an earlier task was broken; otherwise stop and +report. + +- [ ] **2. Write `deploy/crossbar.service`** with exactly this content: + +```ini +[Unit] +Description=crossbar affinity router for llama-server +After=network-online.target +Wants=network-online.target + +[Service] +ExecStart=/usr/local/bin/crossbar -config /etc/crossbar/crossbar.toml +Restart=on-failure +RestartSec=2s +DynamicUser=yes +StateDirectory=crossbar +NoNewPrivileges=yes +ProtectSystem=strict +ProtectHome=yes +PrivateTmp=yes + +[Install] +WantedBy=multi-user.target +``` + +- [ ] **3. Write `README.md`** with these sections, in this order, in plain prose (no + hostnames, addresses or tokens other than the placeholders shown): + 1. `# crossbar` — two sentences: an affinity router in front of several `llama-server` + routers; a client's identity is the first path segment of its base URL; v0 routes each + to the first healthy host on that route's list and streams answers unbuffered. + 2. `## Build` — `make build` (binaries in `bin/`), `make gate`, `make smoke`. + 3. `## Configure` — paste `example.toml` in a fenced block and explain each key in a table: + `listen` (a tailnet address, never `0.0.0.0`), `poll_interval` (60s in production), + `queue_max` (reserved for v1), `hosts..base_url|weight|models`, `routes..hosts` + (preference order) and `default_model`. + 4. `## Run` — copy `bin/crossbar` to `/usr/local/bin/`, the config to + `/etc/crossbar/crossbar.toml`, the unit to `/etc/systemd/system/`, then + `systemctl enable --now crossbar`. + 5. `## Point clients at it` — these two snippets verbatim: + + OpenCode, one provider for every project; each instance is launched as + `CROSSBAR_ROUTE="$(basename "$PWD")-$$" opencode`: + ```jsonc + "provider": { "crossbar": { "npm": "@ai-sdk/openai-compatible", + "options": { "baseURL": "http://crossbar.:7777/{env:CROSSBAR_ROUTE}/v1" }, + "models": { "ornith-1.5-35b-a3b": {} } } } + ``` + Hermes, in `config.yaml`: + ```yaml + custom_providers: + - name: crossbar + base_url: http://crossbar.:7777/hermes-/v1 + models: { ornith-1.5-35b-a3b: {} } + ``` + Note under them: the route name in the URL must exist in `[routes]`; unknown routes are 404. + 6. `## Inspect` — `GET /_crossbar/hosts` and `GET /_crossbar/routes`, one example response each + (take them from the smoke run). + 7. `## What v0 does not do` — leases and stickiness, SQLite, `/slots`, queueing, wake-on-LAN: + see `PLAN.md`. + +- [ ] **4. Run the gate.** `make gate`. Expected last line: `gate: ok`. +- [ ] **5. Log and commit.** Row `v0/05-smoke-readme-deploy`; put the smoke line (`smoke: ok + (stream spread N ms)`) in the Notes column. + +```sh +git add tools/smoke.sh README.md deploy/crossbar.service docs/implementer-log.md +git commit +``` + +## Done when + +- `make smoke` prints `smoke: ok (…)`; `make gate` prints `gate: ok`; `README.md` has the seven + sections; `cmp tools/smoke.sh docs/plans/v0/files/tools/smoke.sh` prints nothing. + +## Stop and report if + +- `make smoke` fails twice in the same way. diff --git a/docs/plans/v0/README.md b/docs/plans/v0/README.md new file mode 100644 index 0000000..72e039a --- /dev/null +++ b/docs/plans/v0/README.md @@ -0,0 +1,62 @@ +# v0 implementation plan: static routing, health, streaming proxy + +> **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:** a `crossbar` binary that reads a TOML file, polls each host's `/health` and +`/v1/models`, and forwards `/{route}/v1/*` to the first healthy host on that route's preference +list, streaming responses chunk by chunk, marking a host down when a proxied request fails, and +showing the health table at `/_crossbar/hosts`. This is `PLAN.md` §10 v0: no leases, no SQLite, +no `/slots`, no queueing. + +**Architecture:** four small packages under `internal/` — `config` (TOML + validation), `health` +(poller + table), `proxy` (route → host → `httputil.ReverseProxy`), `admin` (read-only JSON) — +and two binaries under `cmd/`: `crossbar` and the given `fakeupstream`. Behaviour is pinned by +the tests in `files/`, which were run against a private reference implementation at each task's +end state; the reference is not in this repository and the implementer must not look for it. + +**Tech stack:** Go 1.26, standard library, `github.com/BurntSushi/toml` v1.6.0. Nothing else. + +## Global constraints + +- Everything in `AGENTS.md`. +- Branch `v0`. One task, one fresh OpenCode session, one commit. `gofmt -w` before the gate. +- No source file over 400 lines. Library code never panics on input, never logs a body. +- The given files are copied and never edited. If a copied test fails, the code is wrong. + +## Tasks + +| # | File | Delivers | Tests that define it | +|---|---|---|---| +| 01 | `01-module-gate-config.md` | `go.mod`, the gate, `internal/config` | `internal/config/config_test.go` + `testdata/` | +| 02 | `02-health.md` | `internal/health`: poller and table | `internal/health/health_test.go` | +| 03 | `03-proxy.md` | `internal/proxy`: routing reverse proxy | `internal/proxy/proxy_test.go` | +| 04 | `04-admin-main.md` | `internal/admin`, `cmd/crossbar`, the given `cmd/fakeupstream` | `internal/admin/admin_test.go`, a start/stop check | +| 05 | `05-smoke-readme-deploy.md` | `tools/smoke.sh` run, `README.md`, `deploy/crossbar.service` | `make smoke` | + +At the end: `make gate` prints `gate: ok`, `make smoke` prints `smoke: ok (stream spread N ms)`. + +## For the owner: running a task + +From a clean checkout on `master`: + +```sh +tools/run-plan.sh docs/plans/v0 # all tasks, each in a fresh `opencode run` session +tools/run-plan.sh docs/plans/v0 03 # from task 03 +``` + +The default model is `llama.cpp/ornith-1.5-35b-a3b` (override with `CROSSBAR_MODEL`). The driver +stops at the first task that does not end with a commit, a clean tree and a `done` row in +`docs/implementer-log.md`. Keep the OpenCode TUI closed while it runs. + +## For the reviewer: after task 05 + +1. `git log --oneline master..v0`: five commits with the `Implemented-By` trailer. +2. Copied files are unchanged: + `for f in $(cd docs/plans/v0/files && find . -type f); do cmp "docs/plans/v0/files/$f" "$f"; done` +3. `git diff master..v0 --stat -- PLAN.md AGENTS.md docs/plans` is empty. +4. `make gate` and `make smoke` on straylight. +5. Read every source file against its task. Probe from outside with inputs the tests do not + contain: a route with query strings and encoded characters, a host that hangs after headers, a + 3 MB request body, `/_crossbar/hosts` while a poll is in flight, SIGTERM during a stream. +6. Write findings under "Reviews" in `docs/implementer-log.md`. diff --git a/docs/plans/v0/files/Makefile b/docs/plans/v0/files/Makefile new file mode 100644 index 0000000..1e69fb0 --- /dev/null +++ b/docs/plans/v0/files/Makefile @@ -0,0 +1,17 @@ +# crossbar gate. `make gate` must pass before any task is called done. It needs no network. + +.PHONY: gate build smoke + +gate: + @test -z "$$(gofmt -l . 2>&1)" || { echo "gofmt: these files need formatting:"; gofmt -l .; exit 1; } + go vet ./... + go test -race -count=1 ./... + sh scripts/check-lines.sh + @echo "gate: ok" + +build: + go build -o bin/ ./cmd/... + +# Runs the whole thing against two fake upstreams. Task 05 brings the script. +smoke: build + sh tools/smoke.sh diff --git a/docs/plans/v0/files/cmd/fakeupstream/main.go b/docs/plans/v0/files/cmd/fakeupstream/main.go new file mode 100644 index 0000000..dbb7801 --- /dev/null +++ b/docs/plans/v0/files/cmd/fakeupstream/main.go @@ -0,0 +1,101 @@ +// 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 +// +// /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, one JSON answer otherwise. 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") + flag.Parse() + + ids := strings.Split(*models, ",") + mux := http.NewServeMux() + stamp := func(w http.ResponseWriter) { w.Header().Set("X-Upstream", *name) } + + 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) + 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": map[string]int{"prompt_tokens": 3, "completion_tokens": 3, "total_tokens": 6}, + }) + return + } + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.WriteHeader(http.StatusOK) + fl, _ := w.(http.Flusher) + 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) + if fl != nil { + fl.Flush() + } + time.Sleep(200 * time.Millisecond) + } + 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", *name, *listen, ids) + 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/v0/files/example.toml b/docs/plans/v0/files/example.toml new file mode 100644 index 0000000..7b105cc --- /dev/null +++ b/docs/plans/v0/files/example.toml @@ -0,0 +1,22 @@ +# crossbar example configuration. Replace and the addresses with your own. +listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production +poll_interval = "1s" # 60s in production; 1s makes the smoke run quick +queue_max = 8 + +[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 = 4 }, "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 = 4 } } + +# v0: a route is a preference list; the first healthy host that has the model wins. +[routes.opencode-a] +hosts = ["alpha", "beta"] +default_model = "ornith-1.5-35b-a3b" + +[routes.hermes-x] +hosts = ["beta", "alpha"] diff --git a/docs/plans/v0/files/go.sum b/docs/plans/v0/files/go.sum new file mode 100644 index 0000000..f74b269 --- /dev/null +++ b/docs/plans/v0/files/go.sum @@ -0,0 +1,2 @@ +github.com/BurntSushi/toml v1.6.0 h1:dRaEfpa2VI55EwlIW72hMRHdWouJeRF7TPYhI+AUQjk= +github.com/BurntSushi/toml v1.6.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= diff --git a/docs/plans/v0/files/internal/admin/admin_test.go b/docs/plans/v0/files/internal/admin/admin_test.go new file mode 100644 index 0000000..80a0dfb --- /dev/null +++ b/docs/plans/v0/files/internal/admin/admin_test.go @@ -0,0 +1,99 @@ +package admin_test + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/admin" + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/health" +) + +type fakeHosts map[string]health.Status + +func (f fakeHosts) All() map[string]health.Status { return f } + +func testConfig(t *testing.T) *config.Config { + c, err := config.Parse(strings.NewReader(` +listen = "127.0.0.1:1" +[hosts.alpha] +base_url = "http://alpha:1" +models = { "m" = { } } +[hosts.beta] +base_url = "http://beta:1" +models = { "m" = { } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "m" +`)) + if err != nil { + t.Fatal(err) + } + return c +} + +func TestHosts(t *testing.T) { + when := time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC) + h := admin.Handler(testConfig(t), fakeHosts{ + "alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: when}, + "beta": {Healthy: false, LastErr: "HTTP 503"}, + }) + rec := httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/hosts", nil)) + if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Fatalf("status %d, content-type %q", rec.Code, rec.Header().Get("Content-Type")) + } + var out map[string]admin.HostView + if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { + t.Fatal(err) + } + if a := out["alpha"]; !a.Healthy || len(a.Loaded) != 1 || a.LastOK != "2026-09-25T08:00:00Z" || a.LastErr != "" { + t.Errorf("alpha = %+v", a) + } + if b := out["beta"]; b.Healthy || b.LastOK != "" || b.LastErr != "HTTP 503" || b.Loaded == nil { + t.Errorf("beta = %+v (loaded must be [] not null)", b) + } + if !strings.Contains(rec.Body.String(), `"loaded":[]`) { + t.Errorf("beta.loaded must encode as []: %s", rec.Body.String()) + } +} + +func TestRoutes(t *testing.T) { + h := admin.Handler(testConfig(t), fakeHosts{}) + rec := httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/_crossbar/routes", nil)) + var out map[string]admin.RouteView + if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil { + t.Fatalf("%v: %s", err, rec.Body.String()) + } + r := out["r"] + if len(r.Hosts) != 2 || r.Hosts[0] != "alpha" || r.DefaultModel != "m" { + t.Errorf("routes = %+v", out) + } +} + +func TestMethodsAndUnknown(t *testing.T) { + h := admin.Handler(testConfig(t), fakeHosts{}) + for _, tc := range []struct { + method, path string + want int + }{ + {http.MethodPost, "/_crossbar/hosts", 405}, + {http.MethodDelete, "/_crossbar/routes", 405}, + {http.MethodGet, "/_crossbar/nope", 404}, + {http.MethodGet, "/_crossbar/", 404}, + } { + rec := httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest(tc.method, tc.path, nil)) + if rec.Code != tc.want { + t.Errorf("%s %s = %d, want %d", tc.method, tc.path, rec.Code, tc.want) + } + if !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Errorf("%s %s: errors are JSON too", tc.method, tc.path) + } + } +} diff --git a/docs/plans/v0/files/internal/config/config_test.go b/docs/plans/v0/files/internal/config/config_test.go new file mode 100644 index 0000000..9069c36 --- /dev/null +++ b/docs/plans/v0/files/internal/config/config_test.go @@ -0,0 +1,191 @@ +package config_test + +import ( + "fmt" + "path/filepath" + "strings" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" +) + +func TestGoodFile(t *testing.T) { + c, err := config.Load(filepath.Join("testdata", "good.toml")) + if err != nil { + t.Fatalf("Load: %v", err) + } + if c.Listen != "100.64.0.9:7777" { + t.Errorf("Listen = %q", c.Listen) + } + if c.PollInterval.Duration != 5*time.Second { + t.Errorf("PollInterval = %v", c.PollInterval.Duration) + } + if c.QueueMax != 4 { + t.Errorf("QueueMax = %d", c.QueueMax) + } + alpha := c.Hosts["alpha"] + if alpha.BaseURL != "http://alpha.example:11434" { + t.Errorf("trailing slash not stripped: %q", alpha.BaseURL) + } + if alpha.Weight != 2 { + t.Errorf("alpha.Weight = %v", alpha.Weight) + } + if alpha.Models["ornith-1.5-35b-a3b"].Parallel != 4 || alpha.Models["small-9b"].Parallel != 6 { + t.Errorf("alpha.Models = %+v", alpha.Models) + } + beta := c.Hosts["beta"] + if beta.Weight != 1 { + t.Errorf("beta.Weight default = %v, want 1", beta.Weight) + } + if beta.Models["ornith-1.5-35b-a3b"].Parallel != 1 { + t.Errorf("beta parallel default = %d, want 1", beta.Models["ornith-1.5-35b-a3b"].Parallel) + } + r := c.Routes["opencode-a"] + if len(r.Hosts) != 2 || r.Hosts[0] != "alpha" || r.Hosts[1] != "beta" { + t.Errorf("route hosts = %v", r.Hosts) + } + if r.DefaultModel != "ornith-1.5-35b-a3b" { + t.Errorf("DefaultModel = %q", r.DefaultModel) + } + if c.Routes["hermes-x"].DefaultModel != "" { + t.Errorf("hermes-x DefaultModel should be empty") + } + if !c.Serves("alpha", "small-9b") || c.Serves("beta", "small-9b") || c.Serves("nope", "m") { + t.Errorf("Serves is wrong") + } +} + +func TestDefaults(t *testing.T) { + c, err := config.Parse(strings.NewReader(` +listen = "127.0.0.1:1" +[hosts.a] +base_url = "http://a:1" +models = { "m" = { } } +[routes.r] +hosts = ["a"] +`)) + if err != nil { + t.Fatalf("Parse: %v", err) + } + if c.PollInterval.Duration != config.DefaultPollInterval { + t.Errorf("PollInterval default = %v", c.PollInterval.Duration) + } + if c.QueueMax != config.DefaultQueueMax { + t.Errorf("QueueMax default = %d", c.QueueMax) + } +} + +func TestBadFiles(t *testing.T) { + cases := []struct{ file, field string }{ + {"bad-listen.toml", "listen"}, + {"bad-unknown-host.toml", "routes.r.hosts"}, + {"bad-default-model.toml", "routes.r.default_model"}, + {"bad-unknown-key.toml", "lease_idle"}, + } + for _, tc := range cases { + t.Run(tc.file, func(t *testing.T) { + _, err := config.Load(filepath.Join("testdata", tc.file)) + if err == nil { + t.Fatalf("want error") + } + e, ok := config.IsError(err) + if !ok { + t.Fatalf("want *config.Error, got %T: %v", err, err) + } + if e.Field != tc.field { + t.Errorf("Field = %q, want %q (%v)", e.Field, tc.field, err) + } + if !strings.HasPrefix(err.Error(), "config: "+tc.field+": ") { + t.Errorf("Error() = %q", err.Error()) + } + }) + } +} + +func TestBadValues(t *testing.T) { + base := ` +listen = %q +poll_interval = %q +[hosts.a] +base_url = %q +weight = %v +models = { "m" = { parallel = %d } } +[routes.%s] +hosts = ["a"] +` + cases := []struct { + name string + listen, poll, url, route string + weight float64 + parallel int + field string + }{ + {"empty listen", "", "5s", "http://a:1", "r", 1, 1, "listen"}, + {"no port", "127.0.0.1", "5s", "http://a:1", "r", 1, 1, "listen"}, + {"v6 any", "[::]:7", "5s", "http://a:1", "r", 1, 1, "listen"}, + {"poll too short", "127.0.0.1:7", "500ms", "http://a:1", "r", 1, 1, "poll_interval"}, + {"ftp url", "127.0.0.1:7", "5s", "ftp://a:1", "r", 1, 1, "hosts.a.base_url"}, + {"no host", "127.0.0.1:7", "5s", "http://", "r", 1, 1, "hosts.a.base_url"}, + {"query", "127.0.0.1:7", "5s", "http://a:1/v1?x=1", "r", 1, 1, "hosts.a.base_url"}, + {"negative weight", "127.0.0.1:7", "5s", "http://a:1", "r", -1, 1, "hosts.a.weight"}, + {"negative parallel", "127.0.0.1:7", "5s", "http://a:1", "r", 1, -2, "hosts.a.models.m.parallel"}, + {"route name", "127.0.0.1:7", "5s", "http://a:1", "Bad_Name", 1, 1, "routes.Bad_Name"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + text := fmt.Sprintf(base, tc.listen, tc.poll, tc.url, tc.weight, tc.parallel, tc.route) + _, err := config.Parse(strings.NewReader(text)) + if err == nil { + t.Fatalf("want error for %s", tc.name) + } + e, ok := config.IsError(err) + if !ok { + t.Fatalf("want *config.Error, got %T: %v", err, err) + } + if e.Field != tc.field { + t.Errorf("Field = %q, want %q (%v)", e.Field, tc.field, err) + } + }) + } +} + +func TestMissingSections(t *testing.T) { + for _, tc := range []struct{ name, text, field string }{ + {"no hosts", "listen = \"127.0.0.1:7\"\n[routes.r]\nhosts = [\"a\"]\n", "hosts"}, + {"no routes", "listen = \"127.0.0.1:7\"\n[hosts.a]\nbase_url = \"http://a:1\"\nmodels = { \"m\" = { } }\n", "routes"}, + {"host without models", "listen = \"127.0.0.1:7\"\n[hosts.a]\nbase_url = \"http://a:1\"\n[routes.r]\nhosts = [\"a\"]\n", "hosts.a.models"}, + {"route without hosts", "listen = \"127.0.0.1:7\"\n[hosts.a]\nbase_url = \"http://a:1\"\nmodels = { \"m\" = { } }\n[routes.r]\n", "routes.r.hosts"}, + {"host twice", "listen = \"127.0.0.1:7\"\n[hosts.a]\nbase_url = \"http://a:1\"\nmodels = { \"m\" = { } }\n[routes.r]\nhosts = [\"a\", \"a\"]\n", "routes.r.hosts"}, + } { + t.Run(tc.name, func(t *testing.T) { + _, err := config.Parse(strings.NewReader(tc.text)) + e, ok := config.IsError(err) + if !ok { + t.Fatalf("want *config.Error, got %v", err) + } + if e.Field != tc.field { + t.Errorf("Field = %q, want %q", e.Field, tc.field) + } + }) + } +} + +func TestNotTOML(t *testing.T) { + _, err := config.Parse(strings.NewReader("listen = [unterminated")) + if err == nil { + t.Fatal("want error") + } + if _, ok := config.IsError(err); ok { + t.Errorf("a syntax error is not a validation Error") + } + if !strings.HasPrefix(err.Error(), "config: ") { + t.Errorf("Error() = %q", err.Error()) + } +} + +func TestMissingFile(t *testing.T) { + if _, err := config.Load(filepath.Join("testdata", "does-not-exist.toml")); err == nil { + t.Fatal("want error") + } +} diff --git a/docs/plans/v0/files/internal/config/testdata/bad-default-model.toml b/docs/plans/v0/files/internal/config/testdata/bad-default-model.toml new file mode 100644 index 0000000..4e1cc32 --- /dev/null +++ b/docs/plans/v0/files/internal/config/testdata/bad-default-model.toml @@ -0,0 +1,9 @@ +listen = "127.0.0.1:7777" + +[hosts.alpha] +base_url = "http://alpha.example:11434" +models = { "m" = { } } + +[routes.r] +hosts = ["alpha"] +default_model = "not-served" diff --git a/docs/plans/v0/files/internal/config/testdata/bad-listen.toml b/docs/plans/v0/files/internal/config/testdata/bad-listen.toml new file mode 100644 index 0000000..15e0e6d --- /dev/null +++ b/docs/plans/v0/files/internal/config/testdata/bad-listen.toml @@ -0,0 +1,8 @@ +listen = "0.0.0.0:7777" + +[hosts.alpha] +base_url = "http://alpha.example:11434" +models = { "m" = { } } + +[routes.r] +hosts = ["alpha"] diff --git a/docs/plans/v0/files/internal/config/testdata/bad-unknown-host.toml b/docs/plans/v0/files/internal/config/testdata/bad-unknown-host.toml new file mode 100644 index 0000000..fd5aa92 --- /dev/null +++ b/docs/plans/v0/files/internal/config/testdata/bad-unknown-host.toml @@ -0,0 +1,8 @@ +listen = "127.0.0.1:7777" + +[hosts.alpha] +base_url = "http://alpha.example:11434" +models = { "m" = { } } + +[routes.r] +hosts = ["alpha", "gamma"] diff --git a/docs/plans/v0/files/internal/config/testdata/bad-unknown-key.toml b/docs/plans/v0/files/internal/config/testdata/bad-unknown-key.toml new file mode 100644 index 0000000..e3e94ec --- /dev/null +++ b/docs/plans/v0/files/internal/config/testdata/bad-unknown-key.toml @@ -0,0 +1,9 @@ +listen = "127.0.0.1:7777" +lease_idle = "30m" + +[hosts.alpha] +base_url = "http://alpha.example:11434" +models = { "m" = { } } + +[routes.r] +hosts = ["alpha"] diff --git a/docs/plans/v0/files/internal/config/testdata/good.toml b/docs/plans/v0/files/internal/config/testdata/good.toml new file mode 100644 index 0000000..92ad98e --- /dev/null +++ b/docs/plans/v0/files/internal/config/testdata/good.toml @@ -0,0 +1,19 @@ +listen = "100.64.0.9:7777" +poll_interval = "5s" +queue_max = 4 + +[hosts.alpha] +base_url = "http://alpha.example:11434/" +weight = 2.0 +models = { "ornith-1.5-35b-a3b" = { parallel = 4 }, "small-9b" = { parallel = 6 } } + +[hosts.beta] +base_url = "https://beta.example:8081" +models = { "ornith-1.5-35b-a3b" = { } } + +[routes.opencode-a] +hosts = ["alpha", "beta"] +default_model = "ornith-1.5-35b-a3b" + +[routes.hermes-x] +hosts = ["beta"] diff --git a/docs/plans/v0/files/internal/health/health_test.go b/docs/plans/v0/files/internal/health/health_test.go new file mode 100644 index 0000000..3f3d53f --- /dev/null +++ b/docs/plans/v0/files/internal/health/health_test.go @@ -0,0 +1,154 @@ +package health_test + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/health" +) + +// fake is a llama-server stand-in whose /health can be flipped and whose model list is fixed. +type fake struct { + srv *httptest.Server + down atomic.Bool + models []string + hits atomic.Int32 +} + +func newFake(t *testing.T, models ...string) *fake { + f := &fake{models: models} + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { + f.hits.Add(1) + if f.down.Load() { + http.Error(w, "loading", http.StatusServiceUnavailable) + return + } + _, _ = w.Write([]byte(`{"status":"ok"}`)) + }) + mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { + type m struct { + ID string `json:"id"` + } + var data []m + for _, id := range f.models { + data = append(data, m{ID: id}) + } + _ = json.NewEncoder(w).Encode(map[string]any{"object": "list", "data": data}) + }) + f.srv = httptest.NewServer(mux) + t.Cleanup(f.srv.Close) + return f +} + +func TestFirstPollMakesHealthy(t *testing.T) { + a := newFake(t, "zeta", "alpha", "alpha") + tbl := health.New(map[string]string{"a": a.srv.URL}, time.Hour, nil) + if s, ok := tbl.Get("a"); !ok || s.Healthy || len(s.Loaded) != 0 { + t.Fatalf("before any poll: %+v %v", s, ok) + } + tbl.PollOnce(context.Background()) + s, _ := tbl.Get("a") + if !s.Healthy || s.Consecutive != 1 || s.LastErr != "" || s.LastOK.IsZero() { + t.Errorf("after one good poll: %+v", s) + } + if len(s.Loaded) != 2 || s.Loaded[0] != "alpha" || s.Loaded[1] != "zeta" { + t.Errorf("Loaded = %v, want sorted, unique [alpha zeta]", s.Loaded) + } +} + +func TestFailureThenRecoveryNeedsTwoPolls(t *testing.T) { + a := newFake(t, "m") + tbl := health.New(map[string]string{"a": a.srv.URL}, time.Hour, nil) + ctx := context.Background() + tbl.PollOnce(ctx) + a.down.Store(true) + tbl.PollOnce(ctx) + s, _ := tbl.Get("a") + if s.Healthy || s.Consecutive != 0 || s.LastErr == "" { + t.Fatalf("after failure: %+v", s) + } + if len(s.Loaded) != 1 { + t.Errorf("Loaded is left as last seen; got %v", s.Loaded) + } + a.down.Store(false) + tbl.PollOnce(ctx) + if s, _ := tbl.Get("a"); s.Healthy || s.Consecutive != 1 { + t.Errorf("one good poll after a failure must not be healthy yet: %+v", s) + } + tbl.PollOnce(ctx) + if s, _ := tbl.Get("a"); !s.Healthy || s.Consecutive != 2 || s.LastErr != "" { + t.Errorf("two good polls: %+v", s) + } +} + +func TestMarkDown(t *testing.T) { + a := newFake(t, "m") + tbl := health.New(map[string]string{"a": a.srv.URL}, time.Hour, nil) + tbl.PollOnce(context.Background()) + tbl.MarkDown("a", "connection refused") + s, _ := tbl.Get("a") + if s.Healthy || s.Consecutive != 0 || s.LastErr != "marked down: connection refused" { + t.Errorf("after MarkDown: %+v", s) + } + tbl.MarkDown("nobody", "x") // unknown hosts are ignored, not a panic + tbl.PollOnce(context.Background()) + if s, _ := tbl.Get("a"); s.Healthy { + t.Errorf("one poll after MarkDown must not be healthy: %+v", s) + } +} + +func TestUnreachableAndUnknown(t *testing.T) { + tbl := health.New(map[string]string{"a": "http://127.0.0.1:1"}, time.Hour, &http.Client{Timeout: time.Second}) + tbl.PollOnce(context.Background()) + s, ok := tbl.Get("a") + if !ok || s.Healthy || s.LastErr == "" { + t.Errorf("unreachable host: %+v %v", s, ok) + } + if _, ok := tbl.Get("zzz"); ok { + t.Errorf("unknown host must report ok=false") + } +} + +func TestAllIsACopy(t *testing.T) { + a := newFake(t, "m") + tbl := health.New(map[string]string{"a": a.srv.URL}, time.Hour, nil) + tbl.PollOnce(context.Background()) + all := tbl.All() + all["a"].Loaded[0] = "changed" + if s, _ := tbl.Get("a"); s.Loaded[0] != "m" { + t.Errorf("All must return copies") + } + if len(all) != 1 { + t.Errorf("All = %v", all) + } +} + +func TestRunPollsOnStart(t *testing.T) { + a := newFake(t, "m") + tbl := health.New(map[string]string{"a": a.srv.URL}, 20*time.Millisecond, nil) + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { tbl.Run(ctx); close(done) }() + deadline := time.Now().Add(2 * time.Second) + for a.hits.Load() < 3 && time.Now().Before(deadline) { + time.Sleep(5 * time.Millisecond) + } + cancel() + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("Run did not return after cancel") + } + if a.hits.Load() < 3 { + t.Errorf("Run polled %d times in 2s at 20ms interval", a.hits.Load()) + } + if s, _ := tbl.Get("a"); !s.Healthy { + t.Errorf("not healthy after Run: %+v", s) + } +} diff --git a/docs/plans/v0/files/internal/proxy/proxy_test.go b/docs/plans/v0/files/internal/proxy/proxy_test.go new file mode 100644 index 0000000..8d15174 --- /dev/null +++ b/docs/plans/v0/files/internal/proxy/proxy_test.go @@ -0,0 +1,318 @@ +package proxy_test + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/config" + "git.wntrmute.dev/kyle/crossbar/internal/health" + "git.wntrmute.dev/kyle/crossbar/internal/proxy" +) + +// fakeHealth is a hand-set health table that also records MarkDown calls. +type fakeHealth struct { + mu sync.Mutex + st map[string]health.Status + marked []string +} + +func (f *fakeHealth) Get(name string) (health.Status, bool) { + f.mu.Lock() + defer f.mu.Unlock() + s, ok := f.st[name] + return s, ok +} + +func (f *fakeHealth) MarkDown(name, reason string) { + f.mu.Lock() + defer f.mu.Unlock() + f.marked = append(f.marked, name) + s := f.st[name] + s.Healthy = false + s.LastErr = reason + f.st[name] = s +} + +func (f *fakeHealth) markedHosts() []string { + f.mu.Lock() + defer f.mu.Unlock() + return append([]string{}, f.marked...) +} + +// upstream records what it received and answers with its name. +type upstream struct { + name string + srv *httptest.Server + mu sync.Mutex + reqs []recorded +} + +type recorded struct { + method, path, host, xff string + body string +} + +func newUpstream(t *testing.T, name string) *upstream { + u := &upstream{name: name} + u.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + b, _ := io.ReadAll(r.Body) + u.mu.Lock() + u.reqs = append(u.reqs, recorded{r.Method, r.URL.RequestURI(), r.Host, r.Header.Get("X-Forwarded-For"), string(b)}) + u.mu.Unlock() + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"from":%q}`, name) + })) + t.Cleanup(u.srv.Close) + return u +} + +func (u *upstream) last(t *testing.T) recorded { + u.mu.Lock() + defer u.mu.Unlock() + if len(u.reqs) == 0 { + t.Fatalf("%s: no request received", u.name) + } + return u.reqs[len(u.reqs)-1] +} + +func cfgFor(t *testing.T, alpha, beta string) *config.Config { + c, err := config.Parse(strings.NewReader(fmt.Sprintf(` +listen = "127.0.0.1:1" +[hosts.alpha] +base_url = %q +models = { "shared" = { }, "alpha-only" = { } } +[hosts.beta] +base_url = %q +models = { "shared" = { }, "beta-only" = { } } +[routes.r] +hosts = ["alpha", "beta"] +default_model = "shared" +[routes.beta-first] +hosts = ["beta", "alpha"] +`, alpha, beta))) + if err != nil { + t.Fatal(err) + } + return c +} + +func healthy(loaded ...string) health.Status { + return health.Status{Healthy: true, Loaded: loaded, Consecutive: 1} +} + +func TestSplitRoute(t *testing.T) { + for _, tc := range []struct { + path, route, rest string + ok bool + }{ + {"/a/v1/x", "a", "/v1/x", true}, + {"/a/v1/x?q=1", "a", "/v1/x?q=1", true}, + {"/a", "a", "/", true}, + {"/a/", "a", "/", true}, + {"/opencode-a/v1/chat/completions", "opencode-a", "/v1/chat/completions", true}, + {"/", "", "", false}, + {"//x", "", "", false}, + {"", "", "", false}, + {"noslash/v1", "", "", false}, + } { + route, rest, ok := proxy.SplitRoute(tc.path) + if route != tc.route || rest != tc.rest || ok != tc.ok { + t.Errorf("SplitRoute(%q) = %q %q %v, want %q %q %v", tc.path, route, rest, ok, tc.route, tc.rest, tc.ok) + } + } +} + +func TestChoose(t *testing.T) { + h := &fakeHealth{st: map[string]health.Status{ + "down": {Healthy: false, Loaded: []string{"m"}}, + "alpha": healthy("shared", "alpha-only"), + "beta": healthy("shared", "beta-only"), + }} + hosts := []string{"down", "alpha", "beta"} + if got, ok := proxy.Choose(hosts, "", h); !ok || got != "alpha" { + t.Errorf("no model: %q %v, want alpha (first healthy)", got, ok) + } + if got, ok := proxy.Choose(hosts, "beta-only", h); !ok || got != "beta" { + t.Errorf("beta-only: %q %v, want beta (has the model loaded)", got, ok) + } + if got, ok := proxy.Choose(hosts, "nobody-has-it", h); !ok || got != "alpha" { + t.Errorf("unknown model falls back to the first healthy host: %q %v", got, ok) + } + if got, ok := proxy.Choose([]string{"down", "missing"}, "m", h); ok { + t.Errorf("no healthy host must give ok=false, got %q", got) + } + if got, ok := proxy.Choose(nil, "m", h); ok { + t.Errorf("empty hosts: %q %v", got, ok) + } +} + +func TestRoutesToFirstHealthyAndRewrites(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + h := &fakeHealth{st: map[string]health.Status{"alpha": healthy("shared"), "beta": healthy("shared")}} + p := proxy.New(cfgFor(t, alpha.srv.URL, beta.srv.URL), h, nil) + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodGet, "http://crossbar.local:7777/r/v1/models?x=1", nil) + req.RemoteAddr = "10.9.8.7:5555" + p.ServeHTTP(rec, req) + if rec.Code != 200 || rec.Header().Get(proxy.HostHeader) != "alpha" { + t.Fatalf("status %d host %q body %s", rec.Code, rec.Header().Get(proxy.HostHeader), rec.Body.String()) + } + got := alpha.last(t) + if got.path != "/v1/models?x=1" { + t.Errorf("upstream path = %q, want route stripped and query kept", got.path) + } + if got.host != strings.TrimPrefix(alpha.srv.URL, "http://") { + t.Errorf("Host header = %q, want the upstream's %q", got.host, strings.TrimPrefix(alpha.srv.URL, "http://")) + } + if got.xff != "10.9.8.7" { + t.Errorf("X-Forwarded-For = %q, want the client address", got.xff) + } + if !strings.Contains(rec.Body.String(), `"from":"alpha"`) { + t.Errorf("body = %s", rec.Body.String()) + } +} + +func TestModelPreferenceAndBodyPassThrough(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + h := &fakeHealth{st: map[string]health.Status{"alpha": healthy("shared", "alpha-only"), "beta": healthy("shared", "beta-only")}} + p := proxy.New(cfgFor(t, alpha.srv.URL, beta.srv.URL), h, nil) + body := `{"model":"beta-only","messages":[{"role":"user","content":"hi"}],"stream":false}` + rec := httptest.NewRecorder() + p.ServeHTTP(rec, httptest.NewRequest(http.MethodPost, "/r/v1/chat/completions", strings.NewReader(body))) + if rec.Code != 200 || rec.Header().Get(proxy.HostHeader) != "beta" { + t.Fatalf("status %d host %q", rec.Code, rec.Header().Get(proxy.HostHeader)) + } + if got := beta.last(t); got.body != body || got.method != http.MethodPost { + t.Errorf("upstream got %+v; the body must arrive unchanged after the model peek", got) + } + // Not JSON: no model, the route default ("shared") applies, first healthy wins. + rec = httptest.NewRecorder() + p.ServeHTTP(rec, httptest.NewRequest(http.MethodPost, "/r/v1/embeddings", strings.NewReader("plain text"))) + if rec.Header().Get(proxy.HostHeader) != "alpha" { + t.Errorf("non-JSON body: host %q, want alpha", rec.Header().Get(proxy.HostHeader)) + } + if got := alpha.last(t); got.body != "plain text" { + t.Errorf("non-JSON body must pass through unchanged, got %q", got.body) + } +} + +func TestFailoverOnUpstreamError(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + h := &fakeHealth{st: map[string]health.Status{"alpha": healthy("shared"), "beta": healthy("shared")}} + p := proxy.New(cfgFor(t, alpha.srv.URL, beta.srv.URL), h, nil) + alpha.srv.Close() // health still believes alpha is up + rec := httptest.NewRecorder() + p.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/r/v1/models", nil)) + if rec.Code != http.StatusBadGateway { + t.Fatalf("first request after alpha died: %d, want 502", rec.Code) + } + var e map[string]string + if err := json.Unmarshal(rec.Body.Bytes(), &e); err != nil || e["error"] != "upstream failed" || e["host"] != "alpha" { + t.Errorf("502 body = %s", rec.Body.String()) + } + if m := h.markedHosts(); len(m) != 1 || m[0] != "alpha" { + t.Errorf("MarkDown calls = %v, want [alpha]", m) + } + rec = httptest.NewRecorder() + p.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/r/v1/models", nil)) + if rec.Code != 200 || rec.Header().Get(proxy.HostHeader) != "beta" { + t.Errorf("second request: %d %q, want 200 from beta", rec.Code, rec.Header().Get(proxy.HostHeader)) + } +} + +func TestErrors(t *testing.T) { + alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta") + h := &fakeHealth{st: map[string]health.Status{"alpha": {Healthy: false}, "beta": {Healthy: false}}} + p := proxy.New(cfgFor(t, alpha.srv.URL, beta.srv.URL), h, nil) + for _, tc := range []struct { + name, method, path string + body io.Reader + want int + msg string + }{ + {"bare slash", http.MethodGet, "/", nil, 400, "missing route"}, + {"double slash", http.MethodGet, "//v1/models", nil, 400, "missing route"}, + {"unknown route", http.MethodGet, "/nope/v1/models", nil, 404, "unknown route"}, + {"disallowed path", http.MethodGet, "/r/slots", nil, 404, "not found"}, + {"admin through proxy", http.MethodGet, "/r/_crossbar/hosts", nil, 404, "not found"}, + {"no healthy host", http.MethodGet, "/r/v1/models", nil, 503, "no healthy host"}, + {"body too large", http.MethodPost, "/r/v1/chat/completions", strings.NewReader(strings.Repeat("x", proxy.MaxBody+1)), 413, "body too large"}, + } { + rec := httptest.NewRecorder() + p.ServeHTTP(rec, httptest.NewRequest(tc.method, tc.path, tc.body)) + if rec.Code != tc.want { + t.Errorf("%s: status %d, want %d", tc.name, rec.Code, tc.want) + } + var e map[string]string + if err := json.Unmarshal(rec.Body.Bytes(), &e); err != nil || e["error"] != tc.msg { + t.Errorf("%s: body %s, want error %q", tc.name, rec.Body.String(), tc.msg) + } + if !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") { + t.Errorf("%s: errors are JSON", tc.name) + } + } + if len(h.markedHosts()) != 0 { + t.Errorf("errors before choosing a host must not mark anything down: %v", h.markedHosts()) + } +} + +// TestStreamingIsNotBuffered: the upstream writes one chunk, flushes, and then waits until the +// test has *read* that chunk. If the proxy buffered, the read would never complete. +func TestStreamingIsNotBuffered(t *testing.T) { + release := make(chan struct{}) + up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(200) + fmt.Fprint(w, "data: first\n\n") + w.(http.Flusher).Flush() + select { + case <-release: + case <-time.After(5 * time.Second): + } + fmt.Fprint(w, "data: second\n\n") + })) + t.Cleanup(up.Close) + beta := newUpstream(t, "beta") + h := &fakeHealth{st: map[string]health.Status{"alpha": healthy("shared"), "beta": healthy("shared")}} + front := httptest.NewServer(proxy.New(cfgFor(t, up.URL, beta.srv.URL), h, nil)) + t.Cleanup(front.Close) + + resp, err := http.Post(front.URL+"/r/v1/chat/completions", "application/json", strings.NewReader(`{"model":"shared","stream":true}`)) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + buf := make([]byte, 64) + done := make(chan string, 1) + go func() { + n, err := resp.Body.Read(buf) + if err != nil { + done <- "read error: " + err.Error() + return + } + done <- string(buf[:n]) + }() + select { + case got := <-done: + if !strings.HasPrefix(got, "data: first") { + t.Fatalf("first read = %q", got) + } + case <-time.After(2 * time.Second): + t.Fatal("the first chunk did not arrive before the upstream finished: the proxy buffers") + } + close(release) + rest, _ := io.ReadAll(resp.Body) + if !strings.Contains(string(rest), "data: second") { + t.Errorf("rest = %q", rest) + } + if resp.Header.Get(proxy.HostHeader) != "alpha" { + t.Errorf("host header %q", resp.Header.Get(proxy.HostHeader)) + } +} diff --git a/docs/plans/v0/files/scripts/check-lines.sh b/docs/plans/v0/files/scripts/check-lines.sh new file mode 100644 index 0000000..c8e113e --- /dev/null +++ b/docs/plans/v0/files/scripts/check-lines.sh @@ -0,0 +1,14 @@ +#!/bin/sh +# No Go source file over 400 lines. Reports every offender, then fails. +set -u +limit=400 +bad=0 +for f in $(find . -name '*.go' -not -path './.git/*' -not -path './vendor/*'); do + n=$(wc -l < "$f") || { echo "check-lines: cannot read $f" >&2; exit 2; } + if [ "$n" -gt "$limit" ]; then + echo "check-lines: $f has $n lines (limit $limit)" + bad=1 + fi +done +[ "$bad" -eq 0 ] || exit 1 +echo "check-lines: ok" diff --git a/docs/plans/v0/files/tools/smoke.sh b/docs/plans/v0/files/tools/smoke.sh new file mode 100755 index 0000000..4b21792 --- /dev/null +++ b/docs/plans/v0/files/tools/smoke.sh @@ -0,0 +1,45 @@ +#!/bin/sh +# Smoke run: two fake upstreams, one crossbar, real HTTP. Prints "smoke: ok" or fails. +# Needs: bin/crossbar and bin/fakeupstream (make build), curl. +set -eu +cd "$(dirname "$0")/.." +tmp=$(mktemp -d); trap 'kill $pids 2>/dev/null; rm -rf "$tmp"' EXIT INT TERM +pids="" +bin/fakeupstream -listen 127.0.0.1:18081 -name alpha -models ornith-1.5-35b-a3b,small-9b -down-file "$tmp/alpha.down" >"$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 example.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 + +h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/opencode-a/v1/models") +[ "$h" = "200 alpha" ] || fail "opencode-a should go to alpha, got '$h'" +h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/hermes-x/v1/models") +[ "$h" = "200 beta" ] || fail "hermes-x should go to beta, got '$h'" +h=$(curl -s -o /dev/null -w '%{http_code}' "$base/nope/v1/models") +[ "$h" = "404" ] || fail "unknown route should be 404, got '$h'" + +touch "$tmp/alpha.down"; sleep 2.5 # poll_interval is 1s in example.toml +h=$(curl -s -o /dev/null -w '%{http_code} %header{X-Crossbar-Host}' "$base/opencode-a/v1/models") +[ "$h" = "200 beta" ] || fail "with alpha down, opencode-a should fail over to beta, got '$h'" +curl -s "$base/_crossbar/hosts" | grep -q '"alpha":{"healthy":false' || fail "/_crossbar/hosts does not show alpha unhealthy: $(curl -s $base/_crossbar/hosts)" + +rm "$tmp/alpha.down"; sleep 3.5 # recovery needs two good polls +h=$(curl -s -o /dev/null -w '%header{X-Crossbar-Host}' "$base/opencode-a/v1/models") +[ "$h" = "alpha" ] || fail "alpha should be back after two good polls, got '$h'" + +# Streaming: five chunks 200 ms apart must arrive over >= 0.6 s, not all at once at the end. +start=$(date +%s%N) +first="" +curl -sN -X POST -H 'Content-Type: application/json' -d '{"model":"ornith-1.5-35b-a3b","stream":true,"messages":[]}' \ + "$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 (first ${firstms:-?} ms, last ${lastms:-?} ms): +$(cat "$tmp/stream.txt")" +grep -q 'DONE' "$tmp/stream.txt" || fail "stream did not end with [DONE]" + +grep -q 'route=opencode-a host=alpha' "$tmp/crossbar.log" || fail "no request log line" +echo "smoke: ok (stream spread $((lastms - firstms)) ms)" diff --git a/tools/run-plan.sh b/tools/run-plan.sh new file mode 100755 index 0000000..0fc9954 --- /dev/null +++ b/tools/run-plan.sh @@ -0,0 +1,64 @@ +#!/bin/sh +# Runs the tasks of a plan one after another, each in its own fresh OpenCode session, and stops +# at the first task that does not end cleanly. This is the owner's tool; the implementing model +# never runs it. Adapted from boxmaker's tools/run-plan.sh. +# +# Usage: tools/run-plan.sh docs/plans/v0 [FIRST] FIRST is a task number such as 03 +# Options (environment): +# CROSSBAR_MODEL model to use (default llama.cpp/ornith-1.5-35b-a3b) +# OPENCODE_FLAGS extra flags for `opencode run` (default --pure) +# DRY_RUN=1 print what would be run and stop +set -eu + +plan="${1:?usage: tools/run-plan.sh docs/plans/ [FIRST]}" +first="${2:-00}" +model="${CROSSBAR_MODEL:-llama.cpp/ornith-1.5-35b-a3b}" +flags="${OPENCODE_FLAGS:---pure}" +log=docs/implementer-log.md +milestone=$(basename "$plan") +runs=".state/runs/$milestone" + +[ -d "$plan" ] || { echo "run-plan: $plan is not a directory" >&2; exit 2; } +[ -f "$log" ] || { echo "run-plan: run this from the repository root" >&2; exit 2; } +mkdir -p "$runs" + +for task in "$plan"/[0-9][0-9]-*.md; do + name=$(basename "$task" .md) + number=$(echo "$name" | cut -c1-2) + [ "$number" -ge "$first" ] || continue + prompt="Read \`AGENTS.md\`, then read \`$task\` and do exactly that task." + + if [ "${DRY_RUN:-0}" = 1 ]; then + echo "opencode run -m $model $flags --title $milestone/$name \"$prompt\"" + continue + fi + if [ -n "$(git status --porcelain)" ]; then + echo "run-plan: the working tree is not clean before $name; stopping" >&2 + exit 1 + fi + + before=$(git rev-parse HEAD) + echo "=== $milestone/$name: started $(date '+%H:%M:%S')" + # shellcheck disable=SC2086 + if ! opencode run -m "$model" $flags --title "$milestone/$name" "$prompt" \ + > "$runs/$name.log" 2>&1; then + echo "run-plan: opencode exited with an error in $name; see $runs/$name.log" >&2 + exit 1 + fi + echo "=== $milestone/$name: session ended $(date '+%H:%M:%S')" + + if [ "$(git rev-parse HEAD)" = "$before" ]; then + echo "run-plan: $name made no commit; stopping" >&2 + exit 1 + fi + if [ -n "$(git status --porcelain)" ]; then + echo "run-plan: $name left uncommitted changes; stopping" >&2 + exit 1 + fi + if ! grep -Eq "^\| *$milestone/$name *\|[^|]*\| *done *\|" "$log"; then + echo "run-plan: $name has no 'done' row in $log; stopping" >&2 + exit 1 + fi + echo "=== $milestone/$name: ok ($(git rev-parse --short HEAD))" +done +echo "run-plan: all tasks done"