Wire wake and identity into crossbar; v2 smoke and README
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -32,10 +32,10 @@ func TestParallelAndQueue(t *testing.T) {
|
||||
go func() {
|
||||
rel, waited, err := l.Acquire(ctx, "alpha", "m")
|
||||
if err == nil {
|
||||
defer rel()
|
||||
if waited < 40*time.Millisecond {
|
||||
err = errors.New("third acquire did not wait")
|
||||
}
|
||||
rel() // release before reporting, so the final count check cannot race it
|
||||
}
|
||||
got3 <- err
|
||||
}()
|
||||
|
||||
@@ -69,6 +69,35 @@ func (p *Handler) guard(w http.ResponseWriter, r *http.Request, hosts []string,
|
||||
return newHost, movedHeader(host, newHost), estimate, false
|
||||
}
|
||||
|
||||
// Rule 4a: no host fits. Before refusing, ask a waker to rouse a candidate
|
||||
// whose context may grow when it comes up; a woken host takes the lease.
|
||||
if p.waker != nil {
|
||||
for _, name := range hosts {
|
||||
s, ok := p.health.Get(name)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
// A healthy host does not need waking; only a down host might grow
|
||||
// a larger context when it comes up.
|
||||
if s.Healthy {
|
||||
continue
|
||||
}
|
||||
// A host with a known per-slot context smaller than the estimate
|
||||
// cannot serve it no matter how it wakes.
|
||||
if psc := s.PerSlotCtx(); psc != 0 && psc < estimate {
|
||||
continue
|
||||
}
|
||||
if p.cfg.Hosts[name].Wake == nil || !p.waker.Wake(r.Context(), name) {
|
||||
continue
|
||||
}
|
||||
if err := p.leases.Move(lease.Key{Route: route, FP: fp, Model: model}, name, time.Now()); err != nil {
|
||||
p.writeError(w, http.StatusBadGateway, "upstream failed")
|
||||
return host, "", estimate, true
|
||||
}
|
||||
return name, movedHeader(host, name), estimate, false
|
||||
}
|
||||
}
|
||||
|
||||
// Rule 4: no host fits. Answer 400 with the estimate and the largest
|
||||
// available per-slot context, and record the row.
|
||||
p.refuseCtx(w, host, route, model, fp, started, estimate, largestSlotCtx(hosts, p.health))
|
||||
|
||||
@@ -152,6 +152,30 @@ func (p *Handler) writeError(w http.ResponseWriter, status int, msg string) {
|
||||
_ = json.NewEncoder(w).Encode(map[string]string{"error": msg})
|
||||
}
|
||||
|
||||
// writeNoHealthyHost answers the 503 when no candidate woke. The body names the
|
||||
// hosts that were asked to wake (an empty list, never null), and the row
|
||||
// records the miss.
|
||||
func (p *Handler) writeNoHealthyHost(w http.ResponseWriter, route, model, fp string, started time.Time, tried []string) {
|
||||
if tried == nil {
|
||||
tried = []string{}
|
||||
}
|
||||
p.writeRecord(store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
Model: model,
|
||||
Started: started,
|
||||
TotalMs: time.Since(started).Milliseconds(),
|
||||
Status: http.StatusServiceUnavailable,
|
||||
Err: "no healthy host",
|
||||
})
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusServiceUnavailable)
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||||
"error": "no healthy host",
|
||||
"woke": tried,
|
||||
})
|
||||
}
|
||||
|
||||
// writeRecord writes one accounting row, logging (never returning) a recorder error.
|
||||
func (p *Handler) writeRecord(req store.Request) {
|
||||
if p.rec == nil {
|
||||
|
||||
@@ -7,6 +7,7 @@ package proxy
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
@@ -45,6 +46,11 @@ type Recorder interface {
|
||||
RecordRequest(store.Request) error
|
||||
}
|
||||
|
||||
// Waker rouses a sleeping host. *wake.Waker satisfies it.
|
||||
type Waker interface {
|
||||
Wake(ctx context.Context, host string) bool
|
||||
}
|
||||
|
||||
// Handler forwards requests for a route to one of the route's healthy hosts, choosing by lease when
|
||||
// one is configured and by health alone otherwise.
|
||||
type Handler struct {
|
||||
@@ -54,6 +60,13 @@ type Handler struct {
|
||||
lim *limiter.Limiter
|
||||
rec Recorder
|
||||
log *slog.Logger
|
||||
waker Waker
|
||||
}
|
||||
|
||||
// SetWaker installs the waker the consults when a route has no healthy host left. A nil waker
|
||||
// (the default) leaves the ErrNoHost answer as it was in v0: a plain 503.
|
||||
func (p *Handler) SetWaker(w Waker) {
|
||||
p.waker = w
|
||||
}
|
||||
|
||||
// New builds a Handler. A nil logger becomes slog.Default(). With a nil lease table it behaves like
|
||||
@@ -223,6 +236,12 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
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}) {
|
||||
return
|
||||
}
|
||||
p.writeError(w, http.StatusServiceUnavailable, "no healthy host")
|
||||
case errors.Is(err, lease.ErrPinnedDown):
|
||||
p.writeError(w, http.StatusServiceUnavailable, "pinned host down")
|
||||
@@ -232,6 +251,14 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
// Slot, context guard and forward, holding the slot for the leased host.
|
||||
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, host, reused)
|
||||
}
|
||||
|
||||
// serveLeased queues the request against the leased host's limiter, runs the
|
||||
// context guard, and forwards. The slot is held for the originally leased host
|
||||
// even if the guard relocates the lease: the guard already moved it.
|
||||
func (p *Handler) serveLeased(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, host string, reused bool) {
|
||||
// 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 {
|
||||
@@ -273,3 +300,26 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
p.forward(w, r, route, host, leaseState(reused), rest, fp, model, now, waited, 0, header)
|
||||
}
|
||||
|
||||
// wakeOnErrNoHost answers the request when no host was healthy. It asks, in
|
||||
// route order, each candidate with a wake target to rouse itself; a host that
|
||||
// wakes is leased once more and then served. When none wakes, it answers 503
|
||||
// with the hosts it tried. It returns true when the request has been answered.
|
||||
func (p *Handler) wakeOnErrNoHost(w http.ResponseWriter, r *http.Request, route string, routeCfg config.Route, rest, model, fp string, started time.Time, key lease.Key) bool {
|
||||
var tried []string
|
||||
for _, name := range routeCfg.Hosts {
|
||||
if p.cfg.Hosts[name].Wake == nil {
|
||||
continue
|
||||
}
|
||||
tried = append(tried, name)
|
||||
if !p.waker.Wake(r.Context(), name) {
|
||||
continue
|
||||
}
|
||||
if newHost, _, err := p.leases.Acquire(key, routeCfg.Hosts, time.Now()); err == nil {
|
||||
p.serveLeased(w, r, route, routeCfg, rest, model, fp, started, newHost, true)
|
||||
return true
|
||||
}
|
||||
}
|
||||
p.writeNoHealthyHost(w, route, model, fp, started, tried)
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -85,6 +85,19 @@ func TestDifferentConversationsSpreadByFreeSlots(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
@@ -99,14 +112,20 @@ hosts = ["alpha"]
|
||||
default_model = "shared"
|
||||
`, alpha)
|
||||
codes := make(chan int, 3)
|
||||
for i := 1; i <= 3; i++ {
|
||||
go func(i int) {
|
||||
fire := func(i int) {
|
||||
go func() {
|
||||
resp := r.post("/r/v1/chat/completions", conversation(i, 1))
|
||||
drain(resp)
|
||||
codes <- resp.StatusCode
|
||||
}(i)
|
||||
time.Sleep(30 * time.Millisecond) // arrival order: 1 runs, 2 queues, 3 finds the queue full
|
||||
}()
|
||||
}
|
||||
// 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]++
|
||||
|
||||
Reference in New Issue
Block a user