Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 3 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
94 changes: 86 additions & 8 deletions internal/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -369,6 +369,7 @@ type responseTelemetry struct {
HasPromptChars bool
Started time.Time
metrics *metricsRegistry
finishInFlight func()
}

type telemetryReadCloser struct {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
}
Expand All @@ -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}
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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),
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
74 changes: 74 additions & 0 deletions internal/server/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
Loading