Skip to content
Open
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
4 changes: 4 additions & 0 deletions internal/pluginhost/adapters_usage_translation.go
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,10 @@ func (a *usageAdapter) HandleUsage(ctx context.Context, record coreusage.Record)
RequestedAt: record.RequestedAt,
Latency: record.Latency,
TTFT: record.TTFT,
UpstreamTTFB: record.UpstreamTTFB,
FirstPacket: record.FirstPacket,
ConnSetup: record.ConnSetup,
ConnReused: record.ConnReused,
Failed: record.Failed,
Failure: pluginapi.UsageFailure{
StatusCode: record.Fail.StatusCode,
Expand Down
8 changes: 8 additions & 0 deletions internal/redisqueue/plugin.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,10 @@ func (p *usageQueuePlugin) HandleUsage(ctx context.Context, record coreusage.Rec
Timestamp: timestamp,
LatencyMs: record.Latency.Milliseconds(),
TTFTMs: record.TTFT.Milliseconds(),
UpstreamTTFBMs: record.UpstreamTTFB.Milliseconds(),
FirstPacketMs: record.FirstPacket.Milliseconds(),
ConnSetupMs: record.ConnSetup.Milliseconds(),
ConnReused: record.ConnReused,
Source: record.Source,
AuthIndex: record.AuthIndex,
AccessTokenHash: record.AccessTokenSHA256,
Expand Down Expand Up @@ -178,6 +182,10 @@ type requestDetail struct {
Timestamp time.Time `json:"timestamp"`
LatencyMs int64 `json:"latency_ms"`
TTFTMs int64 `json:"ttft_ms"`
UpstreamTTFBMs int64 `json:"upstream_ttfb_ms,omitempty"`
FirstPacketMs int64 `json:"first_packet_ms,omitempty"`
ConnSetupMs int64 `json:"conn_setup_ms,omitempty"`
ConnReused bool `json:"conn_reused,omitempty"`
Source string `json:"source"`
AuthIndex string `json:"auth_index"`
AccessTokenHash string `json:"access_token_sha256,omitempty"`
Expand Down
130 changes: 130 additions & 0 deletions internal/runtime/executor/helps/upstream_timing_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package helps

import (
"context"
"io"
"net/http"
"net/http/httptest"
"testing"

"github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/usage"
)

func newUpstreamTimingTestReporter(t *testing.T) *UsageReporter {
t.Helper()
reporter := NewUsageReporter(context.Background(), "test-provider", "test-model", nil)
if reporter == nil {
t.Fatal("NewUsageReporter returned nil")
}
return reporter
}

func upstreamTimingTestServer(t *testing.T) *httptest.Server {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "text/plain")
_, _ = w.Write([]byte("hello"))
}))
t.Cleanup(server.Close)
return server
}

func drainAndClose(t *testing.T, resp *http.Response) {
t.Helper()
if _, errRead := io.Copy(io.Discard, resp.Body); errRead != nil {
t.Fatalf("read upstream body: %v", errRead)
}
if errClose := resp.Body.Close(); errClose != nil {
t.Fatalf("close upstream body: %v", errClose)
}
}

func TestRoundTripRecordsUpstreamTiming(t *testing.T) {
server := upstreamTimingTestServer(t)
reporter := newUpstreamTimingTestReporter(t)
client := reporter.TrackHTTPClient(server.Client())

resp, errGet := client.Get(server.URL)
if errGet != nil {
t.Fatalf("get upstream: %v", errGet)
}
drainAndClose(t, resp)

record := reporter.buildRecord(usage.Detail{}, false)
if record.UpstreamTTFB <= 0 {
t.Fatalf("upstream ttfb = %v, want > 0", record.UpstreamTTFB)
}
if record.FirstPacket < record.UpstreamTTFB {
t.Fatalf("first packet %v is smaller than upstream ttfb %v", record.FirstPacket, record.UpstreamTTFB)
}
if record.ConnReused {
t.Fatal("conn reused = true, want false for the first request")
}
if record.ConnSetup < 0 {
t.Fatalf("conn setup = %v, want >= 0", record.ConnSetup)
}
if record.TTFT <= 0 {
t.Fatalf("ttft = %v, want > 0", record.TTFT)
}
}

func TestRoundTripRecordsConnReuse(t *testing.T) {
server := upstreamTimingTestServer(t)
client := server.Client()

firstReporter := newUpstreamTimingTestReporter(t)
firstResp, errFirst := firstReporter.TrackHTTPClient(client).Get(server.URL)
if errFirst != nil {
t.Fatalf("first get: %v", errFirst)
}
drainAndClose(t, firstResp)

secondReporter := newUpstreamTimingTestReporter(t)
secondResp, errSecond := secondReporter.TrackHTTPClient(client).Get(server.URL)
if errSecond != nil {
t.Fatalf("second get: %v", errSecond)
}
drainAndClose(t, secondResp)

record := secondReporter.buildRecord(usage.Detail{}, false)
if !record.ConnReused {
t.Fatal("conn reused = false, want true for the pooled second request")
}
if record.ConnSetup != 0 {
t.Fatalf("conn setup = %v, want 0 when the connection is reused", record.ConnSetup)
}
}

