From 14ade769563aba6008a8689cf690a48385ef5b28 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=96=E7=95=8C?= Date: Sun, 15 Mar 2026 19:09:29 +0800 Subject: [PATCH] ccm,ocm: remove dead code, fix timer leaks, eliminate redundant lookups - Remove unused onBecameUnusable field from CCM credential structs (OCM wires it for WebSocket interruption; CCM has no equivalent) - Replace time.After with time.NewTimer in doHTTPWithRetry and connectorLoop to avoid timer leaks on context cancellation - Pass already-resolved provider to rewriteResponseHeadersForExternalUser instead of re-resolving via credentialForUser - Hoist reverseYamuxConfig to package-level var (immutable, no need to allocate on every call) --- service/ccm/credential.go | 4 +++- service/ccm/credential_default.go | 6 +----- service/ccm/credential_external.go | 6 +----- service/ccm/reverse.go | 12 +++++++----- service/ccm/service_handler.go | 2 +- service/ccm/service_status.go | 7 +------ service/ocm/credential.go | 4 +++- service/ocm/reverse.go | 12 +++++++----- service/ocm/service_handler.go | 2 +- service/ocm/service_status.go | 7 +------ service/ocm/service_websocket.go | 2 +- 11 files changed, 27 insertions(+), 37 deletions(-) diff --git a/service/ccm/credential.go b/service/ccm/credential.go index b73261717..9e1614166 100644 --- a/service/ccm/credential.go +++ b/service/ccm/credential.go @@ -26,10 +26,12 @@ func doHTTPWithRetry(ctx context.Context, client *http.Client, buildRequest func for attempt := range httpRetryMaxAttempts { if attempt > 0 { delay := httpRetryInitialDelay * time.Duration(1<<(attempt-1)) + timer := time.NewTimer(delay) select { case <-ctx.Done(): + timer.Stop() return nil, lastError - case <-time.After(delay): + case <-timer.C: } } request, err := buildRequest() diff --git a/service/ccm/credential_default.go b/service/ccm/credential_default.go index 0bbf18874..021df5d27 100644 --- a/service/ccm/credential_default.go +++ b/service/ccm/credential_default.go @@ -44,8 +44,7 @@ type defaultCredential struct { watcherRetryAt time.Time // Connection interruption - onBecameUnusable func() - interrupted bool + interrupted bool requestContext context.Context cancelRequests context.CancelFunc requestAccess sync.Mutex @@ -353,9 +352,6 @@ func (c *defaultCredential) interruptConnections() { c.cancelRequests() c.requestContext, c.cancelRequests = context.WithCancel(context.Background()) c.requestAccess.Unlock() - if c.onBecameUnusable != nil { - c.onBecameUnusable() - } } func (c *defaultCredential) wrapRequestContext(parent context.Context) *credentialRequestContext { diff --git a/service/ccm/credential_external.go b/service/ccm/credential_external.go index 9a3ca1e16..21e695b24 100644 --- a/service/ccm/credential_external.go +++ b/service/ccm/credential_external.go @@ -40,8 +40,7 @@ type externalCredential struct { usageTracker *AggregatedUsage logger log.ContextLogger - onBecameUnusable func() - interrupted bool + interrupted bool requestContext context.Context cancelRequests context.CancelFunc requestAccess sync.Mutex @@ -494,9 +493,6 @@ func (c *externalCredential) interruptConnections() { c.cancelRequests() c.requestContext, c.cancelRequests = context.WithCancel(context.Background()) c.requestAccess.Unlock() - if c.onBecameUnusable != nil { - c.onBecameUnusable() - } } func (c *externalCredential) doPollUsageRequest(ctx context.Context) (*http.Response, error) { diff --git a/service/ccm/reverse.go b/service/ccm/reverse.go index b6b4c88a0..0da0c567c 100644 --- a/service/ccm/reverse.go +++ b/service/ccm/reverse.go @@ -17,14 +17,14 @@ import ( "github.com/hashicorp/yamux" ) -func reverseYamuxConfig() *yamux.Config { +var defaultYamuxConfig = func() *yamux.Config { config := yamux.DefaultConfig() config.KeepAliveInterval = 15 * time.Second config.ConnectionWriteTimeout = 10 * time.Second config.MaxStreamWindowSize = 512 * 1024 config.LogOutput = io.Discard return config -} +}() type bufferedConn struct { reader *bufio.Reader @@ -108,7 +108,7 @@ func (s *Service) handleReverseConnect(ctx context.Context, w http.ResponseWrite return } - session, err := yamux.Client(conn, reverseYamuxConfig()) + session, err := yamux.Client(conn, defaultYamuxConfig) if err != nil { conn.Close() s.logger.ErrorContext(ctx, "reverse connect: create yamux client for ", receiverCredential.tagName(), ": ", err) @@ -161,9 +161,11 @@ func (c *externalCredential) connectorLoop() { consecutiveFailures++ backoff := connectorBackoff(consecutiveFailures) c.logger.Warn("reverse connection for ", c.tag, " lost: ", err, ", reconnecting in ", backoff) + timer := time.NewTimer(backoff) select { - case <-time.After(backoff): + case <-timer.C: case <-ctx.Done(): + timer.Stop() return } } @@ -236,7 +238,7 @@ func (c *externalCredential) connectorConnect(ctx context.Context) (time.Duratio } } - session, err := yamux.Server(&bufferedConn{reader: reader, Conn: conn}, reverseYamuxConfig()) + session, err := yamux.Server(&bufferedConn{reader: reader, Conn: conn}, defaultYamuxConfig) if err != nil { conn.Close() return 0, E.Cause(err, "create yamux server") diff --git a/service/ccm/service_handler.go b/service/ccm/service_handler.go index fdbb68203..33d6317de 100644 --- a/service/ccm/service_handler.go +++ b/service/ccm/service_handler.go @@ -313,7 +313,7 @@ func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) { // Rewrite response headers for external users if userConfig != nil && userConfig.ExternalCredential != "" { - s.rewriteResponseHeadersForExternalUser(response.Header, userConfig) + s.rewriteResponseHeadersForExternalUser(response.Header, provider, userConfig) } for key, values := range response.Header { diff --git a/service/ccm/service_status.go b/service/ccm/service_status.go index bd8aa4b22..20c4780e5 100644 --- a/service/ccm/service_status.go +++ b/service/ccm/service_status.go @@ -100,12 +100,7 @@ func (s *Service) computeAggregatedUtilization(provider credentialProvider, user totalWeight } -func (s *Service) rewriteResponseHeadersForExternalUser(headers http.Header, userConfig *option.CCMUser) { - provider, err := credentialForUser(s.userConfigMap, s.providers, userConfig.Name) - if err != nil { - return - } - +func (s *Service) rewriteResponseHeadersForExternalUser(headers http.Header, provider credentialProvider, userConfig *option.CCMUser) { avgFiveHour, avgWeekly, totalWeight := s.computeAggregatedUtilization(provider, userConfig) headers.Set("anthropic-ratelimit-unified-5h-utilization", strconv.FormatFloat(avgFiveHour/100, 'f', 6, 64)) diff --git a/service/ocm/credential.go b/service/ocm/credential.go index 80c094cdb..8a56bfcec 100644 --- a/service/ocm/credential.go +++ b/service/ocm/credential.go @@ -29,10 +29,12 @@ func doHTTPWithRetry(ctx context.Context, client *http.Client, buildRequest func for attempt := range httpRetryMaxAttempts { if attempt > 0 { delay := httpRetryInitialDelay * time.Duration(1<<(attempt-1)) + timer := time.NewTimer(delay) select { case <-ctx.Done(): + timer.Stop() return nil, lastError - case <-time.After(delay): + case <-timer.C: } } request, err := buildRequest() diff --git a/service/ocm/reverse.go b/service/ocm/reverse.go index b47826e60..f97df5b87 100644 --- a/service/ocm/reverse.go +++ b/service/ocm/reverse.go @@ -17,14 +17,14 @@ import ( "github.com/hashicorp/yamux" ) -func reverseYamuxConfig() *yamux.Config { +var defaultYamuxConfig = func() *yamux.Config { config := yamux.DefaultConfig() config.KeepAliveInterval = 15 * time.Second config.ConnectionWriteTimeout = 10 * time.Second config.MaxStreamWindowSize = 512 * 1024 config.LogOutput = io.Discard return config -} +}() type bufferedConn struct { reader *bufio.Reader @@ -108,7 +108,7 @@ func (s *Service) handleReverseConnect(ctx context.Context, w http.ResponseWrite return } - session, err := yamux.Client(conn, reverseYamuxConfig()) + session, err := yamux.Client(conn, defaultYamuxConfig) if err != nil { conn.Close() s.logger.ErrorContext(ctx, "reverse connect: create yamux client for ", receiverCredential.tagName(), ": ", err) @@ -161,9 +161,11 @@ func (c *externalCredential) connectorLoop() { consecutiveFailures++ backoff := connectorBackoff(consecutiveFailures) c.logger.Warn("reverse connection for ", c.tag, " lost: ", err, ", reconnecting in ", backoff) + timer := time.NewTimer(backoff) select { - case <-time.After(backoff): + case <-timer.C: case <-ctx.Done(): + timer.Stop() return } } @@ -236,7 +238,7 @@ func (c *externalCredential) connectorConnect(ctx context.Context) (time.Duratio } } - session, err := yamux.Server(&bufferedConn{reader: reader, Conn: conn}, reverseYamuxConfig()) + session, err := yamux.Server(&bufferedConn{reader: reader, Conn: conn}, defaultYamuxConfig) if err != nil { conn.Close() return 0, E.Cause(err, "create yamux server") diff --git a/service/ocm/service_handler.go b/service/ocm/service_handler.go index 7c9242f5a..0a7698d00 100644 --- a/service/ocm/service_handler.go +++ b/service/ocm/service_handler.go @@ -293,7 +293,7 @@ func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) { // Rewrite response headers for external users if userConfig != nil && userConfig.ExternalCredential != "" { - s.rewriteResponseHeadersForExternalUser(response.Header, userConfig) + s.rewriteResponseHeadersForExternalUser(response.Header, provider, userConfig) } for key, values := range response.Header { diff --git a/service/ocm/service_status.go b/service/ocm/service_status.go index 327d3a2da..7057b60f7 100644 --- a/service/ocm/service_status.go +++ b/service/ocm/service_status.go @@ -100,12 +100,7 @@ func (s *Service) computeAggregatedUtilization(provider credentialProvider, user totalWeight } -func (s *Service) rewriteResponseHeadersForExternalUser(headers http.Header, userConfig *option.OCMUser) { - provider, err := credentialForUser(s.userConfigMap, s.providers, userConfig.Name) - if err != nil { - return - } - +func (s *Service) rewriteResponseHeadersForExternalUser(headers http.Header, provider credentialProvider, userConfig *option.OCMUser) { avgFiveHour, avgWeekly, totalWeight := s.computeAggregatedUtilization(provider, userConfig) activeLimitIdentifier := normalizeRateLimitIdentifier(headers.Get("x-codex-active-limit")) diff --git a/service/ocm/service_websocket.go b/service/ocm/service_websocket.go index ce20e1be7..25ded0367 100644 --- a/service/ocm/service_websocket.go +++ b/service/ocm/service_websocket.go @@ -253,7 +253,7 @@ func (s *Service) handleWebSocket( } } if userConfig != nil && userConfig.ExternalCredential != "" { - s.rewriteResponseHeadersForExternalUser(clientResponseHeaders, userConfig) + s.rewriteResponseHeadersForExternalUser(clientResponseHeaders, provider, userConfig) } clientUpgrader := ws.HTTPUpgrader{