Review fixes: record cancelled requests as 499; empty usage is an array
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -0,0 +1,109 @@
|
||||
package proxy_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/store"
|
||||
)
|
||||
|
||||
// A client that goes away mid-stream is still a request that happened: it held a slot, it cost
|
||||
// prefill, and it belongs in the accounting. The row records status 499 and a non-empty err.
|
||||
func TestClientCancelMidStreamIsRecorded(t *testing.T) {
|
||||
slow := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
case "/health":
|
||||
fmt.Fprint(w, `{"status":"ok"}`)
|
||||
case "/v1/models":
|
||||
fmt.Fprint(w, `{"object":"list","data":[{"id":"shared"}]}`)
|
||||
default:
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.WriteHeader(200)
|
||||
fmt.Fprint(w, "data: {\"choices\":[{\"delta\":{\"content\":\"first\"}}]}\n\n")
|
||||
w.(http.Flusher).Flush()
|
||||
select {
|
||||
case <-r.Context().Done():
|
||||
case <-time.After(3 * time.Second):
|
||||
}
|
||||
}
|
||||
}))
|
||||
t.Cleanup(slow.Close)
|
||||
beta := newUpstream(t, "beta")
|
||||
r := newRig(t, twoHosts, &upstream{name: "alpha", srv: slow}, beta)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
body := `{"model":"alpha-only","stream":true,"messages":[{"role":"user","content":"cancel me"}]}`
|
||||
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(body))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
buf := make([]byte, 64)
|
||||
if _, err := resp.Body.Read(buf); err != nil {
|
||||
t.Fatalf("first chunk: %v", err)
|
||||
}
|
||||
cancel()
|
||||
resp.Body.Close()
|
||||
|
||||
deadline := time.Now().Add(3 * time.Second)
|
||||
var counts []store.StatusCount
|
||||
for time.Now().Before(deadline) {
|
||||
counts, _ = r.store.StatusCounts(time.Time{})
|
||||
if len(counts) > 0 {
|
||||
break
|
||||
}
|
||||
time.Sleep(25 * time.Millisecond)
|
||||
}
|
||||
if len(counts) != 1 || counts[0].Status != 499 || counts[0].Route != "r" || counts[0].Count != 1 {
|
||||
t.Fatalf("status counts after a cancelled stream = %+v, want one row: route r, status 499", counts)
|
||||
}
|
||||
rows, _ := r.store.Usage(time.Time{}, store.ByRoute)
|
||||
if len(rows) != 1 || rows[0].Requests != 1 || rows[0].Errors != 1 {
|
||||
t.Errorf("usage = %+v, want 1 request counted as an error", rows)
|
||||
}
|
||||
}
|
||||
|
||||
// The same when the client gives up while waiting in the queue: a 499 row, no slot leaked.
|
||||
func TestClientCancelWhileQueuedIsRecorded(t *testing.T) {
|
||||
alpha := newUpstream(t, "alpha")
|
||||
alpha.delay = 800 * time.Millisecond
|
||||
r := newRig(t, `
|
||||
listen = "127.0.0.1:1"
|
||||
queue_max = 2
|
||||
[hosts.alpha]
|
||||
base_url = %q
|
||||
models = { "shared" = { parallel = 1 } }
|
||||
[routes.r]
|
||||
hosts = ["alpha"]
|
||||
default_model = "shared"
|
||||
`, alpha)
|
||||
go func() { drain(r.post("/r/v1/chat/completions", conversation(1, 1))) }() // holds the one slot
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
|
||||
defer cancel()
|
||||
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, r.front.URL+"/r/v1/chat/completions", strings.NewReader(conversation(2, 1)))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
if _, err := http.DefaultClient.Do(req); err == nil {
|
||||
t.Fatal("the queued request should have been cancelled by its context")
|
||||
}
|
||||
deadline := time.Now().Add(3 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
counts, _ := r.store.StatusCounts(time.Time{})
|
||||
for _, c := range counts {
|
||||
if c.Status == 499 {
|
||||
if r.lim.Queued("alpha", "shared") != 0 {
|
||||
t.Errorf("queued = %d after the waiter cancelled", r.lim.Queued("alpha", "shared"))
|
||||
}
|
||||
return
|
||||
}
|
||||
}
|
||||
time.Sleep(25 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("no 499 row recorded for the request cancelled while queued")
|
||||
}
|
||||
+47
-17
@@ -30,26 +30,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
||||
rev := &forwardState{started: started}
|
||||
rp := newReverseProxy(p.health, host, leaseState, target, rest, rev)
|
||||
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
|
||||
|
||||
// ServeHTTP unwinds with http.ErrAbortHandler when a client leaves mid-stream; recover so the
|
||||
// row the request earned is still written, then re-panic so the server keeps its semantics.
|
||||
defer func() {
|
||||
if pv := recover(); pv != nil {
|
||||
total := time.Since(started).Milliseconds()
|
||||
req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total)
|
||||
if perr, ok := pv.(error); ok && errors.Is(perr, http.ErrAbortHandler) {
|
||||
req.Status = 499
|
||||
req.Err = "client cancelled"
|
||||
} else {
|
||||
req.Err = "upstream error"
|
||||
}
|
||||
p.writeRecord(req)
|
||||
panic(pv)
|
||||
}
|
||||
}()
|
||||
|
||||
rp.ServeHTTP(rec, r)
|
||||
total := time.Since(started)
|
||||
|
||||
req := store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
Model: model,
|
||||
Host: host,
|
||||
Started: started,
|
||||
QueuedMs: waited.Milliseconds(),
|
||||
TTFBMs: ttfbMs(rev),
|
||||
TotalMs: total.Milliseconds(),
|
||||
Status: rec.status,
|
||||
Streamed: rev.streamed,
|
||||
}
|
||||
if rev.tee != nil {
|
||||
prompt, cached, completion := rev.tee.tokens()
|
||||
req.PromptTokens = int64(prompt)
|
||||
req.CachedTokens = int64(cached)
|
||||
req.CompletionTokens = int64(completion)
|
||||
req := forwardRow(route, fp, model, host, started, waited, rev, rec.status, total.Milliseconds())
|
||||
if r.Context().Err() != nil {
|
||||
req.Status = 499
|
||||
req.Err = "client cancelled"
|
||||
}
|
||||
p.writeRecord(req)
|
||||
|
||||
@@ -70,6 +75,31 @@ func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, l
|
||||
)
|
||||
}
|
||||
|
||||
// forwardRow builds the accounting row from the state a forward gathered: what the tee scanned and
|
||||
// what the recorder captured. totalMs is measured from start to the caller's exit, so the forward
|
||||
// path and the recovery path above build identical rows.
|
||||
func forwardRow(route, fp, model, host string, start time.Time, waited time.Duration, rev *forwardState, status int, totalMs int64) store.Request {
|
||||
req := store.Request{
|
||||
Route: route,
|
||||
FP: fp,
|
||||
Model: model,
|
||||
Host: host,
|
||||
Started: start,
|
||||
QueuedMs: waited.Milliseconds(),
|
||||
TTFBMs: ttfbMs(rev),
|
||||
TotalMs: totalMs,
|
||||
Status: status,
|
||||
Streamed: rev.streamed,
|
||||
}
|
||||
if rev.tee != nil {
|
||||
prompt, cached, completion := rev.tee.tokens()
|
||||
req.PromptTokens = int64(prompt)
|
||||
req.CachedTokens = int64(cached)
|
||||
req.CompletionTokens = int64(completion)
|
||||
}
|
||||
return req
|
||||
}
|
||||
|
||||
// leaseState is "reused" when the lease already held the conversation, else "new".
|
||||
func leaseState(reused bool) string {
|
||||
if reused {
|
||||
|
||||
@@ -249,6 +249,16 @@ func (p *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
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
|
||||
}
|
||||
defer release()
|
||||
|
||||
Reference in New Issue
Block a user