181 lines
6.0 KiB
Go
181 lines
6.0 KiB
Go
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 (
|
|
"bufio"
|
|
"net/http"
|
|
"strings"
|
|
"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)
|
|
}
|
|
// The slot is given back just after the answer is sent (a deferred release), so wait for it.
|
|
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
|
|
// 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)
|
|
}
|
|
}
|
|
|
|
// A request counts against its host for as long as its answer is streaming, not only until the
|
|
// first byte: a slot (queueing route) or a tracked place (queue = false) is given back when the
|
|
// stream ends.
|
|
func TestLoadIsHeldForTheWholeStream(t *testing.T) {
|
|
for _, route := range []string{"r", "bm"} {
|
|
t.Run(route, func(t *testing.T) {
|
|
alpha, beta := newUpstream(t, "alpha"), newUpstream(t, "beta")
|
|
r := newRig(t, affinityHosts, alpha, beta)
|
|
body := strings.Replace(conversation(1, 1), `"stream":false`, `"stream":true`, 1)
|
|
resp := r.do(http.MethodPost, "/"+route+"/v1/chat/completions", body)
|
|
defer resp.Body.Close()
|
|
host := resp.Header.Get("X-Crossbar-Host")
|
|
line, err := bufio.NewReader(resp.Body).ReadString('\n')
|
|
if err != nil || !strings.HasPrefix(line, "data:") {
|
|
t.Fatalf("first line %q, err %v", line, err)
|
|
}
|
|
// The first chunk is here; the upstream sends more for another ~30 ms.
|
|
if n := r.lim.InFlight(host, "shared"); n != 1 {
|
|
t.Errorf("in flight on %s after the first chunk = %d, want 1 (released before the stream ended)", host, n)
|
|
}
|
|
drain(resp)
|
|
waitUntil(t, func() bool { return r.lim.InFlight(host, "shared") == 0 })
|
|
})
|
|
}
|
|
}
|