From e3e8cd98f1edfd9a488222290454612a7a145699 Mon Sep 17 00:00:00 2001 From: BMAD CI Fix Agent Date: Sat, 26 Sep 2026 13:07:34 -0500 Subject: [PATCH] Add in-flight request metrics --- CHANGELOG.md | 2 + README.md | 5 +- internal/server/server.go | 94 +++++++++++++++++++++++++++++++--- internal/server/server_test.go | 74 ++++++++++++++++++++++++++ 4 files changed, 165 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e87aac8..ab96344 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 decision is made. - `/metrics` now exposes `devrail_router_prompt_chars` buckets with request labels so prompt-size guardrail behavior can be tuned from real traffic. +- `/metrics` now exposes `devrail_router_inflight_requests` so open streams + can be distinguished from completed or stalled client-side requests. - Opt-in command-backed model profile ensure hooks. - Roadmap for maturing DevRail Router from a single-backend gateway into an observable local inference control plane. diff --git a/README.md b/README.md index a358047..2a0d406 100644 --- a/README.md +++ b/README.md @@ -30,8 +30,9 @@ This repository is in early foundation work. The current service supports: - Docker image and Compose smoke testing with a mock OpenAI-compatible backend - response telemetry for proxied backend calls - streaming response telemetry for first event latency and streamed usage data -- Prometheus-compatible metrics for requests, queue wait, response latency, - first event latency, prompt size, bytes, and token totals +- Prometheus-compatible metrics for in-flight requests, completed requests, + queue wait, response latency, first event latency, prompt size, bytes, and + token totals - consistent OpenAI-shaped errors for router-side failures - request IDs in router responses and logs - a streamed benchmark harness for comparing model aliases with fixed prompts diff --git a/internal/server/server.go b/internal/server/server.go index 2ca5bd6..31dd543 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -369,6 +369,7 @@ type responseTelemetry struct { HasPromptChars bool Started time.Time metrics *metricsRegistry + finishInFlight func() } type telemetryReadCloser struct { @@ -474,6 +475,9 @@ func (body *streamTelemetryReadCloser) log() { } func logTelemetry(telemetry *responseTelemetry) { + if telemetry.finishInFlight != nil { + defer telemetry.finishInFlight() + } telemetry.metrics.record(requestMetrics{ Alias: telemetry.Alias, TargetModel: telemetry.TargetModel, @@ -545,6 +549,18 @@ func instrumentBackendResponse( if isEventStreamResponse(resp) { telemetry.Streaming = true + telemetry.finishInFlight = metrics.startInFlight(responseTelemetryLabels(*telemetry)) + slog.Info( + "backend response started", + "request_id", requestID, + "alias", telemetry.Alias, + "target_model", telemetry.TargetModel, + "backend", telemetry.Backend, + "status", telemetry.Status, + "streaming", telemetry.Streaming, + "route_rule", telemetry.RouteRule, + "prompt_chars", telemetry.PromptChars, + ) resp.Body = &streamTelemetryReadCloser{body: resp.Body, telemetry: telemetry} return } @@ -567,6 +583,18 @@ func instrumentBackendResponse( } } + telemetry.finishInFlight = metrics.startInFlight(responseTelemetryLabels(*telemetry)) + slog.Info( + "backend response started", + "request_id", requestID, + "alias", telemetry.Alias, + "target_model", telemetry.TargetModel, + "backend", telemetry.Backend, + "status", telemetry.Status, + "streaming", telemetry.Streaming, + "route_rule", telemetry.RouteRule, + "prompt_chars", telemetry.PromptChars, + ) resp.Body = &telemetryReadCloser{body: resp.Body, telemetry: telemetry} } @@ -621,6 +649,7 @@ type requestMetrics struct { type metricsRegistry struct { mu sync.Mutex requests map[string]*metricSeries + inFlight map[string]*metricSeries durationSeconds histogram queueWaitSeconds histogram ensureSeconds histogram @@ -665,6 +694,7 @@ func newMetricsRegistry() *metricsRegistry { promptCharBuckets := []float64{1024, 4096, 8192, 16384, 32768, 65536, 80000, 131072, 262144, 524288, 1048576} return &metricsRegistry{ requests: make(map[string]*metricSeries), + inFlight: make(map[string]*metricSeries), durationSeconds: histogram{ Buckets: latencyBuckets, Series: make(map[string]*histogramSeries), @@ -688,6 +718,28 @@ func newMetricsRegistry() *metricsRegistry { } } +func requestMetricsLabels(metrics requestMetrics) metricLabels { + return metricLabels{ + Alias: metrics.Alias, + Backend: metrics.Backend, + TargetModel: metrics.TargetModel, + RouteRule: metrics.RouteRule, + Status: strconv.Itoa(metrics.Status), + Streaming: strconv.FormatBool(metrics.Streaming), + } +} + +func responseTelemetryLabels(telemetry responseTelemetry) metricLabels { + return metricLabels{ + Alias: telemetry.Alias, + Backend: telemetry.Backend, + TargetModel: telemetry.TargetModel, + RouteRule: telemetry.RouteRule, + Status: strconv.Itoa(telemetry.Status), + Streaming: strconv.FormatBool(telemetry.Streaming), + } +} + func requestMetricsFromModel(model config.ModelConfig, backend config.BackendConfig, status int, started time.Time) requestMetrics { return requestMetrics{ Alias: model.ID, @@ -711,14 +763,7 @@ func (registry *metricsRegistry) record(metrics requestMetrics) { return } - labels := metricLabels{ - Alias: metrics.Alias, - Backend: metrics.Backend, - TargetModel: metrics.TargetModel, - RouteRule: metrics.RouteRule, - Status: strconv.Itoa(metrics.Status), - Streaming: strconv.FormatBool(metrics.Streaming), - } + labels := requestMetricsLabels(metrics) duration := time.Since(metrics.Started).Seconds() if metrics.Started.IsZero() { duration = 0 @@ -746,6 +791,33 @@ func (registry *metricsRegistry) record(metrics requestMetrics) { registry.totalTokens += float64(metrics.TotalTokens) } +func (registry *metricsRegistry) startInFlight(labels metricLabels) func() { + if registry == nil { + return nil + } + + registry.mu.Lock() + key := labels.key() + if registry.inFlight[key] == nil { + registry.inFlight[key] = &metricSeries{Labels: labels} + } + registry.inFlight[key].Value++ + registry.mu.Unlock() + + var once sync.Once + return func() { + once.Do(func() { + registry.mu.Lock() + defer registry.mu.Unlock() + series := registry.inFlight[key] + if series == nil || series.Value <= 0 { + return + } + series.Value-- + }) + } +} + func (registry *metricsRegistry) recordEnsure(model config.ModelConfig, status string, duration time.Duration) { if registry == nil { return @@ -793,6 +865,12 @@ func (registry *metricsRegistry) render() string { writeMetricLine(&builder, "devrail_router_requests_total", series.Labels, series.Value) } + writeMetricHelp(&builder, "devrail_router_inflight_requests", "Currently active proxied OpenAI-compatible requests with an upstream response.") + writeMetricType(&builder, "devrail_router_inflight_requests", "gauge") + for _, series := range sortedMetricSeries(registry.inFlight) { + writeMetricLine(&builder, "devrail_router_inflight_requests", series.Labels, series.Value) + } + writeHistogram(&builder, "devrail_router_request_duration_seconds", "End-to-end router request duration in seconds.", registry.durationSeconds) writeHistogram(&builder, "devrail_router_queue_wait_seconds", "Time spent waiting for a model concurrency slot in seconds.", registry.queueWaitSeconds) writeHistogram(&builder, "devrail_router_ensure_duration_seconds", "Time spent ensuring a model profile is ready before proxying.", registry.ensureSeconds) diff --git a/internal/server/server_test.go b/internal/server/server_test.go index d66ff2d..913bdc3 100644 --- a/internal/server/server_test.go +++ b/internal/server/server_test.go @@ -1056,6 +1056,80 @@ func TestMetricsEndpointExposesStreamingFirstEventLatency(t *testing.T) { } } +func TestMetricsEndpointExposesInflightStreamingRequests(t *testing.T) { + started := make(chan struct{}) + release := make(chan struct{}) + backend := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + _, _ = io.WriteString(w, "data: {\"model\":\"target-model\",\"choices\":[{\"delta\":{\"content\":\"ok\"}}]}\n\n") + if flusher, ok := w.(http.Flusher); ok { + flusher.Flush() + } + close(started) + <-release + _, _ = io.WriteString(w, "data: [DONE]\n\n") + })) + t.Cleanup(backend.Close) + + srv := testServerWithBackend(t, backend.URL, config.ModelConfig{ + ID: "local-coder", + Backend: "lmstudio", + TargetModel: "target-model", + }) + req := httptest.NewRequest( + http.MethodPost, + "/v1/chat/completions", + strings.NewReader(`{"model":"local-coder","messages":[],"stream":true}`), + ) + rec := httptest.NewRecorder() + done := make(chan struct{}) + go func() { + defer close(done) + srv.ServeHTTP(rec, req) + }() + + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("timed out waiting for backend stream to start") + } + + wantInflight := `devrail_router_inflight_requests{alias="local-coder",backend="lmstudio",target_model="target-model",route_rule="default",status="200",streaming="true"} 1` + deadline := time.Now().Add(time.Second) + for { + metricsReq := httptest.NewRequest(http.MethodGet, "/metrics", nil) + metricsRec := httptest.NewRecorder() + srv.ServeHTTP(metricsRec, metricsReq) + if strings.Contains(metricsRec.Body.String(), wantInflight) { + break + } + if time.Now().After(deadline) { + t.Fatalf("expected in-flight metric %s, got:\n%s", wantInflight, metricsRec.Body.String()) + } + time.Sleep(10 * time.Millisecond) + } + + close(release) + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("timed out waiting for router request to finish") + } + + metricsReq := httptest.NewRequest(http.MethodGet, "/metrics", nil) + metricsRec := httptest.NewRecorder() + srv.ServeHTTP(metricsRec, metricsReq) + body := metricsRec.Body.String() + for _, want := range []string{ + `devrail_router_inflight_requests{alias="local-coder",backend="lmstudio",target_model="target-model",route_rule="default",status="200",streaming="true"} 0`, + `devrail_router_requests_total{alias="local-coder",backend="lmstudio",target_model="target-model",route_rule="default",status="200",streaming="true"} 1`, + } { + if !strings.Contains(body, want) { + t.Fatalf("expected metrics to contain %s, got:\n%s", want, body) + } + } +} + func TestBackendProxyErrorReturnsOpenAIError(t *testing.T) { t.Parallel()