diff --git a/internal/pluginhost/adapters_usage_translation.go b/internal/pluginhost/adapters_usage_translation.go index b17e9345055..a4eb2c3e9a1 100644 --- a/internal/pluginhost/adapters_usage_translation.go +++ b/internal/pluginhost/adapters_usage_translation.go @@ -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, diff --git a/internal/redisqueue/plugin.go b/internal/redisqueue/plugin.go index 64f0829214a..8dfd09d90bb 100644 --- a/internal/redisqueue/plugin.go +++ b/internal/redisqueue/plugin.go @@ -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, @@ -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"` diff --git a/internal/runtime/executor/helps/upstream_timing_test.go b/internal/runtime/executor/helps/upstream_timing_test.go new file mode 100644 index 00000000000..d992b589dde --- /dev/null +++ b/internal/runtime/executor/helps/upstream_timing_test.go @@ -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) + } +} diff --git a/internal/runtime/executor/helps/usage_helpers.go b/internal/runtime/executor/helps/usage_helpers.go index 77b66d2b54c..9765d834ae4 100644 --- a/internal/runtime/executor/helps/usage_helpers.go +++ b/internal/runtime/executor/helps/usage_helpers.go @@ -3,10 +3,12 @@ package helps import ( "bytes" "context" + "crypto/tls" "errors" "fmt" "io" "net/http" + "net/http/httptrace" "reflect" "strings" "sync" @@ -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 @@ -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) { @@ -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, @@ -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 @@ -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 } diff --git a/sdk/cliproxy/usage/manager.go b/sdk/cliproxy/usage/manager.go index 5bf39920a74..922dd0db6f9 100644 --- a/sdk/cliproxy/usage/manager.go +++ b/sdk/cliproxy/usage/manager.go @@ -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 } diff --git a/sdk/pluginapi/types.go b/sdk/pluginapi/types.go index 397097f4986..43784ba577f 100644 --- a/sdk/pluginapi/types.go +++ b/sdk/pluginapi/types.go @@ -1448,6 +1448,17 @@ type UsageRecord struct { Latency time.Duration // TTFT is the time to first token for streaming requests. TTFT time.Duration + // 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 reports whether the request failed. Failed bool // Failure contains failure details when Failed is true.