func TestObserveUpstreamAttemptFirstWins(t *testing.T) {
reporter := newUpstreamTimingTestReporter(t)
reporter.ObserveUpstreamAttempt(upstreamAttemptTiming{upstreamTTFB: 100, connSetup: 20, connReused: true})
reporter.ObserveUpstreamAttempt(upstreamAttemptTiming{upstreamTTFB: 999, connSetup: 77, connReused: false})

record := reporter.buildRecord(usage.Detail{}, false)
if record.UpstreamTTFB != 100 {
t.Fatalf("upstream ttfb = %v, want first observed 100", record.UpstreamTTFB)
}
if record.ConnSetup != 20 {
t.Fatalf("conn setup = %v, want first observed 20", record.ConnSetup)
}
if !record.ConnReused {
t.Fatal("conn reused = false, want first observed true")
}
}

func TestMarkFirstResponseByteRecordsFirstPacket(t *testing.T) {
server := upstreamTimingTestServer(t)
reporter := newUpstreamTimingTestReporter(t)
client := reporter.TrackHTTPClient(server.Client())

resp, errGet := client.Get(server.URL)
if errGet != nil {
t.Fatalf("get upstream: %v", errGet)
}
drainAndClose(t, resp)

record := reporter.buildRecord(usage.Detail{}, false)
if record.FirstPacket != record.TTFT {
t.Fatalf("first packet = %v, ttft = %v, want equal on the tracked http client path", record.FirstPacket, record.TTFT)
}
}
163 changes: 162 additions & 1 deletion internal/runtime/executor/helps/usage_helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,12 @@ package helps
import (
"bytes"
"context"
"crypto/tls"
"errors"
"fmt"
"io"
"net/http"
"net/http/httptrace"
"reflect"
"strings"
"sync"
Expand Down Expand Up @@ -52,6 +54,10 @@ type UsageReporter struct {
ttftSet bool
once sync.Once

upstreamMu sync.Mutex
upstreamTiming upstreamAttemptTiming
upstreamTimingSet bool

responseModelMu sync.RWMutex
// responseModel holds the latest model name reported by the upstream response.
responseModel string
Expand Down Expand Up @@ -500,7 +506,26 @@ func (r *UsageReporter) MarkFirstResponseByte() {
if start.IsZero() {
return
}
r.setTTFT(time.Since(start))
elapsed := time.Since(start)
r.recordFirstPacket(elapsed)
r.setTTFT(elapsed)
}

// recordFirstPacket stores the first response body byte duration exactly once so
// usage sinks can expose it next to the effective TTFT.
func (r *UsageReporter) recordFirstPacket(elapsed time.Duration) {
if r == nil {
return
}
if elapsed < 0 {
elapsed = 0
}
r.ttftMu.Lock()
if !r.firstPacketSet {
r.firstPacketDuration = elapsed
r.firstPacketSet = true
}
r.ttftMu.Unlock()
}

func (r *UsageReporter) buildAdditionalModelRecord(model string, detail usage.Detail) (usage.Record, bool) {
Expand Down Expand Up @@ -629,6 +654,10 @@ func (r *UsageReporter) buildRecordForModel(model string, detail usage.Detail, f
RequestedAt: r.requestedAt,
Latency: r.latency(),
TTFT: r.ttftDuration(),
UpstreamTTFB: r.upstreamTimings().upstreamTTFB,
FirstPacket: r.firstPacketElapsed(),
ConnSetup: r.upstreamTimings().connSetup,
ConnReused: r.upstreamTimings().connReused,
Failed: failed,
Fail: fail,
Detail: detail,
Expand Down Expand Up @@ -702,6 +731,135 @@ func (r *UsageReporter) ttftDuration() time.Duration {
return 0
}

// upstreamAttemptTiming captures transport-level timings of one upstream HTTP attempt.
type upstreamAttemptTiming struct {
// upstreamTTFB spans request dispatch until the first response header byte.
upstreamTTFB time.Duration
// connSetup accumulates DNS + TCP connect + TLS handshake durations.
connSetup time.Duration
// connReused reports whether a pooled connection was reused.
connReused bool
}

// upstreamTraceCollector gathers httptrace callbacks for one upstream attempt.
type upstreamTraceCollector struct {
mu sync.Mutex
start time.Time
dnsStart time.Time
connectStart time.Time
tlsStart time.Time
timing upstreamAttemptTiming
headerSet bool
}

func newUpstreamTraceCollector() *upstreamTraceCollector {
return &upstreamTraceCollector{start: time.Now()}
}

// trace returns the httptrace.ClientTrace wired into the upstream request context.
func (c *upstreamTraceCollector) trace() *httptrace.ClientTrace {
return &httptrace.ClientTrace{
DNSStart: func(httptrace.DNSStartInfo) {
c.mu.Lock()
c.dnsStart = time.Now()
c.mu.Unlock()
},
DNSDone: func(httptrace.DNSDoneInfo) {
c.mu.Lock()
if !c.dnsStart.IsZero() {
c.timing.connSetup += time.Since(c.dnsStart)
c.dnsStart = time.Time{}
}
c.mu.Unlock()
},
ConnectStart: func(string, string) {
c.mu.Lock()
c.connectStart = time.Now()
c.mu.Unlock()
},
ConnectDone: func(string, string, error) {
c.mu.Lock()
if !c.connectStart.IsZero() {
c.timing.connSetup += time.Since(c.connectStart)
c.connectStart = time.Time{}
}
c.mu.Unlock()
},
TLSHandshakeStart: func() {
c.mu.Lock()
c.tlsStart = time.Now()
c.mu.Unlock()
},
TLSHandshakeDone: func(tls.ConnectionState, error) {
c.mu.Lock()
if !c.tlsStart.IsZero() {
c.timing.connSetup += time.Since(c.tlsStart)
c.tlsStart = time.Time{}
}
c.mu.Unlock()
},
GotConn: func(info httptrace.GotConnInfo) {
c.mu.Lock()
c.timing.connReused = info.Reused
c.mu.Unlock()
},
GotFirstResponseByte: func() {
c.mu.Lock()
if !c.headerSet {
c.timing.upstreamTTFB = time.Since(c.start)
c.headerSet = true
}
c.mu.Unlock()
},
}
}

// snapshot returns the timings collected so far.
func (c *upstreamTraceCollector) snapshot() upstreamAttemptTiming {
c.mu.Lock()
defer c.mu.Unlock()
return c.timing
}

// ObserveUpstreamAttempt records transport-level timings of one upstream HTTP
// attempt. The first observed attempt wins so credential retries cannot
// overwrite the recorded timings.
func (r *UsageReporter) ObserveUpstreamAttempt(timing upstreamAttemptTiming) {
if r == nil {
return
}
r.upstreamMu.Lock()
defer r.upstreamMu.Unlock()
if r.upstreamTimingSet {
return
}
r.upstreamTiming = timing
r.upstreamTimingSet = true
}

// upstreamTimings returns the recorded upstream attempt timings.
func (r *UsageReporter) upstreamTimings() upstreamAttemptTiming {
if r == nil {
return upstreamAttemptTiming{}
}
r.upstreamMu.Lock()
defer r.upstreamMu.Unlock()
return r.upstreamTiming
}

// firstPacketElapsed returns the recorded first response body byte duration.
func (r *UsageReporter) firstPacketElapsed() time.Duration {
if r == nil {
return 0
}
r.ttftMu.RLock()
defer r.ttftMu.RUnlock()
if r.firstPacketSet {
return r.firstPacketDuration
}
return 0
}

type usageTTFTRoundTripper struct {
base http.RoundTripper
reporter *UsageReporter
Expand All @@ -711,7 +869,10 @@ type usageTTFTRoundTripper struct {
func (t usageTTFTRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
cliproxyexecutor.MarkUpstreamAttempt(req.Context())
t.reporter.StartResponseTTFT()
collector := newUpstreamTraceCollector()
req = req.WithContext(httptrace.WithClientTrace(req.Context(), collector.trace()))
resp, errRoundTrip := t.base.RoundTrip(req)
t.reporter.ObserveUpstreamAttempt(collector.snapshot())
if errRoundTrip != nil {
return resp, errRoundTrip
}
Expand Down
17 changes: 14 additions & 3 deletions sdk/cliproxy/usage/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,9 +56,20 @@ type Record struct {
RequestedAt time.Time
Latency time.Duration
TTFT time.Duration
Failed bool
Fail Failure
Detail Detail
// UpstreamTTFB is the duration from dispatching the upstream request until the
// first upstream response header byte arrives at the gateway.
UpstreamTTFB time.Duration
// FirstPacket is the duration from dispatching the upstream request until the
// first upstream response body byte arrives at the gateway.
FirstPacket time.Duration
// ConnSetup is the accumulated DNS + TCP connect + TLS handshake duration of
// the upstream attempt. It is zero when a pooled connection was reused.
ConnSetup time.Duration
// ConnReused reports whether the upstream attempt reused a pooled connection.
ConnReused bool
Failed bool
Fail Failure
Detail Detail
// ResponseHeaders stores a snapshot of upstream response headers for usage sinks.
ResponseHeaders http.Header
}
Expand Down
Loading
Loading