v1 plan: leases, limiter, chooser, fingerprint, SQLite store, accounting, admin — acceptance tests first, no reference
Every given test compiled against a panic-only interface skeleton (go vet clean); nothing was implemented. modernc.org/sqlite v1.59.0 vetted in a scratch module (WAL works); go.sum given. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,293 @@
|
||||
package admin_test
|
||||
|
||||
// v1 admin: read the tables, pin/release a route, drain a host, usage rollups, metrics.
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/admin"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/health"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/lease"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/limiter"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
type fakeHosts struct {
|
||||
st map[string]health.Status
|
||||
draining map[string]bool
|
||||
}
|
||||
|
||||
func (f *fakeHosts) All() map[string]health.Status { return f.st }
|
||||
func (f *fakeHosts) Healthy(n string) bool { return f.st[n].Healthy }
|
||||
func (f *fakeHosts) Draining(n string) bool { return f.draining[n] }
|
||||
func (f *fakeHosts) SetDraining(n string, on bool) { f.draining[n] = on }
|
||||
func (f *fakeHosts) Choose(c []string, model string) (string, bool) {
|
||||
for _, h := range c {
|
||||
if f.st[h].Healthy && !f.draining[h] {
|
||||
return h, true
|
||||
}
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
||||
type rig struct {
|
||||
h http.Handler
|
||||
store *store.Store
|
||||
leases *lease.Table
|
||||
hosts *fakeHosts
|
||||
}
|
||||
|
||||
func newRig(t *testing.T) *rig {
|
||||
cfg, err := config.Parse(strings.NewReader(`
|
||||
listen = "127.0.0.1:1"
|
||||
[hosts.alpha]
|
||||
base_url = "http://alpha:1"
|
||||
models = { "m" = { parallel = 2 } }
|
||||
[hosts.beta]
|
||||
base_url = "http://beta:1"
|
||||
models = { "m" = { parallel = 4 } }
|
||||
[routes.r]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "m"
|
||||
`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
st, err := store.Open(filepath.Join(t.TempDir(), "x.db"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
hosts := &fakeHosts{
|
||||
st: map[string]health.Status{
|
||||
"alpha": {Healthy: true, Loaded: []string{"m"}, LastOK: time.Date(2026, 9, 25, 8, 0, 0, 0, time.UTC)},
|
||||
"beta": {Healthy: false, LastErr: "HTTP 503"},
|
||||
},
|
||||
draining: map[string]bool{},
|
||||
}
|
||||
lt, err := lease.New(st, hosts, hosts, 30*time.Minute)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
lim := limiter.New()
|
||||
lim.Configure("alpha", "m", 2, 8)
|
||||
lim.Configure("beta", "m", 4, 8)
|
||||
return &rig{h: admin.Handler(cfg, hosts, lt, lim, st, hosts), store: st, leases: lt, hosts: hosts}
|
||||
}
|
||||
|
||||
func (r *rig) do(t *testing.T, method, path, body string, hdr ...string) *httptest.ResponseRecorder {
|
||||
req := httptest.NewRequest(method, path, strings.NewReader(body))
|
||||
if body != "" {
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
}
|
||||
for i := 0; i+1 < len(hdr); i += 2 {
|
||||
req.Header.Set(hdr[i], hdr[i+1])
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
r.h.ServeHTTP(rec, req)
|
||||
return rec
|
||||
}
|
||||
|
||||
func TestHostsShowsSlotsAndDrain(t *testing.T) {
|
||||
r := newRig(t)
|
||||
rec := r.do(t, "GET", "/_crossbar/hosts", "")
|
||||
if rec.Code != 200 {
|
||||
t.Fatalf("%d %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
var out map[string]admin.HostView
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
a := out["alpha"]
|
||||
if !a.Healthy || a.FreeSlots != 2 || a.InFlight != 0 || a.Queued != 0 || a.Draining || a.LastOK != "2026-09-25T08:00:00Z" {
|
||||
t.Errorf("alpha = %+v", a)
|
||||
}
|
||||
if b := out["beta"]; b.Healthy || b.LastErr != "HTTP 503" || b.FreeSlots != 4 || b.Loaded == nil {
|
||||
t.Errorf("beta = %+v (loaded must be [] not null)", b)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRoutesShowsLeases(t *testing.T) {
|
||||
r := newRig(t)
|
||||
now := time.Date(2026, 9, 25, 9, 0, 0, 0, time.UTC)
|
||||
if _, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "abc", Model: "m"}, []string{"alpha", "beta"}, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec := r.do(t, "GET", "/_crossbar/routes", "")
|
||||
var out map[string]admin.RouteView
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &out); err != nil {
|
||||
t.Fatalf("%v: %s", err, rec.Body.String())
|
||||
}
|
||||
rv := out["r"]
|
||||
if len(rv.Hosts) != 2 || rv.DefaultModel != "m" || rv.Pinned != "" {
|
||||
t.Errorf("route view = %+v", rv)
|
||||
}
|
||||
if len(rv.Leases) != 1 || rv.Leases[0].FP != "abc" || rv.Leases[0].Host != "alpha" || rv.Leases[0].State != "active" || rv.Leases[0].LastUsed != "2026-09-25T09:00:00Z" {
|
||||
t.Errorf("leases = %+v", rv.Leases)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPinReleaseDrain(t *testing.T) {
|
||||
r := newRig(t)
|
||||
rec := r.do(t, "POST", "/_crossbar/routes/r", `{"host":"beta","pin":true}`)
|
||||
if rec.Code != 200 {
|
||||
t.Fatalf("pin: %d %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err == nil || h != "" {
|
||||
// beta is unhealthy in the rig: a pin to a down host is honoured, not silently moved
|
||||
t.Errorf("acquire on a route pinned to a down host: %q %v, want ErrPinnedDown", h, err)
|
||||
}
|
||||
rec = r.do(t, "GET", "/_crossbar/routes", "")
|
||||
var out map[string]admin.RouteView
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if out["r"].Pinned != "beta" {
|
||||
t.Errorf("Pinned = %q after pin", out["r"].Pinned)
|
||||
}
|
||||
rec = r.do(t, "POST", "/_crossbar/routes/r", `{"release":true}`)
|
||||
if rec.Code != 200 {
|
||||
t.Fatalf("release: %d %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
if h, _, err := r.leases.Acquire(lease.Key{Route: "r", FP: "x", Model: "m"}, []string{"alpha", "beta"}, time.Now()); err != nil || h != "alpha" {
|
||||
t.Errorf("after release: %q %v, want alpha (the only healthy host)", h, err)
|
||||
}
|
||||
for _, tc := range []struct {
|
||||
body string
|
||||
want int
|
||||
}{
|
||||
{`{"host":"nobody","pin":true}`, 404},
|
||||
{`{"pin":true}`, 400},
|
||||
{`not json`, 400},
|
||||
{`{"release":true,"pin":true,"host":"alpha"}`, 400},
|
||||
} {
|
||||
if rec := r.do(t, "POST", "/_crossbar/routes/r", tc.body); rec.Code != tc.want {
|
||||
t.Errorf("POST %s: %d, want %d (%s)", tc.body, rec.Code, tc.want, rec.Body.String())
|
||||
}
|
||||
}
|
||||
if rec := r.do(t, "POST", "/_crossbar/routes/nope", `{"release":true}`); rec.Code != 404 {
|
||||
t.Errorf("unknown route: %d", rec.Code)
|
||||
}
|
||||
|
||||
rec = r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":true}`)
|
||||
if rec.Code != 200 || !r.hosts.Draining("alpha") {
|
||||
t.Fatalf("drain: %d %s draining=%v", rec.Code, rec.Body.String(), r.hosts.Draining("alpha"))
|
||||
}
|
||||
rec = r.do(t, "GET", "/_crossbar/hosts", "")
|
||||
var hv map[string]admin.HostView
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &hv)
|
||||
if !hv["alpha"].Draining {
|
||||
t.Errorf("hosts view must show draining")
|
||||
}
|
||||
if rec := r.do(t, "POST", "/_crossbar/hosts/alpha", `{"drain":false}`); rec.Code != 200 || r.hosts.Draining("alpha") {
|
||||
t.Errorf("undrain: %d draining=%v", rec.Code, r.hosts.Draining("alpha"))
|
||||
}
|
||||
if rec := r.do(t, "POST", "/_crossbar/hosts/nobody", `{"drain":true}`); rec.Code != 404 {
|
||||
t.Errorf("unknown host: %d", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
func seedUsage(t *testing.T, st *store.Store) {
|
||||
t0 := time.Now().UTC().Add(-time.Hour)
|
||||
for i, r := range []store.Request{
|
||||
{Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 1000, PromptTokens: 100, CachedTokens: 80, CompletionTokens: 10},
|
||||
{Route: "r", FP: "a", Model: "m", Host: "alpha", Status: 200, TotalMs: 500, QueuedMs: 30, PromptTokens: 100, CachedTokens: 100, CompletionTokens: 5},
|
||||
{Route: "r2", FP: "b", Model: "m", Host: "beta", Status: 503, TotalMs: 1, Err: "queue full"},
|
||||
} {
|
||||
r.Started = t0.Add(time.Duration(i) * time.Minute)
|
||||
if err := st.RecordRequest(r); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestUsageJSONAndText(t *testing.T) {
|
||||
r := newRig(t)
|
||||
seedUsage(t, r.store)
|
||||
rec := r.do(t, "GET", "/_crossbar/usage?by=route", "")
|
||||
if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") {
|
||||
t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type"))
|
||||
}
|
||||
var rows []store.UsageRow
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil {
|
||||
t.Fatalf("%v: %s", err, rec.Body.String())
|
||||
}
|
||||
if len(rows) != 2 {
|
||||
t.Fatalf("rows = %+v", rows)
|
||||
}
|
||||
for _, row := range rows {
|
||||
if row.Key == "r" && (row.Requests != 2 || row.CachedTokens != 180 || row.QueuedMs != 30) {
|
||||
t.Errorf("r = %+v", row)
|
||||
}
|
||||
if row.Key == "r2" && (row.Requests != 1 || row.Errors != 1) {
|
||||
t.Errorf("r2 = %+v", row)
|
||||
}
|
||||
}
|
||||
rec = r.do(t, "GET", "/_crossbar/usage?by=host&since=24h", "", "Accept", "text/plain")
|
||||
if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") {
|
||||
t.Fatalf("text: %d %q", rec.Code, rec.Header().Get("Content-Type"))
|
||||
}
|
||||
body := rec.Body.String()
|
||||
if !strings.Contains(body, "alpha") || !strings.Contains(body, "beta") || !strings.Contains(strings.ToLower(body), "cache") {
|
||||
t.Errorf("text table = %q", body)
|
||||
}
|
||||
if rec := r.do(t, "GET", "/_crossbar/usage?by=colour", ""); rec.Code != 400 {
|
||||
t.Errorf("bad by: %d", rec.Code)
|
||||
}
|
||||
if rec := r.do(t, "GET", "/_crossbar/usage?since=yesterday", ""); rec.Code != 400 {
|
||||
t.Errorf("bad since: %d", rec.Code)
|
||||
}
|
||||
rec = r.do(t, "GET", "/_crossbar/usage?since=2026-09-25T00:00:00Z&by=model", "")
|
||||
if rec.Code != 200 {
|
||||
t.Errorf("RFC3339 since: %d %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestMetrics(t *testing.T) {
|
||||
r := newRig(t)
|
||||
seedUsage(t, r.store)
|
||||
rec := r.do(t, "GET", "/_crossbar/metrics", "")
|
||||
if rec.Code != 200 || !strings.HasPrefix(rec.Header().Get("Content-Type"), "text/plain") {
|
||||
t.Fatalf("%d %q", rec.Code, rec.Header().Get("Content-Type"))
|
||||
}
|
||||
body := rec.Body.String()
|
||||
for _, want := range []string{
|
||||
`# TYPE crossbar_requests_total counter`,
|
||||
`crossbar_requests_total{route="r",host="alpha",status="200"} 2`,
|
||||
`crossbar_requests_total{route="r2",host="beta",status="503"} 1`,
|
||||
`crossbar_host_healthy{host="alpha"} 1`,
|
||||
`crossbar_host_healthy{host="beta"} 0`,
|
||||
`crossbar_host_free_slots{host="alpha"} 2`,
|
||||
`crossbar_prompt_tokens_total{route="r"} 200`,
|
||||
`crossbar_cached_tokens_total{route="r"} 180`,
|
||||
`crossbar_queue_wait_ms_total{route="r"} 30`,
|
||||
} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("metrics missing %q\n%s", want, body)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMethodsAndUnknown(t *testing.T) {
|
||||
r := newRig(t)
|
||||
for _, tc := range []struct {
|
||||
method, path string
|
||||
want int
|
||||
}{
|
||||
{http.MethodPost, "/_crossbar/hosts", 405},
|
||||
{http.MethodDelete, "/_crossbar/routes", 405},
|
||||
{http.MethodGet, "/_crossbar/routes/r", 405},
|
||||
{http.MethodGet, "/_crossbar/nope", 404},
|
||||
{http.MethodPut, "/_crossbar/usage", 405},
|
||||
} {
|
||||
rec := r.do(t, tc.method, tc.path, "")
|
||||
if rec.Code != tc.want || !strings.HasPrefix(rec.Header().Get("Content-Type"), "application/json") {
|
||||
t.Errorf("%s %s = %d %q, want %d JSON", tc.method, tc.path, rec.Code, rec.Header().Get("Content-Type"), tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,83 @@
|
||||
package choose_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/choose"
|
||||
)
|
||||
|
||||
func infoFor(m map[string]choose.Info) func(string) (choose.Info, bool) {
|
||||
return func(name string) (choose.Info, bool) { i, ok := m[name]; return i, ok }
|
||||
}
|
||||
|
||||
func TestMostFreeSlotsTimesWeightWins(t *testing.T) {
|
||||
info := infoFor(map[string]choose.Info{
|
||||
"alpha": {Healthy: true, Loaded: true, CanServe: true, Free: 3, Weight: 1.0},
|
||||
"beta": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 2.0}, // 4 > 3
|
||||
"gamma": {Healthy: true, Loaded: true, CanServe: true, Free: 4, Weight: 0.5}, // 2
|
||||
})
|
||||
got, ok := choose.Best([]string{"alpha", "beta", "gamma"}, info)
|
||||
if !ok || got != "beta" {
|
||||
t.Errorf("got %q %v, want beta", got, ok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTieGoesToShortestQueueThenListOrder(t *testing.T) {
|
||||
info := infoFor(map[string]choose.Info{
|
||||
"alpha": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 3},
|
||||
"beta": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 1},
|
||||
"gamma": {Healthy: true, Loaded: true, CanServe: true, Free: 2, Weight: 1, Queued: 1},
|
||||
})
|
||||
if got, _ := choose.Best([]string{"alpha", "beta", "gamma"}, info); got != "beta" {
|
||||
t.Errorf("tie on score: shortest queue wins, then list order; got %q", got)
|
||||
}
|
||||
if got, _ := choose.Best([]string{"gamma", "beta"}, info); got != "gamma" {
|
||||
t.Errorf("full tie: first in list order wins; got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadedBeatsMerelyCapable(t *testing.T) {
|
||||
info := infoFor(map[string]choose.Info{
|
||||
"alpha": {Healthy: true, Loaded: false, CanServe: true, Free: 8, Weight: 4},
|
||||
"beta": {Healthy: true, Loaded: true, CanServe: true, Free: 1, Weight: 1},
|
||||
})
|
||||
got, ok := choose.Best([]string{"alpha", "beta"}, info)
|
||||
if !ok || got != "beta" {
|
||||
t.Errorf("a host that has the model loaded wins over one that would have to load it; got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFallsBackToCapableHost(t *testing.T) {
|
||||
info := infoFor(map[string]choose.Info{
|
||||
"alpha": {Healthy: true, Loaded: false, CanServe: true, Free: 1, Weight: 1},
|
||||
"beta": {Healthy: true, Loaded: false, CanServe: false, Free: 9, Weight: 9},
|
||||
})
|
||||
got, ok := choose.Best([]string{"beta", "alpha"}, info)
|
||||
if !ok || got != "alpha" {
|
||||
t.Errorf("only a host configured to serve the model may load it; got %q %v", got, ok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSkipsUnhealthyDrainingUnknownAndFull(t *testing.T) {
|
||||
info := infoFor(map[string]choose.Info{
|
||||
"down": {Healthy: false, Loaded: true, CanServe: true, Free: 9, Weight: 9},
|
||||
"drain": {Healthy: true, Draining: true, Loaded: true, CanServe: true, Free: 9, Weight: 9},
|
||||
"full": {Healthy: true, Loaded: true, CanServe: true, Free: 0, Weight: 9, Queued: 0},
|
||||
"ok": {Healthy: true, Loaded: true, CanServe: true, Free: 1, Weight: 1},
|
||||
})
|
||||
got, ok := choose.Best([]string{"down", "drain", "missing", "full", "ok"}, info)
|
||||
if !ok || got != "ok" {
|
||||
t.Errorf("got %q %v, want ok", got, ok)
|
||||
}
|
||||
// A full host is still better than nothing: it gets the request (it will queue).
|
||||
got, ok = choose.Best([]string{"down", "full"}, info)
|
||||
if !ok || got != "full" {
|
||||
t.Errorf("with only a full host left it must still be chosen; got %q %v", got, ok)
|
||||
}
|
||||
if _, ok := choose.Best([]string{"down", "drain", "missing"}, info); ok {
|
||||
t.Errorf("nothing usable must give ok=false")
|
||||
}
|
||||
if _, ok := choose.Best(nil, info); ok {
|
||||
t.Errorf("empty candidates must give ok=false")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
package config_test
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
)
|
||||
|
||||
const v1Base = `
|
||||
listen = "127.0.0.1:1"
|
||||
[hosts.a]
|
||||
base_url = "http://a:1"
|
||||
models = { "m" = { } }
|
||||
[routes.r]
|
||||
hosts = ["a"]
|
||||
`
|
||||
|
||||
func TestV1Defaults(t *testing.T) {
|
||||
c, err := config.Parse(strings.NewReader(v1Base))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if c.DB != "crossbar.db" {
|
||||
t.Errorf("DB default = %q", c.DB)
|
||||
}
|
||||
if c.LeaseIdle.Duration != 30*time.Minute {
|
||||
t.Errorf("LeaseIdle default = %v", c.LeaseIdle.Duration)
|
||||
}
|
||||
if c.Retention.Duration != 180*24*time.Hour {
|
||||
t.Errorf("Retention default = %v", c.Retention.Duration)
|
||||
}
|
||||
}
|
||||
|
||||
func TestV1Values(t *testing.T) {
|
||||
c, err := config.Parse(strings.NewReader(`
|
||||
db = "/var/lib/crossbar/crossbar.db"
|
||||
lease_idle = "45m"
|
||||
retention = "30d"
|
||||
` + v1Base))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if c.DB != "/var/lib/crossbar/crossbar.db" || c.LeaseIdle.Duration != 45*time.Minute || c.Retention.Duration != 30*24*time.Hour {
|
||||
t.Errorf("got db %q idle %v retention %v", c.DB, c.LeaseIdle.Duration, c.Retention.Duration)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDurationAcceptsDays(t *testing.T) {
|
||||
var d config.Duration
|
||||
for _, tc := range []struct {
|
||||
in string
|
||||
want time.Duration
|
||||
}{
|
||||
{"1d", 24 * time.Hour}, {"7d", 7 * 24 * time.Hour}, {"90m", 90 * time.Minute}, {"2h30m", 150 * time.Minute},
|
||||
} {
|
||||
if err := d.UnmarshalText([]byte(tc.in)); err != nil || d.Duration != tc.want {
|
||||
t.Errorf("UnmarshalText(%q) = %v %v, want %v", tc.in, d.Duration, err, tc.want)
|
||||
}
|
||||
}
|
||||
for _, bad := range []string{"1.5d", "d", "3 days", "1d2h"} {
|
||||
if err := d.UnmarshalText([]byte(bad)); err == nil {
|
||||
t.Errorf("UnmarshalText(%q) must fail", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestV1Validation(t *testing.T) {
|
||||
for _, tc := range []struct{ name, text, field string }{
|
||||
{"empty db", "db = \"\"\n" + v1Base, "db"},
|
||||
{"lease_idle too short", "lease_idle = \"10s\"\n" + v1Base, "lease_idle"},
|
||||
{"retention too short", "retention = \"12h\"\n" + v1Base, "retention"},
|
||||
} {
|
||||
_, err := config.Parse(strings.NewReader(tc.text))
|
||||
e, ok := config.IsError(err)
|
||||
if !ok || e.Field != tc.field {
|
||||
t.Errorf("%s: %v, want *Error on %s", tc.name, err, tc.field)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
package fingerprint_test
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/fingerprint"
|
||||
)
|
||||
|
||||
const conv1 = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Add a config loader."},{"role":"assistant","content":"Sure."},{"role":"user","content":"Now tests."}]}`
|
||||
const conv1later = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Add a config loader."},{"role":"assistant","content":"Sure."},{"role":"user","content":"Now tests."},{"role":"assistant","content":"Done."},{"role":"user","content":"And docs."}]}`
|
||||
const conv2 = `{"model":"m","messages":[{"role":"system","content":"You are the project A assistant."},{"role":"user","content":"Fix the flaky test."}]}`
|
||||
const conv3 = `{"model":"m","messages":[{"role":"system","content":"You are the project B assistant."},{"role":"user","content":"Add a config loader."}]}`
|
||||
|
||||
func TestSameConversationSameKey(t *testing.T) {
|
||||
a := fingerprint.Of([]byte(conv1))
|
||||
b := fingerprint.Of([]byte(conv1later))
|
||||
if a == "" || a != b {
|
||||
t.Errorf("later turns of one conversation must keep the key: %q vs %q", a, b)
|
||||
}
|
||||
if len(a) != 64 || strings.Trim(a, "0123456789abcdef") != "" {
|
||||
t.Errorf("key must be lowercase hex sha256 (64 chars), got %q", a)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDifferentConversationsDifferentKeys(t *testing.T) {
|
||||
a, b, c := fingerprint.Of([]byte(conv1)), fingerprint.Of([]byte(conv2)), fingerprint.Of([]byte(conv3))
|
||||
if a == b {
|
||||
t.Errorf("different first user message must change the key")
|
||||
}
|
||||
if a == c {
|
||||
t.Errorf("different system prompt must change the key")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNoUserMessageIsEmpty(t *testing.T) {
|
||||
for _, body := range []string{
|
||||
`{"model":"m","messages":[{"role":"system","content":"only a system prompt"}]}`,
|
||||
`{"model":"m","messages":[]}`,
|
||||
`{"model":"m"}`,
|
||||
`{"input":"an embeddings request"}`,
|
||||
`not json at all`,
|
||||
``,
|
||||
} {
|
||||
if got := fingerprint.Of([]byte(body)); got != "" {
|
||||
t.Errorf("Of(%q) = %q, want empty", body, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOnlyTheFirstFourKiBCount(t *testing.T) {
|
||||
long := strings.Repeat("x", 5000)
|
||||
a := `{"messages":[{"role":"user","content":"` + long + `A"}]}`
|
||||
b := `{"messages":[{"role":"user","content":"` + long + `B"}]}`
|
||||
if fingerprint.Of([]byte(a)) != fingerprint.Of([]byte(b)) {
|
||||
t.Errorf("bytes after the first 4 KiB of a message must not change the key")
|
||||
}
|
||||
c := `{"messages":[{"role":"user","content":"A` + long + `"}]}`
|
||||
if fingerprint.Of([]byte(a)) == fingerprint.Of([]byte(c)) {
|
||||
t.Errorf("bytes inside the first 4 KiB must change the key")
|
||||
}
|
||||
}
|
||||
|
||||
func TestContentPartsAreFlattened(t *testing.T) {
|
||||
plain := `{"messages":[{"role":"user","content":"hello world"}]}`
|
||||
parts := `{"messages":[{"role":"user","content":[{"type":"text","text":"hello world"}]}]}`
|
||||
if fingerprint.Of([]byte(plain)) != fingerprint.Of([]byte(parts)) {
|
||||
t.Errorf("a content array of text parts must fingerprint like the joined text")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,291 @@
|
||||
package lease_test
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/lease"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
// memPersister is an in-memory Persister that also counts writes.
|
||||
type memPersister struct {
|
||||
mu sync.Mutex
|
||||
leases map[[3]string]store.Lease
|
||||
events []store.LeaseEvent
|
||||
saves int
|
||||
}
|
||||
|
||||
func newPersister() *memPersister { return &memPersister{leases: map[[3]string]store.Lease{}} }
|
||||
|
||||
func (m *memPersister) SaveLease(l store.Lease) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.saves++
|
||||
m.leases[[3]string{l.Route, l.FP, l.Model}] = l
|
||||
return nil
|
||||
}
|
||||
func (m *memPersister) DeleteLease(route, fp, model string) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
delete(m.leases, [3]string{route, fp, model})
|
||||
return nil
|
||||
}
|
||||
func (m *memPersister) ListLeases() ([]store.Lease, error) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
out := []store.Lease{}
|
||||
for _, l := range m.leases {
|
||||
out = append(out, l)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
func (m *memPersister) RecordEvent(e store.LeaseEvent) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.events = append(m.events, e)
|
||||
return nil
|
||||
}
|
||||
func (m *memPersister) reasons() []string {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
var r []string
|
||||
for _, e := range m.events {
|
||||
r = append(r, e.Reason)
|
||||
}
|
||||
return r
|
||||
}
|
||||
|
||||
// world is a hand-set view of hosts plus a chooser that returns a fixed answer.
|
||||
type world struct {
|
||||
mu sync.Mutex
|
||||
healthy map[string]bool
|
||||
draining map[string]bool
|
||||
pick string
|
||||
picks []string // candidates seen by Choose, for assertions
|
||||
}
|
||||
|
||||
func (w *world) Healthy(name string) bool { w.mu.Lock(); defer w.mu.Unlock(); return w.healthy[name] }
|
||||
func (w *world) Draining(name string) bool { w.mu.Lock(); defer w.mu.Unlock(); return w.draining[name] }
|
||||
func (w *world) Choose(candidates []string, model string) (string, bool) {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
w.picks = append([]string{}, candidates...)
|
||||
for _, c := range candidates {
|
||||
if c == w.pick {
|
||||
return c, true
|
||||
}
|
||||
}
|
||||
if len(candidates) > 0 {
|
||||
return candidates[0], true
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
||||
var t0 = time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC)
|
||||
|
||||
func newTable(t *testing.T, p *memPersister, w *world) *lease.Table {
|
||||
tbl, err := lease.New(p, w, w, 30*time.Minute)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return tbl
|
||||
}
|
||||
|
||||
func TestNewLeaseThenSticky(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "beta"}
|
||||
tbl := newTable(t, p, w)
|
||||
k := lease.Key{Route: "r", FP: "conv1", Model: "m"}
|
||||
host, reused, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0)
|
||||
if err != nil || host != "beta" || reused {
|
||||
t.Fatalf("first: %q %v %v", host, reused, err)
|
||||
}
|
||||
w.pick = "alpha" // the chooser would now prefer alpha; the lease must hold
|
||||
for i := 1; i <= 5; i++ {
|
||||
host, reused, err = tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Duration(i)*time.Minute))
|
||||
if err != nil || host != "beta" || !reused {
|
||||
t.Fatalf("turn %d: %q reused=%v %v, want beta reused", i, host, reused, err)
|
||||
}
|
||||
}
|
||||
if got := p.reasons(); len(got) != 1 || got[0] != store.ReasonNew {
|
||||
t.Errorf("events = %v, want one 'new'", got)
|
||||
}
|
||||
snap := tbl.Snapshot()
|
||||
if len(snap) != 1 || snap[0].Host != "beta" || !snap[0].LastUsed.Equal(t0.Add(5*time.Minute)) {
|
||||
t.Errorf("snapshot = %+v", snap)
|
||||
}
|
||||
if p.saves < 2 {
|
||||
t.Errorf("LastUsed must be written through (saves=%d)", p.saves)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnhealthyHostMovesTheLease(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"}
|
||||
tbl := newTable(t, p, w)
|
||||
k := lease.Key{Route: "r", FP: "c", Model: "m"}
|
||||
if host, _, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0); host != "alpha" {
|
||||
t.Fatalf("first: %q", host)
|
||||
}
|
||||
w.mu.Lock()
|
||||
w.healthy["alpha"] = false
|
||||
w.pick = "beta"
|
||||
w.mu.Unlock()
|
||||
host, reused, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Minute))
|
||||
if err != nil || host != "beta" || reused {
|
||||
t.Fatalf("after alpha down: %q reused=%v %v", host, reused, err)
|
||||
}
|
||||
if got := p.reasons(); len(got) != 2 || got[1] != store.ReasonUnhealthy {
|
||||
t.Errorf("events = %v, want [new unhealthy]", got)
|
||||
}
|
||||
if len(w.picks) != 1 || w.picks[0] != "beta" {
|
||||
t.Errorf("Choose must not see the unhealthy host: %v", w.picks)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIdleExpiry(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"}
|
||||
tbl := newTable(t, p, w)
|
||||
k := lease.Key{Route: "r", FP: "c", Model: "m"}
|
||||
tbl.Acquire(k, []string{"alpha", "beta"}, t0)
|
||||
if n := tbl.ExpireIdle(t0.Add(29 * time.Minute)); n != 0 {
|
||||
t.Errorf("expired %d before lease_idle", n)
|
||||
}
|
||||
if n := tbl.ExpireIdle(t0.Add(31 * time.Minute)); n != 1 {
|
||||
t.Errorf("expired %d after lease_idle, want 1", n)
|
||||
}
|
||||
if got := p.reasons(); got[len(got)-1] != store.ReasonIdle {
|
||||
t.Errorf("events = %v, want idle last", got)
|
||||
}
|
||||
w.pick = "beta"
|
||||
if host, reused, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(32*time.Minute)); host != "beta" || reused {
|
||||
t.Errorf("after expiry a new lease is chosen: %q reused=%v", host, reused)
|
||||
}
|
||||
if l, _ := p.ListLeases(); len(l) != 1 {
|
||||
t.Errorf("persister holds %d leases, want 1", len(l))
|
||||
}
|
||||
}
|
||||
|
||||
func TestFingerprintInheritsRouteLease(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "beta"}
|
||||
tbl := newTable(t, p, w)
|
||||
// A request without a fingerprint (no user message) leases the route itself…
|
||||
if host, _, _ := tbl.Acquire(lease.Key{Route: "r", FP: "", Model: "m"}, []string{"alpha", "beta"}, t0); host != "beta" {
|
||||
t.Fatalf("route lease: %q", host)
|
||||
}
|
||||
w.pick = "alpha"
|
||||
// …and a new conversation on that route starts where the route already is.
|
||||
host, reused, err := tbl.Acquire(lease.Key{Route: "r", FP: "conv", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second))
|
||||
if err != nil || host != "beta" || !reused {
|
||||
t.Errorf("fingerprint lease must inherit the route's host: %q reused=%v %v", host, reused, err)
|
||||
}
|
||||
if len(tbl.Snapshot()) != 2 {
|
||||
t.Errorf("both the route lease and the conversation lease exist: %+v", tbl.Snapshot())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPinAndUnpin(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"}
|
||||
tbl := newTable(t, p, w)
|
||||
k := lease.Key{Route: "r", FP: "c", Model: "m"}
|
||||
tbl.Acquire(k, []string{"alpha", "beta"}, t0)
|
||||
if err := tbl.Pin("r", "beta", t0.Add(time.Minute)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
host, _, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(2*time.Minute))
|
||||
if err != nil || host != "beta" {
|
||||
t.Fatalf("pinned route must go to beta: %q %v", host, err)
|
||||
}
|
||||
host, _, err = tbl.Acquire(lease.Key{Route: "r", FP: "other", Model: "m"}, []string{"alpha", "beta"}, t0.Add(2*time.Minute))
|
||||
if err != nil || host != "beta" {
|
||||
t.Fatalf("new conversations on a pinned route go to the pin too: %q %v", host, err)
|
||||
}
|
||||
w.mu.Lock()
|
||||
w.healthy["beta"] = false
|
||||
w.mu.Unlock()
|
||||
if _, _, err := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(3*time.Minute)); !errors.Is(err, lease.ErrPinnedDown) {
|
||||
t.Errorf("a pinned host that is down is ErrPinnedDown, never a silent move: %v", err)
|
||||
}
|
||||
if err := tbl.Pin("r", "nobody", t0); !errors.Is(err, lease.ErrUnknownHost) {
|
||||
t.Errorf("pinning to a host not in the candidates of any lease: %v, want ErrUnknownHost", err)
|
||||
}
|
||||
tbl.Unpin("r")
|
||||
w.mu.Lock()
|
||||
w.healthy["beta"] = true
|
||||
w.mu.Unlock()
|
||||
if host, _, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(4*time.Minute)); host != "beta" {
|
||||
t.Errorf("after unpin the existing lease (on beta) simply continues: %q", host)
|
||||
}
|
||||
if got := p.reasons(); got[len(got)-3] != store.ReasonPin || got[len(got)-1] != store.ReasonRelease {
|
||||
t.Errorf("events = %v, want a pin event and a release event", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDrainKeepsExistingRefusesNew(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, draining: map[string]bool{}, pick: "alpha"}
|
||||
tbl := newTable(t, p, w)
|
||||
k := lease.Key{Route: "r", FP: "c", Model: "m"}
|
||||
tbl.Acquire(k, []string{"alpha", "beta"}, t0)
|
||||
w.mu.Lock()
|
||||
w.draining["alpha"] = true
|
||||
w.mu.Unlock()
|
||||
if host, reused, _ := tbl.Acquire(k, []string{"alpha", "beta"}, t0.Add(time.Minute)); host != "alpha" || !reused {
|
||||
t.Errorf("an existing lease on a draining host continues: %q reused=%v", host, reused)
|
||||
}
|
||||
host, _, err := tbl.Acquire(lease.Key{Route: "r2", FP: "x", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Minute))
|
||||
if err != nil || host != "beta" {
|
||||
t.Errorf("a new lease avoids the draining host: %q %v", host, err)
|
||||
}
|
||||
if len(w.picks) != 1 || w.picks[0] != "beta" {
|
||||
t.Errorf("Choose must not see the draining host: %v", w.picks)
|
||||
}
|
||||
if _, _, err := tbl.Acquire(lease.Key{Route: "r3", FP: "y", Model: "m"}, []string{"alpha"}, t0); !errors.Is(err, lease.ErrNoHost) {
|
||||
t.Errorf("only draining candidates: %v, want ErrNoHost", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseRoute(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"}
|
||||
tbl := newTable(t, p, w)
|
||||
tbl.Acquire(lease.Key{Route: "r", FP: "a", Model: "m"}, []string{"alpha", "beta"}, t0)
|
||||
tbl.Acquire(lease.Key{Route: "r", FP: "b", Model: "m"}, []string{"alpha", "beta"}, t0)
|
||||
tbl.Acquire(lease.Key{Route: "other", FP: "c", Model: "m"}, []string{"alpha", "beta"}, t0)
|
||||
if n := tbl.Release("r"); n != 2 {
|
||||
t.Errorf("Release removed %d, want 2", n)
|
||||
}
|
||||
if n := tbl.Release("r"); n != 0 {
|
||||
t.Errorf("second Release removed %d", n)
|
||||
}
|
||||
if l, _ := p.ListLeases(); len(l) != 1 || l[0].Route != "other" {
|
||||
t.Errorf("persister after release: %+v", l)
|
||||
}
|
||||
w.pick = "beta"
|
||||
if host, reused, _ := tbl.Acquire(lease.Key{Route: "r", FP: "a", Model: "m"}, []string{"alpha", "beta"}, t0); host != "beta" || reused {
|
||||
t.Errorf("after release the route is re-chosen: %q reused=%v", host, reused)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadsFromPersister(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{"alpha": true, "beta": true}, pick: "alpha"}
|
||||
_ = p.SaveLease(store.Lease{Route: "r", FP: "c", Model: "m", Host: "beta", State: store.Active, Created: t0, LastUsed: t0})
|
||||
_ = p.SaveLease(store.Lease{Route: "pinned", FP: "", Model: "", Host: "beta", State: store.Pinned, Created: t0, LastUsed: t0})
|
||||
tbl := newTable(t, p, w)
|
||||
if host, reused, _ := tbl.Acquire(lease.Key{Route: "r", FP: "c", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second)); host != "beta" || !reused {
|
||||
t.Errorf("a restart must not reshuffle: %q reused=%v", host, reused)
|
||||
}
|
||||
if host, _, _ := tbl.Acquire(lease.Key{Route: "pinned", FP: "new", Model: "m"}, []string{"alpha", "beta"}, t0.Add(time.Second)); host != "beta" {
|
||||
t.Errorf("a pin survives a restart: %q", host)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNoCandidates(t *testing.T) {
|
||||
p, w := newPersister(), &world{healthy: map[string]bool{}, pick: ""}
|
||||
tbl := newTable(t, p, w)
|
||||
if _, _, err := tbl.Acquire(lease.Key{Route: "r", FP: "c", Model: "m"}, []string{"alpha"}, t0); !errors.Is(err, lease.ErrNoHost) {
|
||||
t.Errorf("no healthy host: %v, want ErrNoHost", err)
|
||||
}
|
||||
if len(tbl.Snapshot()) != 0 {
|
||||
t.Errorf("a failed acquire must not create a lease")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,178 @@
|
||||
package limiter_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/limiter"
|
||||
)
|
||||
|
||||
func TestParallelAndQueue(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("alpha", "m", 2, 1) // two slots, one waiting place
|
||||
ctx := context.Background()
|
||||
|
||||
rel1, w1, err := l.Acquire(ctx, "alpha", "m")
|
||||
if err != nil || w1 > 50*time.Millisecond {
|
||||
t.Fatalf("first acquire: err %v waited %v", err, w1)
|
||||
}
|
||||
rel2, _, err := l.Acquire(ctx, "alpha", "m")
|
||||
if err != nil {
|
||||
t.Fatalf("second acquire: %v", err)
|
||||
}
|
||||
if l.InFlight("alpha", "m") != 2 || l.FreeSlots("alpha") != 0 {
|
||||
t.Errorf("in flight %d free %d, want 2 and 0", l.InFlight("alpha", "m"), l.FreeSlots("alpha"))
|
||||
}
|
||||
|
||||
// Third waits in the queue.
|
||||
got3 := make(chan error, 1)
|
||||
go func() {
|
||||
rel, waited, err := l.Acquire(ctx, "alpha", "m")
|
||||
if err == nil {
|
||||
defer rel()
|
||||
if waited < 40*time.Millisecond {
|
||||
err = errors.New("third acquire did not wait")
|
||||
}
|
||||
}
|
||||
got3 <- err
|
||||
}()
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
if l.Queued("alpha", "m") != 1 {
|
||||
t.Errorf("queued = %d, want 1", l.Queued("alpha", "m"))
|
||||
}
|
||||
// Fourth finds the queue full and is refused at once.
|
||||
start := time.Now()
|
||||
_, _, err = l.Acquire(ctx, "alpha", "m")
|
||||
if !errors.Is(err, limiter.ErrQueueFull) {
|
||||
t.Fatalf("fourth acquire: %v, want ErrQueueFull", err)
|
||||
}
|
||||
if time.Since(start) > 50*time.Millisecond {
|
||||
t.Errorf("a full queue must refuse immediately, took %v", time.Since(start))
|
||||
}
|
||||
time.Sleep(30 * time.Millisecond)
|
||||
rel1() // frees a slot: the queued third proceeds
|
||||
select {
|
||||
case err := <-got3:
|
||||
if err != nil {
|
||||
t.Fatalf("third: %v", err)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("queued acquire did not proceed after a release")
|
||||
}
|
||||
rel2()
|
||||
if l.InFlight("alpha", "m") != 0 || l.Queued("alpha", "m") != 0 {
|
||||
t.Errorf("after releases: inflight %d queued %d", l.InFlight("alpha", "m"), l.Queued("alpha", "m"))
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseIsIdempotent(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("h", "m", 1, 0)
|
||||
rel, _, err := l.Acquire(context.Background(), "h", "m")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rel()
|
||||
rel() // a second call must not free a slot that was never taken
|
||||
if l.InFlight("h", "m") != 0 {
|
||||
t.Errorf("in flight %d after double release", l.InFlight("h", "m"))
|
||||
}
|
||||
if _, _, err := l.Acquire(context.Background(), "h", "m"); err != nil {
|
||||
t.Errorf("slot must be free again: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCancelWhileQueuedLeaksNothing(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("h", "m", 1, 2)
|
||||
rel, _, err := l.Acquire(context.Background(), "h", "m")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
done := make(chan error, 1)
|
||||
go func() { _, _, err := l.Acquire(ctx, "h", "m"); done <- err }()
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
cancel()
|
||||
select {
|
||||
case err := <-done:
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("cancelled acquire returned %v", err)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("cancelled acquire did not return")
|
||||
}
|
||||
if l.Queued("h", "m") != 0 {
|
||||
t.Errorf("queued = %d after cancel", l.Queued("h", "m"))
|
||||
}
|
||||
rel()
|
||||
if l.InFlight("h", "m") != 0 {
|
||||
t.Errorf("in flight %d, the cancelled waiter must not have taken the slot", l.InFlight("h", "m"))
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueIsFIFO(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("h", "m", 1, 8)
|
||||
rel, _, err := l.Acquire(context.Background(), "h", "m")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var mu sync.Mutex
|
||||
var order []int
|
||||
var wg sync.WaitGroup
|
||||
for i := 1; i <= 4; i++ {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
r, _, err := l.Acquire(context.Background(), "h", "m")
|
||||
if err != nil {
|
||||
t.Errorf("waiter %d: %v", i, err)
|
||||
return
|
||||
}
|
||||
mu.Lock()
|
||||
order = append(order, i)
|
||||
mu.Unlock()
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
r()
|
||||
}(i)
|
||||
time.Sleep(15 * time.Millisecond) // stagger arrivals so the order is defined
|
||||
}
|
||||
rel()
|
||||
wg.Wait()
|
||||
if len(order) != 4 || order[0] != 1 || order[1] != 2 || order[2] != 3 || order[3] != 4 {
|
||||
t.Errorf("waiters proceeded in order %v, want [1 2 3 4]", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnconfiguredPairIsOneSlotNoQueue(t *testing.T) {
|
||||
l := limiter.New()
|
||||
rel, _, err := l.Acquire(context.Background(), "x", "y")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer rel()
|
||||
if _, _, err := l.Acquire(context.Background(), "x", "y"); !errors.Is(err, limiter.ErrQueueFull) {
|
||||
t.Errorf("second acquire on an unconfigured pair: %v, want ErrQueueFull", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFreeSlotsSumsModels(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("h", "a", 4, 0)
|
||||
l.Configure("h", "b", 2, 0)
|
||||
if got := l.FreeSlots("h"); got != 6 {
|
||||
t.Fatalf("free = %d, want 6", got)
|
||||
}
|
||||
rel, _, _ := l.Acquire(context.Background(), "h", "a")
|
||||
defer rel()
|
||||
if got := l.FreeSlots("h"); got != 5 {
|
||||
t.Errorf("free = %d, want 5", got)
|
||||
}
|
||||
if l.FreeSlots("nobody") != 0 {
|
||||
t.Errorf("unknown host has no slots")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,398 @@
|
||||
package proxy_test
|
||||
|
||||
// v1 acceptance tests for the proxy: leases, queueing, accounting, header route override.
|
||||
// They drive the whole handler over real HTTP against fake upstreams; only what a client or an
|
||||
// operator can observe is asserted (status codes, headers, the accounting rows, the health table).
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/health"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/lease"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/limiter"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
// upstream is a llama-server stand-in: streams N chunks with a delay, reports usage/timings in
|
||||
// the final chunk, counts requests, and can be slowed down or killed.
|
||||
type upstream struct {
|
||||
name string
|
||||
srv *httptest.Server
|
||||
hits atomic.Int32
|
||||
delay time.Duration
|
||||
mu sync.Mutex
|
||||
last recorded
|
||||
}
|
||||
|
||||
type recorded struct{ method, path, host, xff, body string }
|
||||
|
||||
func newUpstream(t *testing.T, name string) *upstream {
|
||||
u := &upstream{name: name}
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) })
|
||||
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) {
|
||||
fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"},{"id":"`+name+`-only"}]}`)
|
||||
})
|
||||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||||
u.hits.Add(1)
|
||||
b, _ := io.ReadAll(r.Body)
|
||||
u.mu.Lock()
|
||||
u.last = recorded{r.Method, r.URL.RequestURI(), r.Host, r.Header.Get("X-Forwarded-For"), string(b)}
|
||||
u.mu.Unlock()
|
||||
var req struct {
|
||||
Stream bool `json:"stream"`
|
||||
}
|
||||
_ = json.Unmarshal(b, &req)
|
||||
w.Header().Set("X-Upstream", name)
|
||||
time.Sleep(u.delay)
|
||||
if !req.Stream {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
fmt.Fprintf(w, `{"choices":[{"message":{"role":"assistant","content":"hi from %s"}}],"usage":{"prompt_tokens":100,"completion_tokens":10,"total_tokens":110},"timings":{"prompt_n":100,"cache_n":90,"predicted_n":10,"predicted_ms":50.0}}`, name)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.WriteHeader(200)
|
||||
fl := w.(http.Flusher)
|
||||
for i := 0; i < 3; i++ {
|
||||
fmt.Fprintf(w, "data: {\"choices\":[{\"delta\":{\"content\":\"%s %d \"}}]}\n\n", name, i)
|
||||
fl.Flush()
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
fmt.Fprint(w, `data: {"choices":[],"usage":{"prompt_tokens":200,"completion_tokens":20,"total_tokens":220},"timings":{"prompt_n":200,"cache_n":150,"predicted_n":20,"predicted_ms":80.0}}`+"\n\n")
|
||||
fl.Flush()
|
||||
fmt.Fprint(w, "data: [DONE]\n\n")
|
||||
})
|
||||
u.srv = httptest.NewServer(mux)
|
||||
t.Cleanup(u.srv.Close)
|
||||
return u
|
||||
}
|
||||
|
||||
func (u *upstream) lastReq() recorded { u.mu.Lock(); defer u.mu.Unlock(); return u.last }
|
||||
|
||||
// rig is one crossbar: config, real health table (polled once), real lease table over a real
|
||||
// SQLite store, real limiter, the proxy handler served by httptest.
|
||||
type rig struct {
|
||||
t *testing.T
|
||||
cfg *config.Config
|
||||
health *health.Table
|
||||
store *store.Store
|
||||
leases *lease.Table
|
||||
lim *limiter.Limiter
|
||||
front *httptest.Server
|
||||
}
|
||||
|
||||
// newRig builds crossbar from a config text where %s placeholders are the upstream base URLs.
|
||||
func newRig(t *testing.T, cfgText string, ups ...*upstream) *rig {
|
||||
urls := make([]any, len(ups))
|
||||
for i, u := range ups {
|
||||
urls[i] = u.srv.URL
|
||||
}
|
||||
cfg, err := config.Parse(strings.NewReader(fmt.Sprintf(cfgText, urls...)))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
bases := map[string]string{}
|
||||
for name, h := range cfg.Hosts {
|
||||
bases[name] = h.BaseURL
|
||||
}
|
||||
ht := health.New(bases, time.Hour, nil)
|
||||
ht.PollOnce(t.Context())
|
||||
st, err := store.Open(filepath.Join(t.TempDir(), "crossbar.db"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
lim := limiter.New()
|
||||
for name, h := range cfg.Hosts {
|
||||
for model, m := range h.Models {
|
||||
lim.Configure(name, model, m.Parallel, cfg.QueueMax)
|
||||
}
|
||||
}
|
||||
lt, err := lease.New(st, proxy.HostView(ht, cfg), proxy.Chooser(cfg, ht, lim), cfg.LeaseIdle.Duration)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
p := proxy.New(cfg, ht, lt, lim, st, nil)
|
||||
front := httptest.NewServer(p)
|
||||
t.Cleanup(front.Close)
|
||||
return &rig{t: t, cfg: cfg, health: ht, store: st, leases: lt, lim: lim, front: front}
|
||||
}
|
||||
|
||||
const twoHosts = `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 1
|
||||
lease_idle = "30m"
|
||||
[hosts.alpha]
|
||||
base_url = %q
|
||||
weight = 1.0
|
||||
models = { "shared" = { parallel = 2 }, "alpha-only" = { } }
|
||||
[hosts.beta]
|
||||
base_url = %q
|
||||
weight = 2.0
|
||||
models = { "shared" = { parallel = 2 }, "beta-only" = { } }
|
||||
[routes.r]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "shared"
|
||||
[routes.other]
|
||||
hosts = ["alpha"]
|
||||
`
|
||||
|
||||
func conversation(id, turn int) string {
|
||||
msgs := fmt.Sprintf(`{"role":"system","content":"project"},{"role":"user","content":"conversation %d opening"}`, id)
|
||||
for i := 1; i < turn; i++ {
|
||||
msgs += fmt.Sprintf(`,{"role":"assistant","content":"ok"},{"role":"user","content":"turn %d"}`, i)
|
||||
}
|
||||
return `{"model":"shared","stream":false,"messages":[` + msgs + `]}`
|
||||
}
|
||||
|
||||
func (r *rig) post(path, body string, hdr ...string) *http.Response {
|
||||
req, _ := http.NewRequest(http.MethodPost, r.front.URL+path, strings.NewReader(body))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
for i := 0; i+1 < len(hdr); i += 2 {
|
||||
req.Header.Set(hdr[i], hdr[i+1])
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
r.t.Fatal(err)
|
||||
}
|
||||
return resp
|
||||
}
|
||||
|
||||
func drain(resp *http.Response) string {
|
||||
b, _ := io.ReadAll(resp.Body)
|
||||
resp.Body.Close()
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func TestConversationIsStickyAndLeaseHeaderTellsWhy(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
first := r.post("/r/v1/chat/completions", conversation(1, 1))
|
||||
drain(first)
|
||||
host := first.Header.Get(proxy.HostHeader)
|
||||
if first.StatusCode != 200 || host != "beta" { // beta: same free slots, double weight
|
||||
t.Fatalf("first turn: %d from %q, want 200 from beta", first.StatusCode, host)
|
||||
}
|
||||
if got := first.Header.Get(proxy.LeaseHeader); got != "new" {
|
||||
t.Errorf("%s = %q on the first turn, want new", proxy.LeaseHeader, got)
|
||||
}
|
||||
// Take alpha's slots away as a "better host" signal: it must not matter, the lease holds.
|
||||
for turn := 2; turn <= 6; turn++ {
|
||||
resp := r.post("/r/v1/chat/completions", conversation(1, turn))
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != host || resp.Header.Get(proxy.LeaseHeader) != "reused" {
|
||||
t.Fatalf("turn %d: host %q lease %q, want %q reused", turn, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader), host)
|
||||
}
|
||||
}
|
||||
if alpha.hits.Load() != 0 || beta.hits.Load() != 6 {
|
||||
t.Errorf("hits alpha=%d beta=%d, want 0 and 6", alpha.hits.Load(), beta.hits.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
beta.delay = 300 * time.Millisecond
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
// Two slow conversations occupy beta's two "shared" slots…
|
||||
var wg sync.WaitGroup
|
||||
for i := 1; i <= 2; i++ {
|
||||
wg.Add(1)
|
||||
go func(i int) { defer wg.Done(); drain(r.post("/r/v1/chat/completions", conversation(i, 1))) }(i)
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
// …so a third conversation starting now is sent to alpha (beta has 0 free slots, alpha 2).
|
||||
resp := r.post("/r/v1/chat/completions", conversation(3, 1))
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != "alpha" {
|
||||
t.Errorf("third conversation went to %q, want alpha (free slots beat weight)", resp.Header.Get(proxy.HostHeader))
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func TestQueueFullIs503(t *testing.T) {
|
||||
alpha := newUpstream(t, "alpha")
|
||||
alpha.delay = 400 * time.Millisecond
|
||||
r := newRig(t, `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 1
|
||||
[hosts.alpha]
|
||||
base_url = %q
|
||||
models = { "shared" = { parallel = 1 } }
|
||||
[routes.r]
|
||||
hosts = ["alpha"]
|
||||
default_model = "shared"
|
||||
`, alpha)
|
||||
codes := make(chan int, 3)
|
||||
for i := 1; i <= 3; i++ {
|
||||
go func(i int) {
|
||||
resp := r.post("/r/v1/chat/completions", conversation(i, 1))
|
||||
drain(resp)
|
||||
codes <- resp.StatusCode
|
||||
}(i)
|
||||
time.Sleep(30 * time.Millisecond) // arrival order: 1 runs, 2 queues, 3 finds the queue full
|
||||
}
|
||||
got := map[int]int{}
|
||||
for i := 0; i < 3; i++ {
|
||||
got[<-codes]++
|
||||
}
|
||||
if got[200] != 2 || got[503] != 1 {
|
||||
t.Fatalf("status counts = %v, want two 200 and one 503", got)
|
||||
}
|
||||
rows, err := r.store.Usage(time.Time{}, store.ByRoute)
|
||||
if err != nil || len(rows) != 1 || rows[0].Requests != 3 || rows[0].Errors != 1 {
|
||||
t.Errorf("usage = %+v %v, want 3 requests, 1 error (the 503 is recorded too)", rows, err)
|
||||
}
|
||||
if rows[0].QueuedMs <= 0 {
|
||||
t.Errorf("the queued request must record its wait: %+v", rows[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnhealthyHostReleasesAndMoves(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // lands on beta
|
||||
beta.srv.Close()
|
||||
resp := r.post("/r/v1/chat/completions", conversation(1, 2))
|
||||
drain(resp)
|
||||
if resp.StatusCode != http.StatusBadGateway {
|
||||
t.Fatalf("first request after beta died: %d, want 502", resp.StatusCode)
|
||||
}
|
||||
if s, _ := r.health.Get("beta"); s.Healthy {
|
||||
t.Fatalf("beta must be marked down after the 502")
|
||||
}
|
||||
resp = r.post("/r/v1/chat/completions", conversation(1, 3))
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" || resp.Header.Get(proxy.LeaseHeader) != "new" {
|
||||
t.Errorf("after the move: %d from %q lease %q, want 200 alpha new", resp.StatusCode, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader))
|
||||
}
|
||||
ev, _ := r.store.Events(time.Time{}, 10)
|
||||
var reasons []string
|
||||
for _, e := range ev {
|
||||
reasons = append(reasons, e.Reason)
|
||||
}
|
||||
if len(reasons) != 2 || reasons[0] != store.ReasonNew || reasons[1] != store.ReasonUnhealthy {
|
||||
t.Errorf("lease events = %v, want [new unhealthy]", reasons)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAccountingRowsFromUsageAndTimings(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
drain(r.post("/r/v1/chat/completions", conversation(1, 1))) // non-streamed
|
||||
drain(r.post("/r/v1/chat/completions", strings.Replace(conversation(1, 2), `"stream":false`, `"stream":true`, 1))) // streamed
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
var rows []store.UsageRow
|
||||
for time.Now().Before(deadline) {
|
||||
rows, _ = r.store.Usage(time.Time{}, store.ByHost)
|
||||
if len(rows) == 1 && rows[0].Requests == 2 {
|
||||
break
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
if len(rows) != 1 || rows[0].Requests != 2 {
|
||||
t.Fatalf("usage by host = %+v, want one host with 2 requests (rows may be written after the response completes, within 2 s)", rows)
|
||||
}
|
||||
u := rows[0]
|
||||
if u.PromptTokens != 300 || u.CachedTokens != 240 || u.CompletionTokens != 30 {
|
||||
t.Errorf("tokens = prompt %d cached %d completion %d, want 300/240/30 (100+200, 90+150, 10+20)", u.PromptTokens, u.CachedTokens, u.CompletionTokens)
|
||||
}
|
||||
if u.BusyMs <= 0 || u.Errors != 0 {
|
||||
t.Errorf("busy %d errors %d", u.BusyMs, u.Errors)
|
||||
}
|
||||
if got := u.CacheHitRatio(); got < 0.79 || got > 0.81 {
|
||||
t.Errorf("cache hit ratio = %v, want 0.8", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamIsUnalteredWhileTeed(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
resp := r.post("/r/v1/chat/completions", strings.Replace(conversation(9, 1), `"stream":false`, `"stream":true`, 1))
|
||||
body := drain(resp)
|
||||
want := 0
|
||||
for _, line := range strings.Split(body, "\n") {
|
||||
if strings.HasPrefix(line, "data: ") {
|
||||
want++
|
||||
}
|
||||
}
|
||||
if want != 5 || !strings.HasSuffix(strings.TrimSpace(body), "data: [DONE]") {
|
||||
t.Errorf("client must receive every SSE line untouched (3 deltas, usage, DONE); got %d data lines:\n%s", want, body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHeaderRouteOverride(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
// The header names the route; the path has none.
|
||||
resp := r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other")
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "alpha" {
|
||||
t.Errorf("header route 'other' (alpha only): %d from %q", resp.StatusCode, resp.Header.Get(proxy.HostHeader))
|
||||
}
|
||||
if alpha.lastReq().path != "/v1/chat/completions" {
|
||||
t.Errorf("upstream path = %q", alpha.lastReq().path)
|
||||
}
|
||||
// A path route and a header route that disagree: the header is the operator's intent → 400.
|
||||
resp = r.post("/r/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "other")
|
||||
if drain(resp); resp.StatusCode != 400 {
|
||||
t.Errorf("conflicting route in path and header: %d, want 400", resp.StatusCode)
|
||||
}
|
||||
resp = r.post("/v1/chat/completions", conversation(1, 1), proxy.RouteHeader, "nope")
|
||||
if drain(resp); resp.StatusCode != 404 {
|
||||
t.Errorf("unknown header route: %d, want 404", resp.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
func TestV0BehaviourStillHolds(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
for _, tc := range []struct {
|
||||
method, path string
|
||||
want int
|
||||
msg string
|
||||
}{
|
||||
{http.MethodGet, "/", 400, "missing route"},
|
||||
{http.MethodGet, "/nope/v1/models", 404, "unknown route"},
|
||||
{http.MethodGet, "/r/slots", 404, "not found"},
|
||||
{http.MethodGet, "/r/_crossbar/hosts", 404, "not found"},
|
||||
} {
|
||||
req, _ := http.NewRequest(tc.method, r.front.URL+tc.path, nil)
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
body := drain(resp)
|
||||
var e map[string]string
|
||||
if resp.StatusCode != tc.want || json.Unmarshal([]byte(body), &e) != nil || e["error"] != tc.msg {
|
||||
t.Errorf("%s: %d %s, want %d %q", tc.path, resp.StatusCode, body, tc.want, tc.msg)
|
||||
}
|
||||
}
|
||||
big := strings.Repeat("x", proxy.MaxBody+1)
|
||||
resp := r.post("/r/v1/chat/completions", big)
|
||||
if drain(resp); resp.StatusCode != 413 {
|
||||
t.Errorf("oversize body: %d, want 413", resp.StatusCode)
|
||||
}
|
||||
// GET pass-through with query string, Host and X-Forwarded-For as in v0.
|
||||
resp, err := http.Get(r.front.URL + "/r/v1/models?x=1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
drain(resp)
|
||||
host := resp.Header.Get(proxy.HostHeader)
|
||||
u := map[string]*upstream{"alpha": alpha, "beta": beta}[host]
|
||||
if u == nil || u.lastReq().path != "/v1/models?x=1" || u.lastReq().host != strings.TrimPrefix(u.srv.URL, "http://") || u.lastReq().xff == "" {
|
||||
t.Errorf("GET pass-through: host %q last %+v", host, u.lastReq())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,185 @@
|
||||
package store_test
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
func open(t *testing.T, dir string) *store.Store {
|
||||
s, err := store.Open(filepath.Join(dir, "crossbar.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = s.Close() })
|
||||
return s
|
||||
}
|
||||
|
||||
func TestOpenIsIdempotentAndWAL(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s := open(t, dir)
|
||||
if got := s.JournalMode(); got != "wal" {
|
||||
t.Errorf("journal_mode = %q, want wal", got)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
open(t, dir) // second open on the same file must not fail on existing tables
|
||||
}
|
||||
|
||||
func TestLeasesSurviveReopen(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s := open(t, dir)
|
||||
now := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC)
|
||||
l := store.Lease{Route: "opencode-a", FP: "abc", Model: "m", Host: "alpha", State: store.Active, Created: now, LastUsed: now}
|
||||
if err := s.SaveLease(l); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
l2 := l
|
||||
l2.FP = "def"
|
||||
l2.Host = "beta"
|
||||
l2.State = store.Pinned
|
||||
if err := s.SaveLease(l2); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Saving the same key again replaces, not duplicates.
|
||||
l.Host = "beta"
|
||||
l.LastUsed = now.Add(time.Minute)
|
||||
if err := s.SaveLease(l); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s = open(t, dir)
|
||||
got, err := s.ListLeases()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) != 2 {
|
||||
t.Fatalf("ListLeases = %d rows, want 2: %+v", len(got), got)
|
||||
}
|
||||
byFP := map[string]store.Lease{}
|
||||
for _, x := range got {
|
||||
byFP[x.FP] = x
|
||||
}
|
||||
if a := byFP["abc"]; a.Host != "beta" || !a.LastUsed.Equal(now.Add(time.Minute)) || a.State != store.Active {
|
||||
t.Errorf("abc = %+v", a)
|
||||
}
|
||||
if d := byFP["def"]; d.State != store.Pinned || d.Host != "beta" {
|
||||
t.Errorf("def = %+v", d)
|
||||
}
|
||||
if err := s.DeleteLease("opencode-a", "abc", "m"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, _ = s.ListLeases()
|
||||
if len(got) != 1 || got[0].FP != "def" {
|
||||
t.Errorf("after delete: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventsAndRequestsAndUsage(t *testing.T) {
|
||||
s := open(t, t.TempDir())
|
||||
t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC)
|
||||
must := func(err error) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
must(s.RecordEvent(store.LeaseEvent{TS: t0, Route: "r1", Model: "m", FromHost: "", ToHost: "alpha", Reason: store.ReasonNew}))
|
||||
must(s.RecordEvent(store.LeaseEvent{TS: t0.Add(time.Hour), Route: "r1", Model: "m", FromHost: "alpha", ToHost: "beta", Reason: store.ReasonUnhealthy}))
|
||||
reqs := []store.Request{
|
||||
{Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0, QueuedMs: 0, TTFBMs: 100, TotalMs: 1000, Status: 200, Streamed: true, PromptTokens: 1000, CachedTokens: 900, CompletionTokens: 50},
|
||||
{Route: "r1", FP: "a", Model: "m", Host: "alpha", Started: t0.Add(time.Minute), QueuedMs: 40, TTFBMs: 120, TotalMs: 2000, Status: 200, Streamed: true, PromptTokens: 1100, CachedTokens: 1000, CompletionTokens: 60},
|
||||
{Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(2 * time.Minute), TotalMs: 500, Status: 502, Err: "upstream failed"},
|
||||
{Route: "r2", FP: "b", Model: "m", Host: "beta", Started: t0.Add(-48 * time.Hour), TotalMs: 300, Status: 200, PromptTokens: 10, CompletionTokens: 5},
|
||||
}
|
||||
for _, r := range reqs {
|
||||
must(s.RecordRequest(r))
|
||||
}
|
||||
must(s.RecordHostHealth(store.HostHealth{TS: t0, Host: "alpha", Healthy: true, Loaded: []string{"m"}}))
|
||||
|
||||
rows, err := s.Usage(t0.Add(-time.Hour), store.ByRoute)
|
||||
must(err)
|
||||
if len(rows) != 2 {
|
||||
t.Fatalf("Usage by route since t0-1h: %d rows, want 2 (r1, r2): %+v", len(rows), rows)
|
||||
}
|
||||
byKey := map[string]store.UsageRow{}
|
||||
for _, r := range rows {
|
||||
byKey[r.Key] = r
|
||||
}
|
||||
r1 := byKey["r1"]
|
||||
if r1.Requests != 2 || r1.Errors != 0 || r1.BusyMs != 3000 || r1.QueuedMs != 40 {
|
||||
t.Errorf("r1 = %+v", r1)
|
||||
}
|
||||
if r1.PromptTokens != 2100 || r1.CachedTokens != 1900 || r1.CompletionTokens != 110 {
|
||||
t.Errorf("r1 tokens = %+v", r1)
|
||||
}
|
||||
if got := r1.CacheHitRatio(); got < 0.904 || got > 0.905 {
|
||||
t.Errorf("r1 cache hit ratio = %v, want 1900/2100", got)
|
||||
}
|
||||
r2 := byKey["r2"]
|
||||
if r2.Requests != 1 || r2.Errors != 1 || r2.BusyMs != 500 {
|
||||
t.Errorf("r2 = %+v (the 48h-old request is outside since)", r2)
|
||||
}
|
||||
if r2.CacheHitRatio() != 0 {
|
||||
t.Errorf("no prompt tokens: ratio must be 0, got %v", r2.CacheHitRatio())
|
||||
}
|
||||
byHost, err := s.Usage(time.Time{}, store.ByHost)
|
||||
must(err)
|
||||
if len(byHost) != 2 {
|
||||
t.Errorf("by host, all time: %+v", byHost)
|
||||
}
|
||||
for _, r := range byHost {
|
||||
if r.Key == "beta" && r.Requests != 2 {
|
||||
t.Errorf("beta all-time requests = %d, want 2", r.Requests)
|
||||
}
|
||||
}
|
||||
byModel, err := s.Usage(time.Time{}, store.ByModel)
|
||||
must(err)
|
||||
if len(byModel) != 1 || byModel[0].Key != "m" || byModel[0].Requests != 4 {
|
||||
t.Errorf("by model: %+v", byModel)
|
||||
}
|
||||
ev, err := s.Events(t0.Add(-time.Minute), 10)
|
||||
must(err)
|
||||
if len(ev) != 2 || ev[0].Reason != store.ReasonNew || ev[1].ToHost != "beta" {
|
||||
t.Errorf("events = %+v", ev)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPruneRollsUpOldRequests(t *testing.T) {
|
||||
s := open(t, t.TempDir())
|
||||
t0 := time.Date(2026, 9, 25, 10, 0, 0, 0, time.UTC)
|
||||
old := t0.Add(-200 * 24 * time.Hour)
|
||||
for i := 0; i < 3; i++ {
|
||||
if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: old.Add(time.Duration(i) * time.Minute), TotalMs: 100, Status: 200, PromptTokens: 10, CachedTokens: 5, CompletionTokens: 1}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := s.RecordRequest(store.Request{Route: "r", Model: "m", Host: "h", Started: t0, TotalMs: 100, Status: 200}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
n, err := s.Prune(t0, 180*24*time.Hour)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n != 3 {
|
||||
t.Errorf("Prune removed %d rows, want 3", n)
|
||||
}
|
||||
rows, _ := s.Usage(time.Time{}, store.ByRoute)
|
||||
if len(rows) != 1 || rows[0].Requests != 4 || rows[0].PromptTokens != 30 {
|
||||
t.Errorf("usage must still include pruned traffic through the daily rollup: %+v", rows)
|
||||
}
|
||||
live, _ := s.Usage(old.Add(24*time.Hour), store.ByRoute)
|
||||
if len(live) != 1 || live[0].Requests != 1 {
|
||||
t.Errorf("recent-only usage = %+v", live)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBadPath(t *testing.T) {
|
||||
if _, err := store.Open(filepath.Join(t.TempDir(), "no", "such", "dir", "x.db")); err == nil {
|
||||
t.Fatal("Open must fail when the directory does not exist")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user