diff --git a/internal/redisqueue/plugin.go b/internal/redisqueue/plugin.go index 0b3bf108b..4eb8280b6 100644 --- a/internal/redisqueue/plugin.go +++ b/internal/redisqueue/plugin.go @@ -65,6 +65,7 @@ func (p *usageQueuePlugin) HandleUsage(ctx context.Context, record coreusage.Rec serviceTier = coreusage.ServiceTierFromContext(ctx) } responseServiceTier := strings.TrimSpace(record.ResponseServiceTier) + responseModel := strings.TrimSpace(record.ResponseModel) clientRequestMetadata := internallogging.GetClientRequestMetadata(ctx) sessionID := strings.TrimSpace(record.SessionID) parentSessionID := strings.TrimSpace(record.ParentSessionID) @@ -138,6 +139,7 @@ func (p *usageQueuePlugin) HandleUsage(ctx context.Context, record coreusage.Rec ReasoningEffort: reasoningEffort, ServiceTier: serviceTier, ResponseServiceTier: responseServiceTier, + ResponseModel: responseModel, }) if err != nil { return @@ -162,6 +164,7 @@ type queuedUsageDetail struct { ReasoningEffort string `json:"reasoning_effort"` ServiceTier string `json:"service_tier"` ResponseServiceTier string `json:"response_service_tier,omitempty"` + ResponseModel string `json:"response_model,omitempty"` } type requestDetail struct { diff --git a/internal/redisqueue/plugin_test.go b/internal/redisqueue/plugin_test.go index 248f2a505..78c228565 100644 --- a/internal/redisqueue/plugin_test.go +++ b/internal/redisqueue/plugin_test.go @@ -43,6 +43,7 @@ func TestUsageQueuePluginPayloadIncludesStableFieldsAndSuccess(t *testing.T) { ReasoningEffort: "medium", ServiceTier: "auto", ResponseServiceTier: "default", + ResponseModel: "gpt-5.6-luna", Generate: coreusage.GenerateFlag(true), RequestedAt: time.Date(2026, 4, 25, 0, 0, 0, 0, time.UTC), Latency: 1500 * time.Millisecond, @@ -72,6 +73,7 @@ func TestUsageQueuePluginPayloadIncludesStableFieldsAndSuccess(t *testing.T) { requireStringField(t, payload, "service_tier", "auto") requireMissingField(t, payload, "request_service_tier") requireStringField(t, payload, "response_service_tier", "default") + requireStringField(t, payload, "response_model", "gpt-5.6-luna") requireIntField(t, payload, "accounting_version", coreusage.TokenAccountingSchemaVersion) requireTokenBreakdown(t, payload, coreusage.TokenAccountingQualityComplete, 30) requireTokensBoolField(t, payload, "cache_read_tokens_present", true) diff --git a/internal/runtime/executor/codex_executor_execute.go b/internal/runtime/executor/codex_executor_execute.go index 49ef1a127..ab5489247 100644 --- a/internal/runtime/executor/codex_executor_execute.go +++ b/internal/runtime/executor/codex_executor_execute.go @@ -140,6 +140,7 @@ func (e *CodexExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, re eventData := bytes.TrimSpace(line[5:]) eventData = helps.RestoreCodexMultiAgentV2Response(eventData, optimizeMultiAgentV2) + reporter.ObserveCodexResponseModel(eventData) eventType := gjson.GetBytes(eventData, "type").String() if helps.HasMeaningfulCodexOutputDelta(eventData) { diff --git a/internal/runtime/executor/codex_executor_terminal.go b/internal/runtime/executor/codex_executor_terminal.go index a83f01d22..b91b27694 100644 --- a/internal/runtime/executor/codex_executor_terminal.go +++ b/internal/runtime/executor/codex_executor_terminal.go @@ -591,9 +591,11 @@ func isCodexEmptyPart(payload []byte) bool { } } -// observeCodexTokenEvent inspects a stream payload and marks TTFT on the first substantive token event. +// observeCodexTokenEvent inspects a stream payload, marks TTFT on the first substantive +// token event, and records the model the upstream reports serving. func observeCodexTokenEvent(reporter *helps.UsageReporter, payload []byte) { helps.ObserveResponsesTokenEvent(reporter, payload) + reporter.ObserveCodexResponseModel(payload) } // newCodexBootstrapOverloadErr reports a buffered overload rejection with its real status. diff --git a/internal/runtime/executor/codex_openai_images.go b/internal/runtime/executor/codex_openai_images.go index 59ef7a261..d20b00a91 100644 --- a/internal/runtime/executor/codex_openai_images.go +++ b/internal/runtime/executor/codex_openai_images.go @@ -154,6 +154,7 @@ func (e *CodexExecutor) executeOpenAIImage(ctx context.Context, auth *cliproxyau continue } eventData := bytes.TrimSpace(line[len(dataTag):]) + reporter.ObserveCodexResponseModel(eventData) switch gjson.GetBytes(eventData, "type").String() { case "response.output_item.done": collectCodexOutputItemDone(eventData, outputItemsByIndex, &outputItemsFallback) @@ -278,6 +279,7 @@ func (e *CodexExecutor) executeOpenAIImageStream(ctx context.Context, auth *clip continue } eventData := bytes.TrimSpace(line[len(dataTag):]) + reporter.ObserveCodexResponseModel(eventData) switch gjson.GetBytes(eventData, "type").String() { case "response.output_item.done": collectCodexOutputItemDone(eventData, outputItemsByIndex, &outputItemsFallback) diff --git a/internal/runtime/executor/codex_response_model_test.go b/internal/runtime/executor/codex_response_model_test.go new file mode 100644 index 000000000..97719ddbd --- /dev/null +++ b/internal/runtime/executor/codex_response_model_test.go @@ -0,0 +1,107 @@ +package executor + +import ( + "bytes" + "context" + "testing" + "time" + + "github.com/router-for-me/CLIProxyAPI/v7/internal/config" + "github.com/router-for-me/CLIProxyAPI/v7/internal/runtime/executor/helps" + cliproxyauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" + coreusage "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/usage" +) + +const codexResponseModelTestStream = `event: response.created +data: {"type":"response.created","response":{"id":"resp_1","model":"gpt-5.6-luna"}} + +data: {"type":"response.output_text.delta","delta":"he"} + +data: {"type":"response.output_text.delta","delta":"llo"} + +data: {"type":"response.completed","response":{"id":"resp_1","model":"gpt-5.6-luna","usage":{"input_tokens":5,"output_tokens":7,"total_tokens":12}}} + +data: [DONE]` + +// TestObserveCodexTokenEventRecordsResponseModel guards the wiring between the codex +// stream loops and the usage reporter: every transport funnels through observeCodexTokenEvent. +func TestObserveCodexTokenEventRecordsResponseModel(t *testing.T) { + reporter := helps.NewExecutorUsageReporter(context.Background(), NewCodexExecutor(&config.Config{}), "gpt-6-astra", nil) + + for _, line := range bytes.Split([]byte(codexResponseModelTestStream), []byte("\n")) { + line = bytes.TrimSpace(line) + if !bytes.HasPrefix(line, dataTag) { + continue + } + observeCodexTokenEvent(reporter, bytes.TrimSpace(line[len(dataTag):])) + } + + if got := reporter.ResponseModel(); got != "gpt-5.6-luna" { + t.Fatalf("reporter response model = %q, want %q", got, "gpt-5.6-luna") + } +} + +type codexResponseModelUsageCapture struct { + alias string + records chan coreusage.Record +} + +func (c *codexResponseModelUsageCapture) HandleUsage(_ context.Context, record coreusage.Record) { + if record.Alias != c.alias { + return + } + select { + case c.records <- record: + default: + } +} + +type codexResponseModelNoopUsagePlugin struct{} + +func (codexResponseModelNoopUsagePlugin) HandleUsage(context.Context, coreusage.Record) {} + +func (c *codexResponseModelUsageCapture) await(t *testing.T) coreusage.Record { + t.Helper() + select { + case record := <-c.records: + return record + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for the usage record") + return coreusage.Record{} + } +} + +// TestCodexUsageRecordsCarryResponseModelPerModel checks that the attempt record reports +// the served model while the image generation tool record must not claim it. +func TestCodexUsageRecordsCarryResponseModelPerModel(t *testing.T) { + const alias = "codex-response-model-wiring-test" + capture := &codexResponseModelUsageCapture{alias: alias, records: make(chan coreusage.Record, 4)} + coreusage.RegisterNamedPlugin(t.Name(), capture) + t.Cleanup(func() { + coreusage.RegisterNamedPlugin(t.Name(), codexResponseModelNoopUsagePlugin{}) + }) + + ctx := coreusage.WithRequestedModelAlias(context.Background(), alias) + auth := &cliproxyauth.Auth{ID: "codex-auth-1", Index: "auth-index-7", Provider: "codex"} + reporter := helps.NewExecutorUsageReporter(ctx, NewCodexExecutor(&config.Config{}), "gpt-5.4-mini", auth) + + observeCodexTokenEvent(reporter, []byte(`{"type":"response.completed","response":{"model":"gpt-5.4-mini","usage":{"total_tokens":12}}}`)) + reporter.EnsurePublished(ctx) + reporter.PublishAdditionalModel(ctx, "gpt-image-1.5", coreusage.Detail{TotalTokens: 5}) + + attemptRecord := capture.await(t) + if attemptRecord.Model != "gpt-5.4-mini" { + t.Fatalf("attempt record model = %q, want %q", attemptRecord.Model, "gpt-5.4-mini") + } + if attemptRecord.ResponseModel != "gpt-5.4-mini" { + t.Fatalf("attempt record response model = %q, want %q", attemptRecord.ResponseModel, "gpt-5.4-mini") + } + + imageRecord := capture.await(t) + if imageRecord.Model != "gpt-image-1.5" { + t.Fatalf("image record model = %q, want %q", imageRecord.Model, "gpt-image-1.5") + } + if imageRecord.ResponseModel != "" { + t.Fatalf("image record response model = %q, want empty", imageRecord.ResponseModel) + } +} diff --git a/internal/runtime/executor/helps/response_model.go b/internal/runtime/executor/helps/response_model.go new file mode 100644 index 000000000..1c2eb1d57 --- /dev/null +++ b/internal/runtime/executor/helps/response_model.go @@ -0,0 +1,181 @@ +package helps + +import ( + "strings" + "sync" + "time" + + "github.com/router-for-me/CLIProxyAPI/v7/internal/thinking" + "github.com/tidwall/gjson" +) + +const ( + // maxCodexResponseModelLength is a defensive bound on an upstream-controlled string + // reaching logs and usage records; known codex model ids stay under ~30 bytes. + maxCodexResponseModelLength = 128 + + // codexModelSubstitutionWarnWindow bounds how often one credential and model pair + // warns: on an affected credential every request is substituted. + codexModelSubstitutionWarnWindow = 10 * time.Minute + + // codexModelSubstitutionWarnMaxEntries caps the throttle state, naturally bounded by + // credentials times models; memory safety wins over perfect throttling. + codexModelSubstitutionWarnMaxEntries = 1024 +) + +// extractCodexResponseModelEvent returns the model a codex upstream reports serving, read +// from a raw JSON frame or an SSE line, and whether the event terminates the response. +func extractCodexResponseModelEvent(payload []byte) (model string, terminal bool) { + data := jsonPayload(payload) + if len(data) == 0 { + return "", false + } + // The event type is checked before the payload is validated, so the hot path + // (output deltas) stays a single cheap lookup. + carriesModel, terminal := codexResponseModelEventKind(gjson.GetBytes(data, "type").String()) + if !carriesModel { + return "", false + } + if !gjson.ValidBytes(data) { + return "", false + } + // The value is upstream-controlled: reject non-string and oversized names + // rather than propagating them into logs and usage records. + modelResult := gjson.GetBytes(data, "response.model") + if modelResult.Type != gjson.String { + return "", terminal + } + model = strings.TrimSpace(modelResult.String()) + if len(model) > maxCodexResponseModelLength { + return "", terminal + } + return model, terminal +} + +// codexResponseModelEventKind reports whether a codex event embeds the authoritative +// response object, and whether that event terminates the response. +func codexResponseModelEventKind(eventType string) (carriesModel bool, terminal bool) { + switch strings.TrimSpace(eventType) { + case "response.created", "response.in_progress": + return true, false + case "response.completed", "response.incomplete", "response.done": + return true, true + default: + return false, false + } +} + +// normalizeCodexModelName lower-cases a model id and drops its thinking suffix, +// which never reaches the upstream request body. +func normalizeCodexModelName(model string) string { + return strings.TrimSpace(thinking.ParseSuffix(strings.ToLower(strings.TrimSpace(model))).ModelName) +} + +// IsCodexModelSubstituted reports whether the upstream served a model other than the +// requested one; a dated alias pins a snapshot of the same model and is accepted. +func IsCodexModelSubstituted(requested, served string) bool { + servedModel := normalizeCodexModelName(served) + if servedModel == "" { + return false + } + requestedModel := normalizeCodexModelName(requested) + if requestedModel == "" { + return false + } + if requestedModel == servedModel { + return false + } + return !isCodexDatedModelAlias(requestedModel, servedModel) && + !isCodexDatedModelAlias(servedModel, requestedModel) +} + +// isCodexDatedModelAlias reports whether dated is base plus a release date suffix, +// which upstreams use to pin the exact snapshot of the same model. +func isCodexDatedModelAlias(base, dated string) bool { + prefix := base + "-" + if !strings.HasPrefix(dated, prefix) { + return false + } + return isCodexModelDateSuffix(dated[len(prefix):]) +} + +// isCodexModelDateSuffix reports whether suffix is a YYYY-MM-DD or YYYYMMDD date. +func isCodexModelDateSuffix(suffix string) bool { + switch len(suffix) { + case len("YYYY-MM-DD"): + if suffix[4] != '-' || suffix[7] != '-' { + return false + } + return isCodexModelDigits(suffix[:4]) && isCodexModelDigits(suffix[5:7]) && isCodexModelDigits(suffix[8:]) + case len("YYYYMMDD"): + return isCodexModelDigits(suffix) + default: + return false + } +} + +// isCodexModelDigits reports whether value is a non-empty run of ASCII digits. +func isCodexModelDigits(value string) bool { + if value == "" { + return false + } + for i := 0; i < len(value); i++ { + if value[i] < '0' || value[i] > '9' { + return false + } + } + return true +} + +type codexModelSubstitutionKey struct { + authID string + requested string + served string +} + +// codexModelSubstitutionThrottle records the last warning per key; nowFunc is +// injectable so tests can advance the window without sleeping. +type codexModelSubstitutionThrottle struct { + mu sync.Mutex + nowFunc func() time.Time + lastWarn map[codexModelSubstitutionKey]time.Time +} + +func newCodexModelSubstitutionThrottle(nowFunc func() time.Time) *codexModelSubstitutionThrottle { + if nowFunc == nil { + nowFunc = time.Now + } + return &codexModelSubstitutionThrottle{ + nowFunc: nowFunc, + lastWarn: make(map[codexModelSubstitutionKey]time.Time), + } +} + +// allow reports whether the key may emit a warning now and records the decision. +// Suppressed repeats stay silent instead of moving to a lower level. +func (t *codexModelSubstitutionThrottle) allow(key codexModelSubstitutionKey) bool { + if t == nil { + return true + } + t.mu.Lock() + defer t.mu.Unlock() + now := t.nowFunc() + if last, ok := t.lastWarn[key]; ok && now.Sub(last) < codexModelSubstitutionWarnWindow { + return false + } + if len(t.lastWarn) >= codexModelSubstitutionWarnMaxEntries { + for storedKey, storedAt := range t.lastWarn { + if now.Sub(storedAt) >= codexModelSubstitutionWarnWindow { + delete(t.lastWarn, storedKey) + } + } + if len(t.lastWarn) >= codexModelSubstitutionWarnMaxEntries { + clear(t.lastWarn) + } + } + t.lastWarn[key] = now + return true +} + +// codexModelSubstitutionWarns throttles substitution warnings process-wide. +var codexModelSubstitutionWarns = newCodexModelSubstitutionThrottle(time.Now) diff --git a/internal/runtime/executor/helps/response_model_test.go b/internal/runtime/executor/helps/response_model_test.go new file mode 100644 index 000000000..5d456717d --- /dev/null +++ b/internal/runtime/executor/helps/response_model_test.go @@ -0,0 +1,479 @@ +package helps + +import ( + "context" + "strconv" + "strings" + "sync" + "testing" + "time" + + cliproxyauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" + "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/usage" + log "github.com/sirupsen/logrus" + logtest "github.com/sirupsen/logrus/hooks/test" +) + +func TestExtractCodexResponseModel(t *testing.T) { + tests := []struct { + name string + payload string + want string + }{ + { + name: "sse response created", + payload: `data: {"type":"response.created","response":{"id":"resp_1","model":"gpt-5.6-luna"}}`, + want: "gpt-5.6-luna", + }, + { + name: "sse response created without space after data prefix", + payload: `data:{"type":"response.created","response":{"model":"gpt-6-astra"}}`, + want: "gpt-6-astra", + }, + { + name: "raw json response completed", + payload: `{"type":"response.completed","response":{"model":"gpt-5.6-luna","usage":{"total_tokens":12}}}`, + want: "gpt-5.6-luna", + }, + { + name: "response in progress", + payload: `{"type":"response.in_progress","response":{"model":"gpt-5.6-terra"}}`, + want: "gpt-5.6-terra", + }, + { + name: "response incomplete", + payload: `{"type":"response.incomplete","response":{"model":"gpt-5.6-sol"}}`, + want: "gpt-5.6-sol", + }, + { + name: "response done", + payload: `{"type":"response.done","response":{"model":"gpt-5.3-codex-spark"}}`, + want: "gpt-5.3-codex-spark", + }, + { + name: "model value is trimmed", + payload: `{"type":"response.created","response":{"model":" gpt-6-astra "}}`, + want: "gpt-6-astra", + }, + { + name: "output text delta is ignored", + payload: `data: {"type":"response.output_text.delta","delta":"hi","response":{"model":"gpt-5.6-luna"}}`, + want: "", + }, + { + name: "output item done is ignored", + payload: `{"type":"response.output_item.done","item":{"type":"message"}}`, + want: "", + }, + { + name: "rate limits event is ignored", + payload: `{"type":"codex.rate_limits","rate_limits":{"primary":{"used_percent":12}}}`, + want: "", + }, + { + name: "terminal failure event is ignored", + payload: `{"type":"error","error":{"code":"server_is_overloaded"}}`, + want: "", + }, + { + // /responses/compact answers with a compaction object that carries no + // event type and no response.model, so that route reports nothing. + name: "responses compact object", + payload: `{"id":"resp_1","object":"response.compaction","usage":{"input_tokens":1,"output_tokens":2,"total_tokens":3}}`, + want: "", + }, + { + // The direct /images/generations and /images/edits endpoints answer in the + // Images API shape, which never names the model that served the request. + name: "images api response", + payload: `{"created":1745539200,"data":[{"b64_json":"aGk="}]}`, + want: "", + }, + { + // Direct image streams emit image_generation.* events only. + name: "image generation completed event", + payload: `data: {"type":"image_generation.completed","b64_json":"aGk="}`, + want: "", + }, + { + name: "non string model value", + payload: `{"type":"response.created","response":{"model":123}}`, + want: "", + }, + { + name: "object model value", + payload: `{"type":"response.created","response":{"model":{"id":"gpt-5.6-luna"}}}`, + want: "", + }, + { + name: "oversized model value", + payload: `{"type":"response.created","response":{"model":"` + strings.Repeat("m", maxCodexResponseModelLength+1) + `"}}`, + want: "", + }, + { + name: "model value at the length limit", + payload: `{"type":"response.created","response":{"model":"` + strings.Repeat("m", maxCodexResponseModelLength) + `"}}`, + want: strings.Repeat("m", maxCodexResponseModelLength), + }, + { + name: "done marker", + payload: "data: [DONE]", + want: "", + }, + { + name: "sse event line", + payload: "event: response.created", + want: "", + }, + { + name: "empty payload", + payload: "", + want: "", + }, + { + name: "malformed json", + payload: `{"type":"response.created","response":{"model":"gpt-5.6-luna"`, + want: "", + }, + { + name: "response without model", + payload: `{"type":"response.created","response":{"id":"resp_1"}}`, + want: "", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got, _ := extractCodexResponseModelEvent([]byte(tt.payload)); got != tt.want { + t.Fatalf("extractCodexResponseModelEvent(%q) = %q, want %q", tt.payload, got, tt.want) + } + }) + } +} + +func TestIsCodexModelSubstituted(t *testing.T) { + tests := []struct { + name string + requested string + served string + want bool + }{ + {name: "silent substitution", requested: "gpt-6-astra", served: "gpt-5.6-luna", want: true}, + {name: "same model", requested: "gpt-6-astra", served: "gpt-6-astra", want: false}, + {name: "thinking suffix stripped from requested", requested: "gpt-6-astra(high)", served: "gpt-6-astra", want: false}, + {name: "thinking suffix stripped from served", requested: "gpt-6-astra", served: "gpt-6-astra(high)", want: false}, + {name: "thinking suffix on both sides", requested: "gpt-6-astra(high)", served: "gpt-6-astra(low)", want: false}, + {name: "thinking suffix with substitution", requested: "gpt-6-astra(high)", served: "gpt-5.6-luna", want: true}, + {name: "numeric thinking suffix", requested: "gpt-5.6-terra(16384)", served: "gpt-5.6-terra", want: false}, + {name: "case insensitive requested", requested: "GPT-6-Astra", served: "gpt-6-astra", want: false}, + {name: "case insensitive served", requested: "gpt-6-astra", served: " GPT-6-ASTRA ", want: false}, + {name: "served pins dashed date", requested: "gpt-5.6-terra", served: "gpt-5.6-terra-2026-05-13", want: false}, + {name: "served pins compact date", requested: "gpt-5.6-sol", served: "gpt-5.6-sol-20260513", want: false}, + {name: "requested pins dashed date", requested: "gpt-5.6-terra-2026-05-13", served: "gpt-5.6-terra", want: false}, + {name: "requested pins compact date", requested: "gpt-5.6-sol-20260513", served: "gpt-5.6-sol", want: false}, + {name: "different dates on both sides", requested: "gpt-5.6-terra-2026-05-13", served: "gpt-5.6-terra-2026-06-01", want: true}, + {name: "incomplete date suffix", requested: "gpt-5.6-sol", served: "gpt-5.6-sol-2026-5-13", want: true}, + {name: "non date suffix", requested: "gpt-5.5", served: "gpt-5.5-codex", want: true}, + {name: "date suffix with trailing tag", requested: "gpt-5.6-luna", served: "gpt-5.6-luna-2026-05-13-preview", want: true}, + {name: "spark unchanged", requested: "gpt-5.3-codex-spark", served: "gpt-5.3-codex-spark", want: false}, + {name: "review model unchanged", requested: "codex-auto-review", served: "codex-auto-review", want: false}, + {name: "review model substituted", requested: "codex-auto-review", served: "gpt-5.6-luna", want: true}, + {name: "empty served", requested: "gpt-6-astra", served: "", want: false}, + {name: "blank served", requested: "gpt-6-astra", served: " ", want: false}, + {name: "empty requested", requested: "", served: "gpt-5.6-luna", want: false}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := IsCodexModelSubstituted(tt.requested, tt.served); got != tt.want { + t.Fatalf("IsCodexModelSubstituted(%q, %q) = %v, want %v", tt.requested, tt.served, got, tt.want) + } + }) + } +} + +// setupResponseModelLoggerHook captures warnings from the process-global logrus logger +// and gives the test its own throttle; both are shared, so callers must not use t.Parallel. +func setupResponseModelLoggerHook(t *testing.T) *logtest.Hook { + t.Helper() + logger := log.StandardLogger() + oldLevel := logger.GetLevel() + logger.SetLevel(log.WarnLevel) + savedHooks := logger.ReplaceHooks(make(log.LevelHooks)) + hook := new(logtest.Hook) + logger.AddHook(hook) + savedThrottle := codexModelSubstitutionWarns + codexModelSubstitutionWarns = newCodexModelSubstitutionThrottle(nil) + t.Cleanup(func() { + codexModelSubstitutionWarns = savedThrottle + logger.SetLevel(oldLevel) + logger.ReplaceHooks(savedHooks) + }) + return hook +} + +// codexUsageTestExecutor mirrors the production reporter construction path. +type codexUsageTestExecutor struct{} + +func (codexUsageTestExecutor) Identifier() string { return "codex" } + +func newCodexTestReporter(ctx context.Context, model string, auth *cliproxyauth.Auth) *UsageReporter { + return NewExecutorUsageReporter(ctx, codexUsageTestExecutor{}, model, auth) +} + +func substitutionWarnings(hook *logtest.Hook) []string { + var messages []string + for _, entry := range hook.AllEntries() { + if entry.Level == log.WarnLevel && strings.Contains(entry.Message, "upstream served model") { + messages = append(messages, entry.Message) + } + } + return messages +} + +const codexSubstitutedStream = `data: {"type":"response.created","response":{"id":"resp_1","model":"gpt-5.6-luna"}} +data: {"type":"response.in_progress","response":{"id":"resp_1","model":"gpt-5.6-luna"}} +data: {"type":"response.output_text.delta","delta":"hello"} +data: {"type":"response.completed","response":{"id":"resp_1","model":"gpt-5.6-luna","usage":{"input_tokens":5,"output_tokens":7,"total_tokens":12}}} +data: [DONE]` + +func newSubstitutionTestAuth() *cliproxyauth.Auth { + return &cliproxyauth.Auth{ + ID: "codex-auth-1", + Index: "auth-index-7", + Provider: "codex", + FileName: "/auths/codex-user@example.com.json", + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterRecordsSubstitutedCodexResponseModelAndWarnsOnce(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + reporter := newCodexTestReporter(ctx, "gpt-6-astra", newSubstitutionTestAuth()) + + var detail usage.Detail + for _, line := range strings.Split(codexSubstitutedStream, "\n") { + reporter.ObserveCodexResponseModel([]byte(line)) + if parsed, ok := ParseCodexUsage(JSONPayload([]byte(line))); ok { + detail = parsed + } + } + + if got := reporter.ResponseModel(); got != "gpt-5.6-luna" { + t.Fatalf("reporter response model = %q, want %q", got, "gpt-5.6-luna") + } + record := reporter.buildRecord(detail, false) + if record.ResponseModel != "gpt-5.6-luna" { + t.Fatalf("record response model = %q, want %q", record.ResponseModel, "gpt-5.6-luna") + } + if record.Model != "gpt-6-astra" { + t.Fatalf("record model = %q, want %q", record.Model, "gpt-6-astra") + } + // Observing events must not log anything; the warning belongs to the attempt. + if warnings := substitutionWarnings(hook); len(warnings) != 0 { + t.Fatalf("expected no warning before the attempt published, got %#v", warnings) + } + + reporter.Publish(ctx, detail) + reporter.EnsurePublished(ctx) + + warnings := substitutionWarnings(hook) + if len(warnings) != 1 { + t.Fatalf("substitution warnings = %d, want 1: %#v", len(warnings), warnings) + } + want := `codex executor: upstream served model "gpt-5.6-luna" for requested model "gpt-6-astra" (auth_index=auth-index-7)` + if warnings[0] != want { + t.Fatalf("warning = %q, want %q", warnings[0], want) + } + // The credential file name carries the account e-mail and must stay out of logs. + if strings.Contains(warnings[0], "example.com") || strings.Contains(warnings[0], "auth_file") { + t.Fatalf("warning leaked credential details: %q", warnings[0]) + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterDoesNotWarnWhenCodexResponseModelMatches(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + reporter := newCodexTestReporter(ctx, "gpt-6-astra(high)", nil) + + reporter.ObserveCodexResponseModel([]byte(`data: {"type":"response.created","response":{"model":"gpt-6-astra-2026-05-13"}}`)) + reporter.Publish(ctx, usage.Detail{TotalTokens: 3}) + + if got := reporter.ResponseModel(); got != "gpt-6-astra-2026-05-13" { + t.Fatalf("reporter response model = %q, want %q", got, "gpt-6-astra-2026-05-13") + } + if warnings := substitutionWarnings(hook); len(warnings) != 0 { + t.Fatalf("expected no substitution warning, got %#v", warnings) + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterIgnoresPayloadsWithoutResponseModel(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + reporter := newCodexTestReporter(ctx, "gpt-6-astra", nil) + + reporter.ObserveCodexResponseModel([]byte(`data: {"type":"response.created","response":{"model":"gpt-5.6-luna"}}`)) + reporter.ObserveCodexResponseModel([]byte(`data: {"type":"response.output_text.delta","delta":"hello"}`)) + reporter.Publish(ctx, usage.Detail{TotalTokens: 3}) + + if got := reporter.ResponseModel(); got != "gpt-5.6-luna" { + t.Fatalf("reporter response model = %q, want %q", got, "gpt-5.6-luna") + } + warnings := substitutionWarnings(hook) + if len(warnings) != 1 { + t.Fatalf("substitution warnings = %d, want 1: %#v", len(warnings), warnings) + } + // A reporter without a credential must still produce a usable label. + if !strings.HasSuffix(warnings[0], "(auth_index=nil)") { + t.Fatalf("warning = %q, want an auth_index=nil suffix", warnings[0]) + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterWarnsOnceUnderConcurrentObservationsAndPublishes(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + reporter := newCodexTestReporter(ctx, "gpt-6-astra", newSubstitutionTestAuth()) + created := []byte(`data: {"type":"response.created","response":{"model":"gpt-5.6-luna"}}`) + completed := []byte(`data: {"type":"response.completed","response":{"model":"gpt-5.6-luna","usage":{"total_tokens":4}}}`) + + const workers = 32 + start := make(chan struct{}) + var wg sync.WaitGroup + wg.Add(workers) + for i := range workers { + go func(worker int) { + defer wg.Done() + <-start + if worker%2 == 0 { + reporter.ObserveCodexResponseModel(created) + } else { + reporter.ObserveCodexResponseModel(completed) + } + reporter.EnsurePublished(ctx) + }(i) + } + close(start) + wg.Wait() + + if got := reporter.ResponseModel(); got != "gpt-5.6-luna" { + t.Fatalf("reporter response model = %q, want %q", got, "gpt-5.6-luna") + } + if warnings := substitutionWarnings(hook); len(warnings) != 1 { + t.Fatalf("substitution warnings = %d, want 1: %#v", len(warnings), warnings) + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterThrottlesRepeatedSubstitutionWarnings(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + now := time.Unix(1_700_000_000, 0) + codexModelSubstitutionWarns = newCodexModelSubstitutionThrottle(func() time.Time { return now }) + + publishSubstitutedAttempt := func(authID, authIndex string) { + auth := &cliproxyauth.Auth{ID: authID, Index: authIndex, Provider: "codex"} + reporter := newCodexTestReporter(ctx, "gpt-6-astra", auth) + reporter.ObserveCodexResponseModel([]byte(`{"type":"response.completed","response":{"model":"gpt-5.6-luna"}}`)) + reporter.Publish(ctx, usage.Detail{TotalTokens: 3}) + } + + publishSubstitutedAttempt("codex-auth-1", "auth-index-7") + publishSubstitutedAttempt("codex-auth-1", "auth-index-7") + if warnings := substitutionWarnings(hook); len(warnings) != 1 { + t.Fatalf("substitution warnings = %d, want 1 inside the window: %#v", len(warnings), warnings) + } + + // A different credential is an independent signal. + publishSubstitutedAttempt("codex-auth-2", "auth-index-8") + if warnings := substitutionWarnings(hook); len(warnings) != 2 { + t.Fatalf("substitution warnings = %d, want 2 after a second credential: %#v", len(warnings), warnings) + } + + now = now.Add(codexModelSubstitutionWarnWindow - time.Second) + publishSubstitutedAttempt("codex-auth-1", "auth-index-7") + if warnings := substitutionWarnings(hook); len(warnings) != 2 { + t.Fatalf("substitution warnings = %d, want 2 before the window elapsed: %#v", len(warnings), warnings) + } + + now = now.Add(time.Second) + publishSubstitutedAttempt("codex-auth-1", "auth-index-7") + if warnings := substitutionWarnings(hook); len(warnings) != 3 { + t.Fatalf("substitution warnings = %d, want 3 once the window elapsed: %#v", len(warnings), warnings) + } +} + +// Relies on the global logrus logger and warning throttle; do not add t.Parallel. +func TestUsageReporterThrottlesSubstitutionWarningsAcrossServedModelCase(t *testing.T) { + hook := setupResponseModelLoggerHook(t) + ctx := context.Background() + now := time.Unix(1_700_000_000, 0) + codexModelSubstitutionWarns = newCodexModelSubstitutionThrottle(func() time.Time { return now }) + + auth := &cliproxyauth.Auth{ID: "codex-auth-1", Index: "auth-index-7", Provider: "codex"} + for _, served := range []string{"gpt-5.6-luna", "GPT-5.6-LUNA"} { + reporter := newCodexTestReporter(ctx, "gpt-6-astra", auth) + reporter.ObserveCodexResponseModel([]byte(`{"type":"response.completed","response":{"model":"` + served + `"}}`)) + reporter.Publish(ctx, usage.Detail{TotalTokens: 3}) + } + + // Both attempts describe one substitution pair, so the casing must not open a + // second throttle window while the WARN keeps the raw upstream name. + warnings := substitutionWarnings(hook) + if len(warnings) != 1 { + t.Fatalf("substitution warnings = %d, want 1 for one normalized pair: %#v", len(warnings), warnings) + } + want := `codex executor: upstream served model "gpt-5.6-luna" for requested model "gpt-6-astra" (auth_index=auth-index-7)` + if warnings[0] != want { + t.Fatalf("warning = %q, want %q", warnings[0], want) + } +} + +func TestCodexModelSubstitutionThrottleBoundsStoredEntries(t *testing.T) { + now := time.Unix(1_700_000_000, 0) + throttle := newCodexModelSubstitutionThrottle(func() time.Time { return now }) + + for i := range codexModelSubstitutionWarnMaxEntries + 16 { + key := codexModelSubstitutionKey{ + authID: "codex-auth-" + strconv.Itoa(i), + requested: "gpt-6-astra", + served: "gpt-5.6-luna", + } + if !throttle.allow(key) { + t.Fatalf("first warning for key %d was suppressed", i) + } + } + throttle.mu.Lock() + stored := len(throttle.lastWarn) + throttle.mu.Unlock() + if stored > codexModelSubstitutionWarnMaxEntries { + t.Fatalf("stored throttle entries = %d, want at most %d", stored, codexModelSubstitutionWarnMaxEntries) + } +} + +func TestUsageReporterAdditionalModelRecordOmitsResponseModel(t *testing.T) { + ctx := context.Background() + reporter := newCodexTestReporter(ctx, "gpt-5.4-mini", nil) + reporter.ObserveCodexResponseModel([]byte(`data: {"type":"response.completed","response":{"model":"gpt-5.4-mini"}}`)) + + mainRecord := reporter.buildRecord(usage.Detail{TotalTokens: 12}, false) + if mainRecord.ResponseModel != "gpt-5.4-mini" { + t.Fatalf("main record response model = %q, want %q", mainRecord.ResponseModel, "gpt-5.4-mini") + } + + additionalRecord, ok := reporter.buildAdditionalModelRecord("gpt-image-1.5", usage.Detail{TotalTokens: 5}) + if !ok { + t.Fatal("expected an additional model record") + } + if additionalRecord.Model != "gpt-image-1.5" { + t.Fatalf("additional record model = %q, want %q", additionalRecord.Model, "gpt-image-1.5") + } + // The upstream response model describes the main text model, never the image + // generation tool model, so consumers must not see a fake substitution here. + if additionalRecord.ResponseModel != "" { + t.Fatalf("additional record response model = %q, want empty", additionalRecord.ResponseModel) + } +} diff --git a/internal/runtime/executor/helps/usage_helpers.go b/internal/runtime/executor/helps/usage_helpers.go index c2a43db29..2bd22fb30 100644 --- a/internal/runtime/executor/helps/usage_helpers.go +++ b/internal/runtime/executor/helps/usage_helpers.go @@ -10,6 +10,7 @@ import ( "reflect" "strings" "sync" + "sync/atomic" "time" "github.com/gin-gonic/gin" @@ -50,6 +51,13 @@ type UsageReporter struct { ttftStart time.Time ttftSet bool once sync.Once + + responseModelMu sync.RWMutex + // responseModel holds the latest model name reported by the upstream response. + responseModel string + // responseModelFinal marks that a terminal event already reported the served + // model, so later frames skip parsing entirely. + responseModelFinal atomic.Bool } type usageExecutor interface { @@ -175,6 +183,69 @@ func (r *UsageReporter) accessTokenFingerprint() string { return r.accessTokenHash } +// ObserveCodexResponseModel stores the model reported by a codex upstream event and +// ignores payloads without one; the substitution warning is emitted at publish time. +func (r *UsageReporter) ObserveCodexResponseModel(payload []byte) { + if r == nil || r.responseModelFinal.Load() { + return + } + served, terminal := extractCodexResponseModelEvent(payload) + if served == "" { + return + } + r.responseModelMu.Lock() + r.responseModel = served + r.responseModelMu.Unlock() + if terminal { + r.responseModelFinal.Store(true) + } +} + +// warnCodexModelSubstitution warns about a silent upstream model swap, throttled per +// credential and model pair, and labels the credential by index only, never by account. +func (r *UsageReporter) warnCodexModelSubstitution(ctx context.Context) { + if r == nil { + return + } + served := r.ResponseModel() + if served == "" || !IsCodexModelSubstituted(r.model, served) { + return + } + // The throttle key uses the same normalized names as the substitution check, so + // aliases of one pair share a window instead of each warning on its own. + requested := normalizeCodexModelName(r.model) + servedNormalized := normalizeCodexModelName(served) + if !codexModelSubstitutionWarns.allow(codexModelSubstitutionKey{ + authID: r.authID, + requested: requested, + served: servedNormalized, + }) { + return + } + LogWithRequestID(ctx).Warnf("codex executor: upstream served model %q for requested model %q (auth_index=%s)", served, r.model, r.authIndexForLog()) +} + +// authIndexForLog labels the credential without exposing its file name or account. +func (r *UsageReporter) authIndexForLog() string { + if r == nil { + return "nil" + } + if authIndex := strings.TrimSpace(r.authIndex); authIndex != "" { + return authIndex + } + return "nil" +} + +// ResponseModel returns the latest model reported by the upstream response. +func (r *UsageReporter) ResponseModel() string { + if r == nil { + return "" + } + r.responseModelMu.RLock() + defer r.responseModelMu.RUnlock() + return r.responseModel +} + func ExecutorTypeName(executor any) string { if executor == nil { return "" @@ -392,7 +463,7 @@ func (r *UsageReporter) publishWithOutcome(ctx context.Context, detail usage.Det } detail = normalizeUsageDetailTotal(detail, r.provider, r.executorType) r.once.Do(func() { - r.publishRecord(ctx, r.buildRecord(detail, failed, fail)) + r.publishAttemptRecord(ctx, r.buildRecord(detail, failed, fail)) }) } @@ -420,10 +491,17 @@ func (r *UsageReporter) EnsurePublished(ctx context.Context) { return } r.once.Do(func() { - r.publishRecord(ctx, r.buildRecord(usage.Detail{}, false, usage.Failure{})) + r.publishAttemptRecord(ctx, r.buildRecord(usage.Detail{}, false, usage.Failure{})) }) } +// publishAttemptRecord emits the record for one upstream attempt and the +// observability warnings that belong to the attempt rather than to a single event. +func (r *UsageReporter) publishAttemptRecord(ctx context.Context, record usage.Record) { + r.publishRecord(ctx, record) + r.warnCodexModelSubstitution(ctx) +} + func (r *UsageReporter) publishRecord(ctx context.Context, record usage.Record) { record.ResponseHeaders = internallogging.GetResponseHeaders(ctx) usage.PublishRecord(ctx, record) @@ -444,6 +522,12 @@ func (r *UsageReporter) buildRecordForModel(model string, detail usage.Detail, f if r == nil { return usage.Record{Model: model, Detail: detail, Failed: failed, Fail: fail, Generate: usage.GenerateFlag(true)} } + // Additional-model records describe a side model (image generation tool usage) that + // the upstream response model never refers to, so they must stay empty. + responseModel := "" + if model == r.model { + responseModel = r.ResponseModel() + } return usage.Record{ Provider: r.provider, BaseURL: r.baseURL, @@ -461,6 +545,7 @@ func (r *UsageReporter) buildRecordForModel(model string, detail usage.Detail, f ReasoningEffort: r.reasoning, ServiceTier: r.serviceTier, ResponseServiceTier: strings.TrimSpace(detail.ResponseServiceTier), + ResponseModel: responseModel, Generate: usage.GenerateFlag(r.generate), Stream: r.stream, RequestedAt: r.requestedAt, diff --git a/sdk/cliproxy/usage/manager.go b/sdk/cliproxy/usage/manager.go index 5fd800999..5bf39920a 100644 --- a/sdk/cliproxy/usage/manager.go +++ b/sdk/cliproxy/usage/manager.go @@ -45,6 +45,8 @@ type Record struct { RequestServiceTier string // ResponseServiceTier stores the final tier reported by the upstream response. ResponseServiceTier string + // ResponseModel stores the model name reported by the upstream response, empty when unknown. + ResponseModel string // Generate reports whether the client requested actual generation. // nil or true means generation is enabled; only an explicit false disables generation. // Use GenerateFlag to set the value and GenerateEnabled to read it with the default.