From aea2eeae2c4099d7d616b95b42b3bfb8c71b5be7 Mon Sep 17 00:00:00 2001 From: Kyle Isom Date: Fri, 25 Sep 2026 19:33:46 -0700 Subject: [PATCH] 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 --- docs/implementer-log.md | 1 + internal/config/config.go | 61 ---------- internal/config/config_v23_test.go | 75 ++++++++++++ internal/config/route.go | 89 ++++++++++++++ internal/limiter/limiter.go | 17 ++- internal/limiter/track_test.go | 91 +++++++++++++++ internal/proxy/affinity_test.go | 180 +++++++++++++++++++++++++++++ internal/proxy/proxy.go | 57 +++++---- 8 files changed, 486 insertions(+), 85 deletions(-) create mode 100644 internal/config/config_v23_test.go create mode 100644 internal/limiter/track_test.go create mode 100644 internal/proxy/affinity_test.go diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 0854222..f02f921 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -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..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 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). | ? | diff --git a/internal/config/config.go b/internal/config/config.go index 6559877..c59afc3 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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 -} diff --git a/internal/config/config_v23_test.go b/internal/config/config_v23_test.go new file mode 100644 index 0000000..c336107 --- /dev/null +++ b/internal/config/config_v23_test.go @@ -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`) + } +} diff --git a/internal/config/route.go b/internal/config/route.go index 43f3ce4..c6443cf 100644 --- a/internal/config/route.go +++ b/internal/config/route.go @@ -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 // "-*" 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 +} diff --git a/internal/limiter/limiter.go b/internal/limiter/limiter.go index 267e603..29ea8b4 100644 --- a/internal/limiter/limiter.go +++ b/internal/limiter/limiter.go @@ -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) diff --git a/internal/limiter/track_test.go b/internal/limiter/track_test.go new file mode 100644 index 0000000..9e8464f --- /dev/null +++ b/internal/limiter/track_test.go @@ -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) + } +} diff --git a/internal/proxy/affinity_test.go b/internal/proxy/affinity_test.go new file mode 100644 index 0000000..476889e --- /dev/null +++ b/internal/proxy/affinity_test.go @@ -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 }) + }) + } +} diff --git a/internal/proxy/proxy.go b/internal/proxy/proxy.go index d3e26c5..d6b7aff 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -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()