From d0a0305bcbe84d048750231fc1c91de2e6503d3c Mon Sep 17 00:00:00 2001 From: BMAD CI Fix Agent Date: Sat, 5 Sep 2026 17:40:03 -0500 Subject: [PATCH] feat(telemetry): observe streamed backend responses --- CHANGELOG.md | 2 + README.md | 1 + docs/roadmap.md | 2 +- internal/server/server.go | 122 +++++++++++++++++++++++++++++---- internal/server/server_test.go | 52 ++++++++++++++ 5 files changed, 164 insertions(+), 15 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4d4c438..3c3c7ec 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Roadmap for maturing DevRail Router from a single-backend gateway into an observable local inference control plane. - Request IDs are now added to router responses and structured logs. +- Streaming backend responses now log first event latency and usage tokens when + the backend sends OpenAI-compatible streamed usage chunks. ### Changed diff --git a/README.md b/README.md index 46edebb..4134695 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,7 @@ This repository is in early foundation work. The current service supports: - Linux/systemd install script and unit - 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 - consistent OpenAI-shaped errors for router-side failures - request IDs in router responses and logs diff --git a/docs/roadmap.md b/docs/roadmap.md index b7995bb..aa8ce11 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -27,6 +27,7 @@ opencode, OpenClaw, and Hermes. - Bounded per-alias queueing and concurrency limits. - Command-backed profile ensure hooks for LM Studio and similar hosts. - Response telemetry for route, status, duration, bytes, and usage tokens. +- Streaming telemetry for first event latency and streamed usage chunks. - Linux tarball packaging with systemd installation. - Docker and Vagrant smoke tests. - Consistent OpenAI-shaped router errors. @@ -36,7 +37,6 @@ Useful next work: - Add readiness checks for configured backends. - Add configurable upstream transport timeouts. -- Capture streaming completion telemetry without buffering streams. - Publish example configs for common LM Studio, Ollama, and vLLM setups. ## Phase 2: Native Backend Adapters diff --git a/internal/server/server.go b/internal/server/server.go index 4ab2601..4a09989 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -299,6 +299,8 @@ type responseTelemetry struct { Backend string UpstreamModel string Status int + Streaming bool + FirstEvent time.Time PromptTokens int CompletionTokens int TotalTokens int @@ -329,23 +331,104 @@ func (body *telemetryReadCloser) Close() error { func (body *telemetryReadCloser) log() { body.once.Do(func() { - slog.Info( - "backend response completed", - "request_id", body.telemetry.RequestID, - "alias", body.telemetry.Alias, - "target_model", body.telemetry.TargetModel, - "upstream_model", body.telemetry.UpstreamModel, - "backend", body.telemetry.Backend, - "status", body.telemetry.Status, - "duration_ms", time.Since(body.telemetry.Started).Milliseconds(), - "bytes", body.telemetry.Bytes, - "prompt_tokens", body.telemetry.PromptTokens, - "completion_tokens", body.telemetry.CompletionTokens, - "total_tokens", body.telemetry.TotalTokens, - ) + logTelemetry(body.telemetry) }) } +func (telemetry responseTelemetry) firstEventMilliseconds() int64 { + if telemetry.FirstEvent.IsZero() { + return 0 + } + + return telemetry.FirstEvent.Sub(telemetry.Started).Milliseconds() +} + +type streamTelemetryReadCloser struct { + body io.ReadCloser + telemetry *responseTelemetry + lineBuffer string + once sync.Once +} + +func (body *streamTelemetryReadCloser) Read(p []byte) (int, error) { + n, err := body.body.Read(p) + body.telemetry.Bytes += int64(n) + body.observe(p[:n]) + if errors.Is(err, io.EOF) { + body.log() + } + return n, err +} + +func (body *streamTelemetryReadCloser) Close() error { + err := body.body.Close() + body.log() + return err +} + +func (body *streamTelemetryReadCloser) observe(chunk []byte) { + if len(chunk) == 0 { + return + } + + body.lineBuffer += string(chunk) + for { + index := strings.IndexByte(body.lineBuffer, '\n') + if index < 0 { + break + } + line := strings.TrimRight(body.lineBuffer[:index], "\r") + body.lineBuffer = body.lineBuffer[index+1:] + body.observeLine(line) + } + + if len(body.lineBuffer) > 64*1024 { + body.lineBuffer = "" + } +} + +func (body *streamTelemetryReadCloser) observeLine(line string) { + line = strings.TrimSpace(line) + if !strings.HasPrefix(line, "data:") { + return + } + + data := strings.TrimSpace(strings.TrimPrefix(line, "data:")) + if data == "" || data == "[DONE]" { + return + } + + if body.telemetry.FirstEvent.IsZero() { + body.telemetry.FirstEvent = time.Now() + } + applyOpenAIUsageTelemetry([]byte(data), body.telemetry) +} + +func (body *streamTelemetryReadCloser) log() { + body.once.Do(func() { + logTelemetry(body.telemetry) + }) +} + +func logTelemetry(telemetry *responseTelemetry) { + slog.Info( + "backend response completed", + "request_id", telemetry.RequestID, + "alias", telemetry.Alias, + "target_model", telemetry.TargetModel, + "upstream_model", telemetry.UpstreamModel, + "backend", telemetry.Backend, + "status", telemetry.Status, + "streaming", telemetry.Streaming, + "first_event_ms", telemetry.firstEventMilliseconds(), + "duration_ms", time.Since(telemetry.Started).Milliseconds(), + "bytes", telemetry.Bytes, + "prompt_tokens", telemetry.PromptTokens, + "completion_tokens", telemetry.CompletionTokens, + "total_tokens", telemetry.TotalTokens, + ) +} + func instrumentBackendResponse(resp *http.Response, started time.Time, model config.ModelConfig, backend config.BackendConfig, requestID string) { if resp.Body == nil { return @@ -360,6 +443,12 @@ func instrumentBackendResponse(resp *http.Response, started time.Time, model con Started: started, } + if isEventStreamResponse(resp) { + telemetry.Streaming = true + resp.Body = &streamTelemetryReadCloser{body: resp.Body, telemetry: telemetry} + return + } + if isJSONResponse(resp) { raw, err := io.ReadAll(resp.Body) if closeErr := resp.Body.Close(); closeErr != nil { @@ -381,6 +470,11 @@ func instrumentBackendResponse(resp *http.Response, started time.Time, model con resp.Body = &telemetryReadCloser{body: resp.Body, telemetry: telemetry} } +func isEventStreamResponse(resp *http.Response) bool { + contentType := strings.ToLower(resp.Header.Get("Content-Type")) + return strings.Contains(contentType, "text/event-stream") +} + func isJSONResponse(resp *http.Response) bool { contentType := strings.ToLower(resp.Header.Get("Content-Type")) return strings.Contains(contentType, "application/json") diff --git a/internal/server/server_test.go b/internal/server/server_test.go index 28c36a2..808d0a9 100644 --- a/internal/server/server_test.go +++ b/internal/server/server_test.go @@ -455,6 +455,58 @@ func TestBackendResponseTelemetryLogsUsage(t *testing.T) { } } +func TestStreamingBackendTelemetryLogsUsage(t *testing.T) { + var logs bytes.Buffer + originalLogger := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&logs, nil))) + t.Cleanup(func() { + slog.SetDefault(originalLogger) + }) + + 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\":\"o\"}}]}\n\n") + _, _ = io.WriteString(w, "data: {\"model\":\"target-model\",\"choices\":[{\"delta\":{\"content\":\"k\"}}],\"usage\":{\"prompt_tokens\":27,\"completion_tokens\":2,\"total_tokens\":29}}\n\n") + _, _ = 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() + + srv.ServeHTTP(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("unexpected status: %d", rec.Code) + } + if !strings.Contains(rec.Body.String(), "data: [DONE]") { + t.Fatalf("stream body was not proxied: %s", rec.Body.String()) + } + + logText := logs.String() + for _, want := range []string{ + `"msg":"backend response completed"`, + `"streaming":true`, + `"first_event_ms":`, + `"upstream_model":"target-model"`, + `"prompt_tokens":27`, + `"completion_tokens":2`, + `"total_tokens":29`, + } { + if !strings.Contains(logText, want) { + t.Fatalf("expected log to contain %s, got logs:\n%s", want, logText) + } + } +} + func TestBackendProxyErrorReturnsOpenAIError(t *testing.T) { t.Parallel()