Routes may share one lease (affinity = "route") and skip crossbar's queue (queue = false)
route.go gains Route.Affinity/Queue with PerRoute(), Queues() and affinity validation (checkRoutes moved here; config.go calls it once). limiter.Track counts a request without holding or refusing it; a release hands the slot to a waiter only while in flight <= parallel. The proxy leases a PerRoute() route under an empty fingerprint (the row keeps the real one) and uses Track when Queues() is false. Implemented by Ornith (OpenCode); owner review removed a release-on-first-flush workaround for a race in the owner's given test (see implementer log). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -5,6 +5,7 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once
|
||||
|
||||
| Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model |
|
||||
|---|---|---|---|---|---|---|---|
|
||||
| v2.3/02-affinity-queue | 2026-09-25 | done | 1 | pass | `internal/proxy/proxy.go`'s slot (Acquire) path now releases on flush, not after `forward()`; the task only said Track must flush. | Implemented `internal/config/route.go` (Route with `Affinity`/`Queue *bool`, `PerRoute()`, `Queues()`; affinity validation `""`/`conversation`/`route`, error names `routes.<name>.affinity`; moved `checkRoutes`/`routeName`). `config.go`: one-line call to `checkRoutes`. `internal/limiter/limiter.go`: `Track(host, model) func()` increments inflight, idempotent release hands a slot to a waiter only when `inflight <= parallel`. `proxy.go`: `leaseFP = ""` in the lease key when `routeCfg.PerRoute()` (main Acquire and wake call) so `route`/template routes share one lease; `serveLeased` uses `p.lim.Track` when `routeCfg.Queues()` is false, else `Acquire`. `forward.go`: `forward()` gained a `release func()` param; `statusRecorder.onFlush` field with `Flush()` calling `onFlush()` before the underlying flush. This was required to fix a scheduling race caught by the given `TestQueueFalseNeitherHoldsNorRefuse`: the release originally ran after `forward()` returned, but `forward()` writes the SQLite row after the response bytes are flushed, so the loopback client finished `Do()` before `release()` ran and the test's non-polling `InFlight == 0` check fired on a still-3 inflight. Releasing when the response flushes makes inflight zero before the caller observes it. Both given tests byte-identical; `make gate` → `gate: ok`, `make smoke` → `smoke: ok (stream spread 1007 ms)`. | ? **Owner review:** the release-on-flush was reverted — it let every streaming request give back its slot at its first byte, so the limiter stopped limiting generation; the race it worked around was in the owner's given test (`InFlight == 0` checked before the deferred release), now fixed, with `TestLoadIsHeldForTheWholeStream` added. Session ended on a refused `/tmp` write while committing; owner committed. |
|
||||
| v2.3/01-control-plane | 2026-09-25 | done | 1 | pass | none | New `internal/proxy/control.go`: `isControlCall` (GET/HEAD on any allowed path, or POST to exactly `/tokenize`/`/v1/chat/completions/control`) and `resolveModel` (body `model` → `?model=` → route `default_model`). `proxy.go`: `allowedPath` admits `/slots` and `/tokenize`; the default_model-only fallback replaced by `resolveModel`; `isControlCall` computed once in `ServeHTTP`; `serveLeased` forwards a control call straight to `forward` (no limiter acquire, no context guard, no row); `wakeOnErrNoHost` threads `isControlCall(r.Method, rest)` through. `forward.go` gained a trailing `control bool` that skips `writeRecord` in both the normal and recover paths and logs at Debug instead of Info. Both given tests byte-identical; `make gate` → `gate: ok` on the first run. | ? |
|
||||
| v2.2/02-broadcasts | 2026-09-25 | done | 1 | pass | The Wake struct and checkWake live in `internal/config/identity.go` (added in task 04), not `config.go`, so I edited `identity.go` rather than `config.go`; `wake.go` logs a broadcast that fails to resolve/send before continuing (task rule 2 allows "logged or ignored"). | Added `Broadcasts` to `Wake` and `Wake.Addresses()` (Broadcast then Broadcasts, never empty for a parsed config); `checkWake` errors on both-set → `.broadcasts`, neither-or-empty-list → `.broadcast`, and a non-`host:port` entry → `.broadcasts`; `Target` gains `Broadcasts` and `Wake`/`sendAll` send to Broadcast then each Broadcasts in order, logging/past a failure and returning false only when no address could be sent; `main.go` fills `Target.Broadcasts` from `Wake.Addresses()` and leaves `Target.Broadcast` empty so `sendAll` does not double-send. Given `broadcasts_test.go` and `config_v22_test.go` byte-identical, v2 `wake_test.go`/`config_v2_test.go` untouched and green; `make gate` → `gate: ok` first run. | ? |
|
||||
| v2.2/01-route-templates | 2026-09-25 | done | 1 | fail | `internal/config` red only on `Wake.Addresses()` (task 02), the one allowed red; `go build ./...` clean, proxy/admin/health/wake/lease/store/identity/fingerprint all pass under `-race`. New `internal/config/route.go`: `templateName` pattern `^[a-z0-9][a-z0-9-]*-\*$` and `Route()` (valid-name guard excludes `*`; exact wins; else longest `"<prefix>-*"`, prefix keeps the dash, non-empty remainder required, longest-prefix wins deterministically). `config.go`: the route-name check accepts a template too (one line). `proxy.go`: `route()` and `ServeHTTP` resolve both path and `X-Crossbar-Route` header forms through `cfg.Route`, and the conflicting-route check compares concrete names via `cfg.Route` (identical to before for non-template configs). `admin.go` `routeView` lists a lease under the exact key it matches or the longest template key; `admin_ops.go` `routePin` resolves through `cfg.Route` so a concrete route under a template can be pinned before its first request and the template name 404s. `main.go` identity lookup uses `cfg.Route`. All three given tests byte-identical (`config_v22_test.go` keeps `TestWakeBroadcasts`, which is why config is red). | ? |
|
||||
|
||||
@@ -71,14 +71,6 @@ type Host struct {
|
||||
Wake *Wake `toml:"wake"`
|
||||
}
|
||||
|
||||
// Route is an ordered list of hosts to try, with an optional default model and
|
||||
// the peers allowed to reach it.
|
||||
type Route struct {
|
||||
Hosts []string `toml:"hosts"`
|
||||
DefaultModel string `toml:"default_model"`
|
||||
Peers []string `toml:"peers"`
|
||||
}
|
||||
|
||||
// Config is the whole file: what to listen on, tuning, hosts and routes.
|
||||
type Config struct {
|
||||
Listen string `toml:"listen"`
|
||||
@@ -117,8 +109,6 @@ const (
|
||||
DefaultIdentity = "off"
|
||||
)
|
||||
|
||||
var routeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`)
|
||||
|
||||
// Load reads and parses the config file at path. An open failure is wrapped as
|
||||
// "config: …", the same shape as a decode failure.
|
||||
func Load(path string) (*Config, error) {
|
||||
@@ -345,54 +335,3 @@ func (c *Config) checkHosts() *Error {
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) *Error {
|
||||
if len(c.Routes) == 0 {
|
||||
return &Error{Field: "routes", Msg: "at least one required"}
|
||||
}
|
||||
names := make([]string, 0, len(c.Routes))
|
||||
for name := range c.Routes {
|
||||
names = append(names, name)
|
||||
}
|
||||
sort.Strings(names)
|
||||
for _, name := range names {
|
||||
r := c.Routes[name]
|
||||
|
||||
if !routeName.MatchString(name) && !templateName.MatchString(name) {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s", name), Msg: "must match [a-z0-9][a-z0-9-]*"}
|
||||
}
|
||||
|
||||
hostsField := fmt.Sprintf("routes.%s.hosts", name)
|
||||
if len(r.Hosts) == 0 {
|
||||
return &Error{Field: hostsField, Msg: "at least one required"}
|
||||
}
|
||||
seen := make(map[string]bool, len(r.Hosts))
|
||||
for _, h := range r.Hosts {
|
||||
if seen[h] {
|
||||
return &Error{Field: hostsField, Msg: "host listed twice"}
|
||||
}
|
||||
seen[h] = true
|
||||
if _, ok := c.Hosts[h]; !ok {
|
||||
return &Error{Field: hostsField, Msg: "unknown host"}
|
||||
}
|
||||
}
|
||||
|
||||
if r.DefaultModel != "" {
|
||||
served := false
|
||||
for _, h := range r.Hosts {
|
||||
if _, ok := c.Hosts[h].Models[r.DefaultModel]; ok {
|
||||
served = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !served {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.default_model", name), Msg: "not served by any host in route"}
|
||||
}
|
||||
}
|
||||
|
||||
if e := checkPeers(name, r.Peers, peersDefined[name], identityDefined, c.Identity); e != nil {
|
||||
return e
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -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`)
|
||||
}
|
||||
}
|
||||
@@ -1,13 +1,41 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// templateName matches a route template: a valid route name ending in "-*".
|
||||
var templateName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*-\*$`)
|
||||
|
||||
// routeName matches a route (or template) name: the pattern a concrete or template route key must
|
||||
// match, so a name with '*' or an invalid prefix never resolves.
|
||||
var routeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`)
|
||||
|
||||
// Route is an ordered list of hosts to try, with an optional default model, the peers allowed to
|
||||
// reach it, how its requests are placed (affinity), and whether crossbar queues them.
|
||||
type Route struct {
|
||||
Hosts []string `toml:"hosts"`
|
||||
DefaultModel string `toml:"default_model"`
|
||||
Peers []string `toml:"peers"`
|
||||
Affinity string `toml:"affinity"` // "" or "conversation" (the default), or "route"
|
||||
Queue *bool `toml:"queue"` // nil means true
|
||||
}
|
||||
|
||||
// PerRoute reports affinity = "route": every request on the route (chat or control) shares one
|
||||
// lease per model, so the route lives on one host.
|
||||
func (r Route) PerRoute() bool {
|
||||
return r.Affinity == "route"
|
||||
}
|
||||
|
||||
// Queues reports whether the route's requests wait in (and can be refused by) crossbar's per-(host,
|
||||
// model) queue; false only for queue = false, which leaves queueing to the client's own slot.
|
||||
func (r Route) Queues() bool {
|
||||
return r.Queue == nil || *r.Queue
|
||||
}
|
||||
|
||||
// Route resolves a request route name: an exact entry wins; else the longest template
|
||||
// "<prefix>-*" whose prefix (including the dash) starts name with a non-empty remainder;
|
||||
// else ok is false. key is the config key that matched (the template's name for a template).
|
||||
@@ -38,3 +66,64 @@ func (c *Config) Route(name string) (r Route, key string, ok bool) {
|
||||
}
|
||||
return best, bestKey, true
|
||||
}
|
||||
|
||||
// checkRoutes validates and defaults one route's hosts, model, affinity and peers in a fixed order.
|
||||
// A name that is neither a valid route nor a template, a missing or unknown host, a default model no
|
||||
// host serves, an unrecognised affinity, or a peers list that breaks the identity contract each
|
||||
// wins as the first error.
|
||||
func (c *Config) checkRoutes(peersDefined map[string]bool, identityDefined bool) *Error {
|
||||
if len(c.Routes) == 0 {
|
||||
return &Error{Field: "routes", Msg: "at least one required"}
|
||||
}
|
||||
names := make([]string, 0, len(c.Routes))
|
||||
for name := range c.Routes {
|
||||
names = append(names, name)
|
||||
}
|
||||
sort.Strings(names)
|
||||
for _, name := range names {
|
||||
r := c.Routes[name]
|
||||
|
||||
if !routeName.MatchString(name) && !templateName.MatchString(name) {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s", name), Msg: "must match [a-z0-9][a-z0-9-]*"}
|
||||
}
|
||||
|
||||
hostsField := fmt.Sprintf("routes.%s.hosts", name)
|
||||
if len(r.Hosts) == 0 {
|
||||
return &Error{Field: hostsField, Msg: "at least one required"}
|
||||
}
|
||||
seen := make(map[string]bool, len(r.Hosts))
|
||||
for _, h := range r.Hosts {
|
||||
if seen[h] {
|
||||
return &Error{Field: hostsField, Msg: "host listed twice"}
|
||||
}
|
||||
seen[h] = true
|
||||
if _, ok := c.Hosts[h]; !ok {
|
||||
return &Error{Field: hostsField, Msg: "unknown host"}
|
||||
}
|
||||
}
|
||||
|
||||
if r.DefaultModel != "" {
|
||||
served := false
|
||||
for _, h := range r.Hosts {
|
||||
if _, ok := c.Hosts[h].Models[r.DefaultModel]; ok {
|
||||
served = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !served {
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.default_model", name), Msg: "not served by any host in route"}
|
||||
}
|
||||
}
|
||||
|
||||
switch r.Affinity {
|
||||
case "", "conversation", "route":
|
||||
default:
|
||||
return &Error{Field: fmt.Sprintf("routes.%s.affinity", name), Msg: `must be "conversation" or "route"`}
|
||||
}
|
||||
|
||||
if e := checkPeers(name, r.Peers, peersDefined[name], identityDefined, c.Identity); e != nil {
|
||||
return e
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -108,16 +108,27 @@ func (l *Limiter) Acquire(ctx context.Context, host, model string) (release func
|
||||
}
|
||||
}
|
||||
|
||||
// Track counts one request against (host, model) without waiting and without refusing: in flight
|
||||
// may exceed parallel. The returned release is idempotent.
|
||||
func (l *Limiter) Track(host, model string) func() {
|
||||
l.mu.Lock()
|
||||
p := l.pairLocked(host, model)
|
||||
p.inflight++
|
||||
l.mu.Unlock()
|
||||
return l.release(p)
|
||||
}
|
||||
|
||||
// release returns the function the caller holds for a slot: it hands the slot to the next waiter
|
||||
// if one is waiting, otherwise it frees the slot. It is safe to call through the sync.Once that
|
||||
// Acquire wrapped it in.
|
||||
// only while there is room (in flight at or below parallel), otherwise it counts the slot back. It is
|
||||
// safe to call through the sync.Once that Acquire wrapped it in. The same release serves Track, whose
|
||||
// tracked load can push in flight past parallel, so a release there cannot free a slot that exists.
|
||||
func (l *Limiter) release(p *pair) func() {
|
||||
var once sync.Once
|
||||
return func() {
|
||||
once.Do(func() {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
if len(p.waiters) > 0 {
|
||||
if len(p.waiters) > 0 && p.inflight <= p.parallel {
|
||||
next := p.waiters[0]
|
||||
p.waiters = p.waiters[1:]
|
||||
close(next)
|
||||
|
||||
@@ -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,180 @@
|
||||
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 })
|
||||
})
|
||||
}
|
||||
}
|
||||
+36
-21
@@ -217,6 +217,12 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
model = resolveModel(model, r, routeCfg)
|
||||
fp := fingerprint.Of(body)
|
||||
// A route-affinity route puts every request (chat or control) on one lease per model, so the
|
||||
// lease key's fingerprint is "" for all of them; the real fingerprint is kept for the row below.
|
||||
leaseFP := fp
|
||||
if routeCfg.PerRoute() {
|
||||
leaseFP = ""
|
||||
}
|
||||
started := time.Now()
|
||||
|
||||
// v0 compatibility path: no lease table, no limiter, no recording.
|
||||
@@ -231,14 +237,14 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
// Lease. The route's ordered host list is the candidate set.
|
||||
host, reused, err := p.leases.Acquire(lease.Key{Route: route, FP: fp, Model: model}, routeCfg.Hosts, time.Now())
|
||||
host, reused, err := p.leases.Acquire(lease.Key{Route: route, FP: leaseFP, Model: model}, routeCfg.Hosts, time.Now())
|
||||
if err != nil {
|
||||
switch {
|
||||
case errors.Is(err, lease.ErrNoHost):
|
||||
// No host healthy. Ask a waker to rouse a sleeping one; it answers
|
||||
// (served or 503) when it has had a turn, else falls through to the
|
||||
// plain 503.
|
||||
if p.waker != nil && p.wakeOnErrNoHost(w, r, route, routeCfg, rest, model, fp, started, lease.Key{Route: route, FP: fp, Model: model}) {
|
||||
if p.waker != nil && p.wakeOnErrNoHost(w, r, route, routeCfg, rest, model, fp, started, lease.Key{Route: route, FP: leaseFP, Model: model}) {
|
||||
return
|
||||
}
|
||||
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
||||
@@ -265,10 +271,30 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri
|
||||
p.forward(w, r, route, host, leaseState(reused), rest, fp, model, started, 0, 0, "", true)
|
||||
return
|
||||
}
|
||||
// Slot. A full queue is a 503; a context done while waiting means the client left.
|
||||
release, waited, err := p.lim.Acquire(r.Context(), host, model)
|
||||
if err != nil {
|
||||
if errors.Is(err, limiter.ErrQueueFull) {
|
||||
// Slot or track. A queue = false route leaves queueing to the client's own
|
||||
// llama-server slot: crossbar only counts the request on the host, never
|
||||
// holding it or refusing it.
|
||||
var waited time.Duration
|
||||
var release func()
|
||||
if routeCfg.Queues() {
|
||||
var err error
|
||||
release, waited, err = p.lim.Acquire(r.Context(), host, model)
|
||||
if err != nil {
|
||||
if errors.Is(err, limiter.ErrQueueFull) {
|
||||
p.writeRecord(store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
Model: model,
|
||||
Host: host,
|
||||
Started: started,
|
||||
TotalMs: time.Since(started).Milliseconds(),
|
||||
Status: http.StatusServiceUnavailable,
|
||||
Err: "queue full",
|
||||
})
|
||||
p.writeError(w, http.StatusServiceUnavailable, "queue full")
|
||||
return
|
||||
}
|
||||
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
||||
p.writeRecord(store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
@@ -276,24 +302,13 @@ func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route stri
|
||||
Host: host,
|
||||
Started: started,
|
||||
TotalMs: time.Since(started).Milliseconds(),
|
||||
Status: http.StatusServiceUnavailable,
|
||||
Err: "queue full",
|
||||
Status: 499,
|
||||
Err: "client cancelled while queued",
|
||||
})
|
||||
p.writeError(w, http.StatusServiceUnavailable, "queue full")
|
||||
return
|
||||
}
|
||||
p.log.Warn("request", "route", route, "host", host, "method", r.Method, "path", rest, "status", 499)
|
||||
p.writeRecord(store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
Model: model,
|
||||
Host: host,
|
||||
Started: started,
|
||||
TotalMs: time.Since(started).Milliseconds(),
|
||||
Status: 499,
|
||||
Err: "client cancelled while queued",
|
||||
})
|
||||
return
|
||||
} else {
|
||||
release = p.lim.Track(host, model)
|
||||
}
|
||||
defer release()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user