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 @@ -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

Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion docs/roadmap.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down
122 changes: 108 additions & 14 deletions internal/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,8 @@ type responseTelemetry struct {
Backend string
UpstreamModel string
Status int
Streaming bool
FirstEvent time.Time
PromptTokens int
CompletionTokens int
TotalTokens int
Expand Down Expand Up @@ -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
Expand All @@ -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 {
Expand All @@ -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")
Expand Down
52 changes: 52 additions & 0 deletions internal/server/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
Loading