v2.3 plan: clients that manage their own slots (control calls, route affinity, queue = false, route listeners, llama-server ctx error); given tests
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,75 @@
|
||||
package config_test
|
||||
|
||||
// v2.3 task 02: the affinity and queue route keys.
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
)
|
||||
|
||||
const affinityBase = `
|
||||
listen = "127.0.0.1:1"
|
||||
[hosts.a]
|
||||
base_url = "http://a:1"
|
||||
models = { "m" = { } }
|
||||
[routes.plain]
|
||||
hosts = ["a"]
|
||||
[routes.convo]
|
||||
hosts = ["a"]
|
||||
affinity = "conversation"
|
||||
[routes.boxmaker]
|
||||
hosts = ["a"]
|
||||
affinity = "route"
|
||||
queue = false
|
||||
[routes."bm-*"]
|
||||
hosts = ["a"]
|
||||
affinity = "route"
|
||||
queue = false
|
||||
[routes.queued]
|
||||
hosts = ["a"]
|
||||
queue = true
|
||||
`
|
||||
|
||||
func TestAffinityAndQueueKeys(t *testing.T) {
|
||||
c, err := config.Parse(strings.NewReader(affinityBase))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, tc := range []struct {
|
||||
route string
|
||||
perRoute, queues bool
|
||||
}{
|
||||
{"plain", false, true}, // defaults: conversation affinity, queueing on
|
||||
{"convo", false, true},
|
||||
{"boxmaker", true, false},
|
||||
{"bm-agent-1", true, false}, // a template's keys reach its concrete routes
|
||||
{"queued", false, true},
|
||||
} {
|
||||
r, _, ok := c.Route(tc.route)
|
||||
if !ok {
|
||||
t.Fatalf("route %q not found", tc.route)
|
||||
}
|
||||
if r.PerRoute() != tc.perRoute || r.Queues() != tc.queues {
|
||||
t.Errorf("%s: PerRoute %v Queues %v, want %v %v", tc.route, r.PerRoute(), r.Queues(), tc.perRoute, tc.queues)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAffinityRejectsUnknownValues(t *testing.T) {
|
||||
for _, bad := range []string{`"session"`, `"Route"`, `1`} {
|
||||
text := strings.Replace(affinityBase, `affinity = "conversation"`, "affinity = "+bad, 1)
|
||||
_, err := config.Parse(strings.NewReader(text))
|
||||
if err == nil || !strings.Contains(err.Error(), "routes.convo.affinity") {
|
||||
t.Errorf("affinity = %s: err %v, want one naming routes.convo.affinity", bad, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueMustBeABool(t *testing.T) {
|
||||
text := strings.Replace(affinityBase, "queue = true", `queue = "no"`, 1)
|
||||
if _, err := config.Parse(strings.NewReader(text)); err == nil {
|
||||
t.Error(`queue = "no" parsed; want an error`)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package config_test
|
||||
|
||||
// v2.3 task 03: a concrete route may own a dedicated listener. Every request that arrives on it is
|
||||
// that route, with the upstream path unprefixed, for clients that cannot put a route in the path
|
||||
// or a header (Boxmaker's inferproxy rewrites nothing).
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/config"
|
||||
)
|
||||
|
||||
const listenBase = `
|
||||
listen = "127.0.0.1:7777"
|
||||
[hosts.a]
|
||||
base_url = "http://a:1"
|
||||
models = { "m" = { } }
|
||||
[routes.bm-a]
|
||||
hosts = ["a"]
|
||||
listen = "127.0.0.1:7801"
|
||||
[routes.bm-b]
|
||||
hosts = ["a"]
|
||||
listen = "127.0.0.1:7802"
|
||||
[routes.plain]
|
||||
hosts = ["a"]
|
||||
`
|
||||
|
||||
func TestRouteListen(t *testing.T) {
|
||||
c, err := config.Parse(strings.NewReader(listenBase))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for route, want := range map[string]string{"bm-a": "127.0.0.1:7801", "bm-b": "127.0.0.1:7802", "plain": ""} {
|
||||
if got := c.Routes[route].Listen; got != want {
|
||||
t.Errorf("%s listen = %q, want %q", route, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRouteListenRejected(t *testing.T) {
|
||||
for _, tc := range []struct{ name, text, want string }{
|
||||
{"not host:port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"7801"`, 1), "routes.bm-a.listen"},
|
||||
{"bad port", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:http"`, 1), "routes.bm-a.listen"},
|
||||
{"port zero", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:0"`, 1), "routes.bm-a.listen"},
|
||||
{"same as another route", strings.Replace(listenBase, `"127.0.0.1:7802"`, `"127.0.0.1:7801"`, 1), "listen"},
|
||||
{"same as the main listener", strings.Replace(listenBase, `"127.0.0.1:7801"`, `"127.0.0.1:7777"`, 1), "routes.bm-a.listen"},
|
||||
{"on a template", listenBase + "[routes.\"t-*\"]\nhosts = [\"a\"]\nlisten = \"127.0.0.1:7803\"\n", "t-*"},
|
||||
} {
|
||||
if _, err := config.Parse(strings.NewReader(tc.text)); err == nil || !strings.Contains(err.Error(), tc.want) {
|
||||
t.Errorf("%s: err %v, want one containing %q", tc.name, err, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
package identity_test
|
||||
|
||||
// v2.3 task 03: on a route's dedicated listener the route is fixed, so the gate is that route's
|
||||
// peers for every request, whatever path or X-Crossbar-Route header the caller sends.
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/identity"
|
||||
)
|
||||
|
||||
func TestRouteMiddleware(t *testing.T) {
|
||||
inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) })
|
||||
checker := identity.NewChecker(fakeResolver{"100.64.0.5": "talos", "100.64.0.9": "titan"})
|
||||
locked := identity.RouteMiddleware(checker, []string{"talos"}, inner)
|
||||
open := identity.RouteMiddleware(checker, nil, inner)
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
h http.Handler
|
||||
path, hdr string
|
||||
addr string
|
||||
want int
|
||||
}{
|
||||
{"right peer", locked, "/v1/chat/completions", "", "100.64.0.5:5", 204},
|
||||
{"wrong peer", locked, "/v1/chat/completions", "", "100.64.0.9:5", 403},
|
||||
{"not a peer", locked, "/slots", "", "203.0.113.1:5", 403},
|
||||
{"a path that looks like an open route is still this route", locked, "/open/v1/models", "", "100.64.0.9:5", 403},
|
||||
{"a header naming another route does not change the gate", locked, "/v1/models", "open", "100.64.0.9:5", 403},
|
||||
{"admin-looking path is gated too (no admin on this listener)", locked, "/_crossbar/hosts", "", "100.64.0.9:5", 403},
|
||||
{"open route, anyone", open, "/v1/models", "", "203.0.113.1:5", 204},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
req := httptest.NewRequest(http.MethodGet, tc.path, nil)
|
||||
req.RemoteAddr = tc.addr
|
||||
if tc.hdr != "" {
|
||||
req.Header.Set("X-Crossbar-Route", tc.hdr)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
tc.h.ServeHTTP(rec, req)
|
||||
if rec.Code != tc.want {
|
||||
t.Errorf("%s = %d, want %d (%s)", tc.path, rec.Code, tc.want, rec.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
package limiter_test
|
||||
|
||||
// v2.3 task 02: Track counts a request without holding or refusing it. A route with queue = false
|
||||
// leaves queueing to llama-server's own slots, but its requests are still load on the host, so the
|
||||
// routes that do queue must see them.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/limiter"
|
||||
)
|
||||
|
||||
func TestTrackNeverWaitsAndCounts(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("alpha", "m", 1, 0) // one slot, no waiting room
|
||||
|
||||
start := time.Now()
|
||||
rel1 := l.Track("alpha", "m")
|
||||
rel2 := l.Track("alpha", "m")
|
||||
rel3 := l.Track("alpha", "m")
|
||||
if d := time.Since(start); d > 50*time.Millisecond {
|
||||
t.Fatalf("Track waited %v", d)
|
||||
}
|
||||
if n := l.InFlight("alpha", "m"); n != 3 {
|
||||
t.Fatalf("in flight = %d, want 3 (Track may pass parallel)", n)
|
||||
}
|
||||
if n := l.FreeSlots("alpha"); n != 0 {
|
||||
t.Errorf("free slots = %d, want 0", n)
|
||||
}
|
||||
// A queueing request sees the host full: no waiting room, so it is refused.
|
||||
if _, _, err := l.Acquire(context.Background(), "alpha", "m"); err == nil {
|
||||
t.Error("Acquire on an over-tracked pair succeeded; want ErrQueueFull")
|
||||
}
|
||||
rel1()
|
||||
rel1() // idempotent
|
||||
rel2()
|
||||
rel3()
|
||||
if n := l.InFlight("alpha", "m"); n != 0 {
|
||||
t.Errorf("in flight after release = %d, want 0", n)
|
||||
}
|
||||
}
|
||||
|
||||
// A waiter gets a slot only once in flight is back under parallel: releasing a tracked request
|
||||
// while the pair is still over its limit must not hand the slot on.
|
||||
func TestTrackReleaseHandsOverOnlyUnderTheLimit(t *testing.T) {
|
||||
l := limiter.New()
|
||||
l.Configure("alpha", "m", 1, 1)
|
||||
relA := l.Track("alpha", "m")
|
||||
relB := l.Track("alpha", "m") // in flight 2, parallel 1
|
||||
|
||||
got := make(chan func(), 1)
|
||||
go func() {
|
||||
rel, _, err := l.Acquire(context.Background(), "alpha", "m")
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
close(got)
|
||||
return
|
||||
}
|
||||
got <- rel
|
||||
}()
|
||||
waitUntil(t, func() bool { return l.Queued("alpha", "m") == 1 })
|
||||
|
||||
relA() // in flight 1 == parallel: still no free slot
|
||||
select {
|
||||
case <-got:
|
||||
t.Fatal("waiter got a slot while in flight was still at parallel")
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
if n := l.InFlight("alpha", "m"); n != 1 {
|
||||
t.Fatalf("in flight = %d after one release, want 1", n)
|
||||
}
|
||||
|
||||
relB() // now the slot is free: hand it to the waiter
|
||||
select {
|
||||
case rel := <-got:
|
||||
if rel == nil {
|
||||
t.Fatal("waiter failed")
|
||||
}
|
||||
if n := l.InFlight("alpha", "m"); n != 1 {
|
||||
t.Errorf("in flight = %d with the waiter running, want 1", n)
|
||||
}
|
||||
rel()
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("waiter never got the freed slot")
|
||||
}
|
||||
if n := l.InFlight("alpha", "m"); n != 0 {
|
||||
t.Errorf("in flight at the end = %d, want 0", n)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,153 @@
|
||||
package proxy_test
|
||||
|
||||
// v2.3 task 02: affinity = "route" puts every request on the route (every conversation, every
|
||||
// control call) on one lease, so one host; queue = false counts the route's requests on the host
|
||||
// without ever holding or refusing them, because the client pins its own llama-server slot and
|
||||
// the server's queue is the one that must show it.
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
const affinityHosts = `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 0
|
||||
lease_idle = "30m"
|
||||
[hosts.alpha]
|
||||
base_url = %q
|
||||
weight = 1.0
|
||||
models = { "shared" = { parallel = 1 } }
|
||||
[hosts.beta]
|
||||
base_url = %q
|
||||
weight = 1.0
|
||||
models = { "shared" = { parallel = 1 } }
|
||||
[routes.r]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "shared"
|
||||
[routes.bm]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "shared"
|
||||
affinity = "route"
|
||||
queue = false
|
||||
[routes."agent-*"]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "shared"
|
||||
affinity = "route"
|
||||
`
|
||||
|
||||
func TestRouteAffinityPutsEverythingOnOneHost(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
|
||||
seen := map[string]int{}
|
||||
note := func(what string, resp *http.Response) {
|
||||
body := drain(resp)
|
||||
if resp.StatusCode != 200 {
|
||||
t.Fatalf("%s: %d %s", what, resp.StatusCode, body)
|
||||
}
|
||||
seen[resp.Header.Get("X-Crossbar-Host")]++
|
||||
}
|
||||
// Different conversations (different fingerprints), then control calls without any.
|
||||
for id := 1; id <= 4; id++ {
|
||||
note("chat", r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(id, 1)))
|
||||
}
|
||||
note("slots", r.do(http.MethodGet, "/bm/slots?model=shared", ""))
|
||||
note("props", r.do(http.MethodGet, "/bm/props?model=shared", ""))
|
||||
note("control", r.do(http.MethodPost, "/bm/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`))
|
||||
if len(seen) != 1 {
|
||||
t.Fatalf("route-affinity requests spread over %v, want one host", seen)
|
||||
}
|
||||
|
||||
// Templated concrete routes each get their own route lease, and each is internally sticky.
|
||||
for _, route := range []string{"agent-a", "agent-b", "agent-c"} {
|
||||
hosts := map[string]bool{}
|
||||
for id := 1; id <= 3; id++ {
|
||||
resp := r.do(http.MethodPost, "/"+route+"/v1/chat/completions", conversation(id, 1))
|
||||
drain(resp)
|
||||
hosts[resp.Header.Get("X-Crossbar-Host")] = true
|
||||
}
|
||||
if len(hosts) != 1 {
|
||||
t.Errorf("%s spread over %v, want one host", route, hosts)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueFalseNeitherHoldsNorRefuses(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
alpha.delay, beta.delay = 400*time.Millisecond, 400*time.Millisecond
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
|
||||
// parallel = 1 and queue_max = 0: a queueing route would refuse the second and third.
|
||||
var wg sync.WaitGroup
|
||||
codes := make(chan int, 3)
|
||||
start := time.Now()
|
||||
for id := 1; id <= 3; id++ {
|
||||
wg.Add(1)
|
||||
go func(id int) {
|
||||
defer wg.Done()
|
||||
resp := r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(id, 1))
|
||||
drain(resp)
|
||||
codes <- resp.StatusCode
|
||||
}(id)
|
||||
}
|
||||
// While they run, the host carries all three and a queueing route sees it full.
|
||||
var host string
|
||||
waitUntil(t, func() bool {
|
||||
for _, h := range []string{"alpha", "beta"} {
|
||||
if r.lim.InFlight(h, "shared") == 3 {
|
||||
host = h
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
})
|
||||
if n := r.lim.FreeSlots(host); n != 0 {
|
||||
t.Errorf("free slots on %s = %d while bm runs three, want 0", host, n)
|
||||
}
|
||||
wg.Wait()
|
||||
close(codes)
|
||||
for c := range codes {
|
||||
if c != 200 {
|
||||
t.Errorf("queue = false request: %d, want 200", c)
|
||||
}
|
||||
}
|
||||
// Concurrent, not serialised behind one slot: three 400 ms answers well under 1.2 s.
|
||||
if d := time.Since(start); d > 1100*time.Millisecond {
|
||||
t.Errorf("three queue = false requests took %v; they were held", d)
|
||||
}
|
||||
if n := r.lim.InFlight(host, "shared"); n != 0 {
|
||||
t.Errorf("in flight after = %d, want 0", n)
|
||||
}
|
||||
// Accounting is unchanged: each chat is still a row.
|
||||
waitUntil(t, func() bool { return r.rows("bm") == 3 })
|
||||
}
|
||||
|
||||
// The default is unchanged: two conversations on a conversation-affinity route may land on
|
||||
// different hosts (they start where there is most room).
|
||||
func TestConversationAffinityStillSpreads(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
alpha.delay, beta.delay = 300*time.Millisecond, 300*time.Millisecond
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
var wg sync.WaitGroup
|
||||
var mu sync.Mutex
|
||||
hosts := map[string]bool{}
|
||||
for id := 1; id <= 2; id++ {
|
||||
wg.Add(1)
|
||||
go func(id int) {
|
||||
defer wg.Done()
|
||||
resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(id, 1))
|
||||
drain(resp)
|
||||
mu.Lock()
|
||||
hosts[resp.Header.Get("X-Crossbar-Host")] = true
|
||||
mu.Unlock()
|
||||
}(id)
|
||||
time.Sleep(50 * time.Millisecond) // let the first take its slot so the second sees one host full
|
||||
}
|
||||
wg.Wait()
|
||||
if len(hosts) != 2 {
|
||||
t.Errorf("two concurrent conversations on route r used %v, want both hosts", hosts)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,176 @@
|
||||
package proxy_test
|
||||
|
||||
// v2.3 task 01: control-plane requests. A client that manages its own slots (Boxmaker) polls
|
||||
// /slots, reads /props, tokenizes and steers a running completion through
|
||||
// /v1/chat/completions/control. Those calls follow the route's lease like any other request but
|
||||
// must never wait for, or take, a slot: /control is sent while the client's own stream holds one.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// controlClient gives every control call a short deadline: a call that queues behind a full host
|
||||
// is the bug, and it must fail the test rather than hang it.
|
||||
var controlClient = &http.Client{Timeout: 2 * time.Second}
|
||||
|
||||
func (r *rig) do(method, path, body string) *http.Response {
|
||||
r.t.Helper()
|
||||
var rd *strings.Reader
|
||||
if body != "" {
|
||||
rd = strings.NewReader(body)
|
||||
}
|
||||
var req *http.Request
|
||||
var err error
|
||||
if rd != nil {
|
||||
req, err = http.NewRequest(method, r.front.URL+path, rd)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
} else {
|
||||
req, err = http.NewRequest(method, r.front.URL+path, nil)
|
||||
}
|
||||
if err != nil {
|
||||
r.t.Fatal(err)
|
||||
}
|
||||
resp, err := controlClient.Do(req)
|
||||
if err != nil {
|
||||
r.t.Fatalf("%s %s: %v", method, path, err)
|
||||
}
|
||||
return resp
|
||||
}
|
||||
|
||||
func (r *rig) rows(route string) int64 {
|
||||
r.t.Helper()
|
||||
counts, err := r.store.StatusCounts(time.Time{})
|
||||
if err != nil {
|
||||
r.t.Fatal(err)
|
||||
}
|
||||
var n int64
|
||||
for _, c := range counts {
|
||||
if c.Route == route {
|
||||
n += c.Count
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// A GET names its model in the query string: /slots?model=alpha-only must reach the host that
|
||||
// has alpha-only loaded, not whichever host the route's default model would pick.
|
||||
func TestGetModelComesFromTheQuery(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
|
||||
resp := r.do(http.MethodGet, "/r/slots?model=alpha-only", "")
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != "alpha" {
|
||||
t.Fatalf("GET /r/slots?model=alpha-only: %d on %q, want 200 on alpha", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"))
|
||||
}
|
||||
if got := alpha.lastReq(); got.method != "GET" || got.path != "/slots?model=alpha-only" {
|
||||
t.Errorf("alpha saw %s %s, want GET /slots?model=alpha-only", got.method, got.path)
|
||||
}
|
||||
// The same for beta-only, so a lucky default cannot pass the test.
|
||||
resp = r.do(http.MethodGet, "/r/slots?model=beta-only", "")
|
||||
drain(resp)
|
||||
if resp.Header.Get("X-Crossbar-Host") != "beta" {
|
||||
t.Errorf("GET /r/slots?model=beta-only went to %q, want beta", resp.Header.Get("X-Crossbar-Host"))
|
||||
}
|
||||
}
|
||||
|
||||
// /slots and /tokenize are proxied; the per-slot actions under /slots/ (save, restore, erase)
|
||||
// are not.
|
||||
func TestControlPathsAllowed(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
for _, tc := range []struct {
|
||||
method, path, body string
|
||||
want int
|
||||
}{
|
||||
{http.MethodGet, "/r/slots", "", 200},
|
||||
{http.MethodGet, "/r/slots?model=shared", "", 200},
|
||||
{http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`, 200},
|
||||
{http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`, 200},
|
||||
{http.MethodGet, "/r/slots/0", "", 404},
|
||||
{http.MethodPost, "/r/slots/0?action=erase", "", 404},
|
||||
{http.MethodPost, "/r/slots/0?action=save", `{"filename":"x"}`, 404},
|
||||
} {
|
||||
resp := r.do(tc.method, tc.path, tc.body)
|
||||
body := drain(resp)
|
||||
if resp.StatusCode != tc.want {
|
||||
t.Errorf("%s %s: %d %s, want %d", tc.method, tc.path, resp.StatusCode, body, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// With every slot on both hosts taken and the queue full, control-plane calls still go straight
|
||||
// through: no 503, no wait, no slot taken, no accounting row.
|
||||
func TestControlRequestsNeverTakeASlot(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
|
||||
// Take every "shared" slot (parallel 2 on each host) and the one queue place per host.
|
||||
var releases []func()
|
||||
for _, host := range []string{"alpha", "beta"} {
|
||||
for i := 0; i < 2; i++ {
|
||||
rel, _, err := r.lim.Acquire(context.Background(), host, "shared")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
releases = append(releases, rel)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go func() { _, _, _ = r.lim.Acquire(ctx, host, "shared") }()
|
||||
waitUntil(t, func() bool { return r.lim.Queued(host, "shared") == 1 })
|
||||
}
|
||||
defer func() {
|
||||
for _, rel := range releases {
|
||||
rel()
|
||||
}
|
||||
}()
|
||||
|
||||
for _, tc := range []struct{ method, path, body string }{
|
||||
{http.MethodGet, "/r/slots?model=shared", ""},
|
||||
{http.MethodGet, "/r/props?model=shared", ""},
|
||||
{http.MethodHead, "/r/props?model=shared", ""},
|
||||
{http.MethodGet, "/r/v1/models", ""},
|
||||
{http.MethodPost, "/r/tokenize", `{"model":"shared","content":"hello"}`},
|
||||
{http.MethodPost, "/r/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`},
|
||||
} {
|
||||
resp := r.do(tc.method, tc.path, tc.body)
|
||||
body := drain(resp)
|
||||
if resp.StatusCode != 200 {
|
||||
t.Errorf("%s %s with the host full: %d %s, want 200", tc.method, tc.path, resp.StatusCode, body)
|
||||
}
|
||||
}
|
||||
for _, host := range []string{"alpha", "beta"} {
|
||||
if n := r.lim.InFlight(host, "shared"); n != 2 {
|
||||
t.Errorf("%s in flight = %d after control calls, want 2 (control takes no slot)", host, n)
|
||||
}
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond) // a row is written after the answer; give a stray one time to land
|
||||
if n := r.rows("r"); n != 0 {
|
||||
t.Errorf("control calls wrote %d accounting rows, want 0", n)
|
||||
}
|
||||
|
||||
// A chat completion on the same full route still queues or is refused as before: the bypass
|
||||
// is for control calls only.
|
||||
resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1))
|
||||
drain(resp)
|
||||
if resp.StatusCode != http.StatusServiceUnavailable {
|
||||
t.Errorf("chat on a full route: %d, want 503 (queue full)", resp.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
// A chat completion is not a control call just because its path starts the same way.
|
||||
func TestChatIsNotControl(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
resp := r.do(http.MethodPost, "/r/v1/chat/completions", conversation(1, 1))
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 {
|
||||
t.Fatalf("chat: %d", resp.StatusCode)
|
||||
}
|
||||
waitUntil(t, func() bool { return r.rows("r") == 1 }) // the row lands just after the answer
|
||||
}
|
||||
@@ -0,0 +1,164 @@
|
||||
package proxy_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
|
||||
)
|
||||
|
||||
// ctxUpstream is a fake router that reports a context size in /props and echoes completions.
|
||||
func ctxUpstream(t *testing.T, name string, nCtx, slots int) *upstream {
|
||||
u := &upstream{name: name}
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"status":"ok"}`) })
|
||||
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"data":[{"id":"shared"}]}`) })
|
||||
mux.HandleFunc("/props", func(w http.ResponseWriter, r *http.Request) {
|
||||
fmt.Fprintf(w, `{"default_generation_settings":{"n_ctx":%d},"total_slots":%d}`, nCtx, slots)
|
||||
})
|
||||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||||
u.hits.Add(1)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
fmt.Fprint(w, `{"choices":[{"message":{"role":"assistant","content":"ok"}}],"usage":{"prompt_tokens":1,"completion_tokens":1}}`)
|
||||
})
|
||||
u.srv = httptest.NewServer(mux)
|
||||
t.Cleanup(u.srv.Close)
|
||||
return u
|
||||
}
|
||||
|
||||
const ctxHosts = `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 2
|
||||
[hosts.small]
|
||||
base_url = %q
|
||||
weight = 10.0
|
||||
models = { "shared" = { parallel = 2 } }
|
||||
[hosts.big]
|
||||
base_url = %q
|
||||
weight = 1.0
|
||||
models = { "shared" = { parallel = 1 } }
|
||||
[routes.r]
|
||||
hosts = ["small", "big"]
|
||||
default_model = "shared"
|
||||
`
|
||||
|
||||
// bodyOfTokens builds a chat body whose byte size implies roughly n tokens under the guard's
|
||||
// estimate (bytes/4 × 1.2): n tokens ≈ 3.33 n bytes ≈ 2n/3 five-byte words.
|
||||
func bodyOfTokens(n int) string {
|
||||
text := strings.Repeat("word ", n*2/3)
|
||||
return fmt.Sprintf(`{"model":"shared","stream":false,"messages":[{"role":"user","content":"%s"}]}`, text)
|
||||
}
|
||||
|
||||
func TestOversizedPromptMovesToAHostWhereItFits(t *testing.T) {
|
||||
small := ctxUpstream(t, "small", 8192, 2) // 4096 per slot
|
||||
big := ctxUpstream(t, "big", 131072, 1) // 131072 per slot
|
||||
r := newRig(t, ctxHosts, small, big)
|
||||
// A small prompt starts on `small` (weight 10).
|
||||
resp := r.post("/r/v1/chat/completions", bodyOfTokens(100))
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != "small" {
|
||||
t.Fatalf("small prompt went to %q, want small", resp.Header.Get(proxy.HostHeader))
|
||||
}
|
||||
// A new conversation with ~10k tokens does not fit small's 4096-token slot: it must be
|
||||
// placed on big, with the reason visible in a header.
|
||||
resp = r.post("/r/v1/chat/completions", bodyOfTokens(10000))
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" {
|
||||
t.Fatalf("oversized prompt: %d from %q, want 200 from big", resp.StatusCode, resp.Header.Get(proxy.HostHeader))
|
||||
}
|
||||
if got := resp.Header.Get(proxy.CtxHeader); !strings.HasPrefix(got, "moved") {
|
||||
t.Errorf("%s = %q, want moved:… ", proxy.CtxHeader, got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOversizedPromptWithNoFitIs400(t *testing.T) {
|
||||
small := ctxUpstream(t, "small", 8192, 2)
|
||||
tiny := ctxUpstream(t, "big", 4096, 2) // also too small
|
||||
r := newRig(t, ctxHosts, small, tiny)
|
||||
resp := r.post("/r/v1/chat/completions", bodyOfTokens(10000))
|
||||
body := drain(resp)
|
||||
if resp.StatusCode != http.StatusBadRequest {
|
||||
t.Fatalf("status %d body %s, want 400", resp.StatusCode, body)
|
||||
}
|
||||
// v2.3: llama-server's own shape for this error, so a client handles crossbar's refusal the
|
||||
// way it handles the server's (Boxmaker keys on error.type; the error JSON must come first).
|
||||
if !strings.HasPrefix(body, `{"error":`) {
|
||||
t.Errorf("body must start with the error object: %s", body)
|
||||
}
|
||||
var e struct {
|
||||
Error struct {
|
||||
Code int `json:"code"`
|
||||
Type string `json:"type"`
|
||||
Message string `json:"message"`
|
||||
NPromptTokens float64 `json:"n_prompt_tokens"`
|
||||
NCtx float64 `json:"n_ctx"`
|
||||
} `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(body), &e); err != nil || e.Error.Code != 400 || e.Error.Type != "exceed_context_size_error" || e.Error.Message != "prompt too large" {
|
||||
t.Fatalf("body = %s, want {\"error\":{\"code\":400,\"type\":\"exceed_context_size_error\",\"message\":\"prompt too large\",…}}", body)
|
||||
}
|
||||
if est := e.Error.NPromptTokens; est < 8000 || est > 13000 {
|
||||
t.Errorf("n_prompt_tokens = %v, want roughly 10000 tokens", est)
|
||||
}
|
||||
if max := e.Error.NCtx; max != 4096 {
|
||||
t.Errorf("n_ctx = %v, want the largest per-slot context among the route's hosts (4096)", max)
|
||||
}
|
||||
if ct := resp.Header.Get("Content-Type"); !strings.HasPrefix(ct, "application/json") {
|
||||
t.Errorf("Content-Type = %q, want application/json", ct)
|
||||
}
|
||||
if small.hits.Load()+tiny.hits.Load() != 0 {
|
||||
t.Errorf("a refused prompt must not reach any upstream")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnknownContextNeverBlocks(t *testing.T) {
|
||||
// /props missing on both hosts: NCtx 0 means "unknown", and the guard must stay out of the way.
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, alpha, beta)
|
||||
resp := r.post("/r/v1/chat/completions", bodyOfTokens(50000))
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get(proxy.CtxHeader) != "" {
|
||||
t.Errorf("unknown context: %d %q, want 200 and no ctx header", resp.StatusCode, resp.Header.Get(proxy.CtxHeader))
|
||||
}
|
||||
}
|
||||
|
||||
// grow appends later turns to a conversation body without touching its system prompt or first
|
||||
// user message, so the fingerprint — and therefore the lease — stays the same.
|
||||
func grow(body string, words int) string {
|
||||
turn := `,{"role":"assistant","content":"ok"},{"role":"user","content":"` + strings.Repeat("x ", words) + `"}`
|
||||
return strings.Replace(body, `]}`, turn+`]}`, 1)
|
||||
}
|
||||
|
||||
func TestStickyLeaseSurvivesGrowthUntilItDoesNotFit(t *testing.T) {
|
||||
small := ctxUpstream(t, "small", 8192, 2)
|
||||
big := ctxUpstream(t, "big", 131072, 1)
|
||||
r := newRig(t, ctxHosts, small, big)
|
||||
body := bodyOfTokens(100)
|
||||
resp := r.post("/r/v1/chat/completions", body)
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != "small" {
|
||||
t.Fatal("setup: first turn must be on small")
|
||||
}
|
||||
// Same conversation, a later turn well under 4096 tokens: stays.
|
||||
resp = r.post("/r/v1/chat/completions", grow(body, 500))
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != "small" || resp.Header.Get(proxy.LeaseHeader) != "reused" {
|
||||
t.Errorf("turn 2: %q %q, want small reused", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader))
|
||||
}
|
||||
// A turn that outgrows the slot moves the lease — once — and the move is visible in the header.
|
||||
huge := grow(body, 30000)
|
||||
resp = r.post("/r/v1/chat/completions", huge)
|
||||
drain(resp)
|
||||
if resp.StatusCode != 200 || resp.Header.Get(proxy.HostHeader) != "big" || !strings.HasPrefix(resp.Header.Get(proxy.CtxHeader), "moved") {
|
||||
t.Fatalf("outgrown turn: %d %q ctx=%q, want 200 from big with a moved header", resp.StatusCode, resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.CtxHeader))
|
||||
}
|
||||
resp = r.post("/r/v1/chat/completions", huge)
|
||||
drain(resp)
|
||||
if resp.Header.Get(proxy.HostHeader) != "big" || resp.Header.Get(proxy.LeaseHeader) != "reused" {
|
||||
t.Errorf("after the move the lease is on big: %q %q", resp.Header.Get(proxy.HostHeader), resp.Header.Get(proxy.LeaseHeader))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
package proxy_test
|
||||
|
||||
// v2.3 task 03: Handler.ForRoute serves one route with unprefixed paths, for a route's dedicated
|
||||
// listener.
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
|
||||
)
|
||||
|
||||
// dedicated serves r's route on its own test server, sharing r's health, leases, limiter and
|
||||
// store, as main does for a route with listen set.
|
||||
func dedicated(t *testing.T, r *rig, route string) *httptest.Server {
|
||||
p := proxy.New(r.cfg, r.health, r.leases, r.lim, r.store, nil)
|
||||
srv := httptest.NewServer(p.ForRoute(route))
|
||||
t.Cleanup(srv.Close)
|
||||
return srv
|
||||
}
|
||||
|
||||
func call(t *testing.T, method, url, body string, hdr ...string) (*http.Response, string) {
|
||||
t.Helper()
|
||||
var req *http.Request
|
||||
if body != "" {
|
||||
req, _ = http.NewRequest(method, url, strings.NewReader(body))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
} else {
|
||||
req, _ = http.NewRequest(method, url, nil)
|
||||
}
|
||||
for i := 0; i+1 < len(hdr); i += 2 {
|
||||
req.Header.Set(hdr[i], hdr[i+1])
|
||||
}
|
||||
resp, err := controlClient.Do(req)
|
||||
if err != nil {
|
||||
t.Fatalf("%s %s: %v", method, url, err)
|
||||
}
|
||||
return resp, drain(resp)
|
||||
}
|
||||
|
||||
func TestForRouteServesUnprefixedPaths(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
srv := dedicated(t, r, "bm")
|
||||
|
||||
resp, body := call(t, http.MethodPost, srv.URL+"/v1/chat/completions", conversation(1, 1))
|
||||
if resp.StatusCode != 200 {
|
||||
t.Fatalf("chat on the dedicated listener: %d %s", resp.StatusCode, body)
|
||||
}
|
||||
host := resp.Header.Get("X-Crossbar-Host")
|
||||
up := map[string]*upstream{"alpha": alpha, "beta": beta}[host]
|
||||
if up == nil || up.lastReq().path != "/v1/chat/completions" {
|
||||
t.Fatalf("upstream %q saw %+v, want /v1/chat/completions unchanged", host, up.lastReq())
|
||||
}
|
||||
resp, _ = call(t, http.MethodGet, srv.URL+"/slots?model=shared", "")
|
||||
if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host || up.lastReq().path != "/slots?model=shared" {
|
||||
t.Errorf("/slots: %d on %q (last %+v), want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), up.lastReq(), host)
|
||||
}
|
||||
resp, _ = call(t, http.MethodPost, srv.URL+"/v1/chat/completions/control", `{"id":"chatcmpl-1","action":"reasoning_end","model":"shared"}`)
|
||||
if resp.StatusCode != 200 || resp.Header.Get("X-Crossbar-Host") != host {
|
||||
t.Errorf("/control: %d on %q, want 200 on %q", resp.StatusCode, resp.Header.Get("X-Crossbar-Host"), host)
|
||||
}
|
||||
// The chat is accounted to the route the listener serves.
|
||||
waitUntil(t, func() bool { return r.rows("bm") == 1 })
|
||||
// The same route through the main listener shares the lease: same host.
|
||||
resp = r.do(http.MethodPost, "/bm/v1/chat/completions", conversation(2, 1))
|
||||
drain(resp)
|
||||
if resp.Header.Get("X-Crossbar-Host") != host {
|
||||
t.Errorf("main listener /bm went to %q, dedicated to %q; one route, one lease", resp.Header.Get("X-Crossbar-Host"), host)
|
||||
}
|
||||
}
|
||||
|
||||
func TestForRouteRefusals(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
srv := dedicated(t, r, "bm")
|
||||
for _, tc := range []struct {
|
||||
name, method, path string
|
||||
hdr []string
|
||||
want int
|
||||
msg string
|
||||
}{
|
||||
{"a prefixed path is not stripped", http.MethodGet, "/bm/v1/models", nil, 404, "not found"},
|
||||
{"no admin here", http.MethodGet, "/_crossbar/hosts", nil, 404, "not found"},
|
||||
{"root", http.MethodGet, "/", nil, 404, "not found"},
|
||||
{"header naming another route", http.MethodGet, "/v1/models", []string{"X-Crossbar-Route", "r"}, 400, "conflicting route"},
|
||||
} {
|
||||
resp, body := call(t, tc.method, srv.URL+tc.path, "", tc.hdr...)
|
||||
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.name, resp.StatusCode, body, tc.want, tc.msg)
|
||||
}
|
||||
}
|
||||
// A header naming this same route is harmless.
|
||||
resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "", "X-Crossbar-Route", "bm")
|
||||
if resp.StatusCode != 200 {
|
||||
t.Errorf("header naming the listener's own route: %d %s, want 200", resp.StatusCode, body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestForRouteUnknownRoute(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
r := newRig(t, affinityHosts, alpha, beta)
|
||||
srv := dedicated(t, r, "nope")
|
||||
resp, body := call(t, http.MethodGet, srv.URL+"/v1/models", "")
|
||||
if resp.StatusCode != 404 || !strings.Contains(body, "unknown route") {
|
||||
t.Errorf("ForRoute(unknown): %d %s, want 404 unknown route", resp.StatusCode, body)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,292 @@
|
||||
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).
|
||||
// The rig, the fake upstream and the request helpers live in helpers_test.go.
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/proxy"
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
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())
|
||||
}
|
||||
}
|
||||
|
||||
// spreadHosts: beta is preferred (weight 10) until both of its "shared" slots are busy; then
|
||||
// alpha (2 free × 1) beats beta (0 free × 10), and a new conversation must start on alpha.
|
||||
const spreadHosts = `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 4
|
||||
lease_idle = "30m"
|
||||
[hosts.alpha]
|
||||
base_url = %q
|
||||
weight = 1.0
|
||||
models = { "shared" = { parallel = 2 } }
|
||||
[hosts.beta]
|
||||
base_url = %q
|
||||
weight = 10.0
|
||||
models = { "shared" = { parallel = 2 } }
|
||||
[routes.r]
|
||||
hosts = ["alpha", "beta"]
|
||||
default_model = "shared"
|
||||
`
|
||||
|
||||
func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) {
|
||||
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
||||
beta.delay = 400 * time.Millisecond
|
||||
r := newRig(t, spreadHosts, 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)
|
||||
// arrive one after the other so both pick beta (10 > 2): wait until beta holds i slots
|
||||
waitUntil(t, func() bool { return r.lim.InFlight("beta", "shared") == i })
|
||||
}
|
||||
// …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()
|
||||
if beta.hits.Load() != 2 || alpha.hits.Load() != 1 {
|
||||
t.Errorf("hits beta=%d alpha=%d, want 2 and 1", beta.hits.Load(), alpha.hits.Load())
|
||||
}
|
||||
}
|
||||
|
||||
// waitUntil polls cond every 5 ms for up to two seconds and fails the test if it never holds.
|
||||
func waitUntil(t *testing.T, cond func() bool) {
|
||||
t.Helper()
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if cond() {
|
||||
return
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("condition not reached within two seconds")
|
||||
}
|
||||
|
||||
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)
|
||||
fire := func(i int) {
|
||||
go func() {
|
||||
resp := r.post("/r/v1/chat/completions", conversation(i, 1))
|
||||
drain(resp)
|
||||
codes <- resp.StatusCode
|
||||
}()
|
||||
}
|
||||
// Arrival order is enforced by watching the limiter, not by sleeping: 1 runs, 2 queues,
|
||||
// 3 finds the queue full.
|
||||
fire(1)
|
||||
waitUntil(t, func() bool { return r.lim.InFlight("alpha", "shared") == 1 })
|
||||
fire(2)
|
||||
waitUntil(t, func() bool { return r.lim.Queued("alpha", "shared") == 1 })
|
||||
fire(3)
|
||||
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 are written after each response completes; allow the store a moment to catch up.
|
||||
var rows []store.UsageRow
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
rows, _ = r.store.Usage(time.Time{}, store.ByRoute)
|
||||
if len(rows) == 1 && rows[0].Requests == 3 {
|
||||
break
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
if len(rows) != 1 || rows[0].Requests != 3 || rows[0].Errors != 1 {
|
||||
t.Fatalf("usage = %+v, want 3 requests, 1 error (the 503 is recorded too)", rows)
|
||||
}
|
||||
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"},
|
||||
// v2.3: /slots itself is proxied (a control-plane path); its per-slot actions are not.
|
||||
{http.MethodGet, "/r/slots/0", 404, "not found"},
|
||||
{http.MethodGet, "/r/metrics", 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())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user