Files
kyle 97f7cdffd9 Route by lease, queue per host and model, record every request
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
2026-09-25 05:15:21 -07:00

167 lines
4.6 KiB
Go

package proxy
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httputil"
"net/url"
"strings"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
// forward builds the reverse proxy for one host, tees the response, records the accounting row, and
// logs. leaseState is "new" or "reused"; waited is the time spent in the queue.
func (p *Handler) forward(w http.ResponseWriter, r *http.Request, route, host, leaseState, rest, fp, model string, started time.Time, waited time.Duration) {
hostCfg, ok := p.cfg.Hosts[host]
if !ok {
p.writeError(w, http.StatusBadGateway, "upstream failed")
return
}
target, err := url.Parse(hostCfg.BaseURL)
if err != nil {
p.writeError(w, http.StatusBadGateway, "upstream failed")
return
}
rev := &forwardState{started: started}
rp := newReverseProxy(p.health, host, leaseState, target, rest, rev)
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
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)
}
p.writeRecord(req)
fp8 := fp
if len(fp8) > 8 {
fp8 = fp8[:8]
}
p.log.Info("request",
"route", route,
"host", host,
"method", r.Method,
"path", rest,
"status", rec.status,
"lease", leaseState,
"queued_ms", waited.Milliseconds(),
"fp", fp8,
"ms", total.Milliseconds(),
)
}
// leaseState is "reused" when the lease already held the conversation, else "new".
func leaseState(reused bool) string {
if reused {
return "reused"
}
return "new"
}
// ttfbMs is the time from request start to the response head; zero when the head never arrived.
func ttfbMs(rev *forwardState) int64 {
if rev.ttfb.IsZero() || rev.ttfb.Before(rev.started) {
return 0
}
return rev.ttfb.Sub(rev.started).Milliseconds()
}
// forwardState carries, across one forward, when the request started, when the head arrived, whether
// the response streamed, and the tee that scanned it.
type forwardState struct {
started time.Time
ttfb time.Time
streamed bool
tee *tee
}
// statusRecorder records the status written and forwards Flush so the reverse proxy can stream.
type statusRecorder struct {
http.ResponseWriter
status int
}
func (r *statusRecorder) WriteHeader(code int) {
r.status = code
r.ResponseWriter.WriteHeader(code)
}
func (r *statusRecorder) Flush() {
if f, ok := r.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
// writeError answers with a JSON {"error":"…"} body.
func (p *Handler) writeError(w http.ResponseWriter, status int, msg string) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(map[string]string{"error": msg})
}
// writeRecord writes one accounting row, logging (never returning) a recorder error.
func (p *Handler) writeRecord(req store.Request) {
if p.rec == nil {
return
}
if err := p.rec.RecordRequest(req); err != nil {
p.log.Error("record request", "err", err)
}
}
// newReverseProxy forwards to a single host, rewriting the path to target.Path+rest and keeping the
// original query string. It flushes after every write so long server-sent-event streams are not
// buffered, tees the response for usage/timings, and marks the host down on any transport error
// other than a client disconnect.
func newReverseProxy(h Health, host, leaseState string, target *url.URL, rest string, rev *forwardState) *httputil.ReverseProxy {
return &httputil.ReverseProxy{
Rewrite: func(pr *httputil.ProxyRequest) {
pr.SetURL(target)
pr.Out.URL.Path = target.Path + rest
pr.Out.URL.RawPath = ""
pr.Out.Host = target.Host
pr.SetXForwarded()
},
FlushInterval: -1,
ModifyResponse: func(resp *http.Response) error {
resp.Header.Set(HostHeader, host)
resp.Header.Set(LeaseHeader, leaseState)
rev.ttfb = time.Now()
rev.streamed = strings.HasPrefix(resp.Header.Get("Content-Type"), "text/event-stream")
t := newTee(resp.Body, rev.streamed)
resp.Body = t
rev.tee = t
return nil
},
ErrorHandler: func(w http.ResponseWriter, req *http.Request, err error) {
if errors.Is(err, context.Canceled) {
return
}
h.MarkDown(host, err.Error())
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusBadGateway)
_ = json.NewEncoder(w).Encode(map[string]string{"error": "upstream failed", "host": host})
},
}
}