v0 plan: AGENTS.md, gate, five task files with given tests, run-plan driver
Tests were run against a private reference implementation: gate ok after every task in order, smoke ok (stream spread ~1000 ms). The reference is not in the repository. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@@ -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: <first undecoded key, keys sorted>, 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.<name>.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.<name>.weight`: `0` becomes `1`; negative is an error.
|
||||
- `hosts.<name>.models`: at least one. For each model in sorted order,
|
||||
`hosts.<name>.models.<model>.parallel`: `0` becomes `1`; negative is an error.
|
||||
- `routes`: at least one. For each route in sorted name order:
|
||||
- `routes.<name>`: the name must match `^[a-z0-9][a-z0-9-]*$`.
|
||||
- `routes.<name>.hosts`: at least one; every entry must be a configured host; no host twice.
|
||||
- `routes.<name>.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.
|
||||
@@ -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":"<model>", …}, …]}` 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 <base>/health` must return 200,
|
||||
else the poll fails with `LastErr` starting `health: ` (for example `health: HTTP 503`, or the
|
||||
client error); then `GET <base>/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.
|
||||
@@ -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":"<msg>"}` 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":"<name>"}`.
|
||||
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.
|
||||
@@ -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: <err>` to
|
||||
stderr and exit 1.
|
||||
- Logger: `slog.New(slog.NewTextHandler(os.Stderr, nil))`.
|
||||
- `health.New(<name -> 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.
|
||||
@@ -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.<name>.base_url|weight|models`, `routes.<name>.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.<tailnet>:7777/{env:CROSSBAR_ROUTE}/v1" },
|
||||
"models": { "ornith-1.5-35b-a3b": {} } } }
|
||||
```
|
||||
Hermes, in `config.yaml`:
|
||||
```yaml
|
||||
custom_providers:
|
||||
- name: crossbar
|
||||
base_url: http://crossbar.<tailnet>:7777/hermes-<agent>/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.
|
||||
@@ -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`.
|
||||
@@ -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
|
||||
@@ -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: <name>.
|
||||
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)
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
# crossbar example configuration. Replace <tailnet> and the addresses with your own.
|
||||
listen = "127.0.0.1:17777" # never 0.0.0.0 — bind the tailnet address in production
|
||||
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.<tailnet>: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.<tailnet>: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"]
|
||||
@@ -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=
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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"
|
||||
@@ -0,0 +1,8 @@
|
||||
listen = "0.0.0.0:7777"
|
||||
|
||||
[hosts.alpha]
|
||||
base_url = "http://alpha.example:11434"
|
||||
models = { "m" = { } }
|
||||
|
||||
[routes.r]
|
||||
hosts = ["alpha"]
|
||||
@@ -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"]
|
||||
@@ -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"]
|
||||
@@ -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"]
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -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"
|
||||
Executable
+45
@@ -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)"
|
||||
Reference in New Issue
Block a user