From 25f40d8cf8dfa8765e45873060cf41056cd1b112 Mon Sep 17 00:00:00 2001 From: Viggo95 <15306225869@163.com> Date: Fri, 18 Sep 2026 09:52:00 +0800 Subject: [PATCH] feat(codex): record upstream response model and warn on silent model substitution Codex upstreams can silently serve a different model than the one requested (HTTP 200, with response.model naming the substitute). The proxy kept no record of it: nothing logged, nothing reported, only the pass-through response body. - Add Record.ResponseModel to sdk/cliproxy/usage, aligned with the existing ResponseServiceTier field, and emit it from the redis usage queue as the optional response_model payload field alongside response_service_tier. Only the record for the requested model carries it: additional-model records (image generation tool usage) describe a side model the upstream response never refers to, and would otherwise look like a substitution downstream. - Add internal/runtime/executor/helps/response_model.go with extractCodexResponseModelEvent (SSE frames and raw JSON, restricted to the events that embed the authoritative response object, rejecting non-string and oversized upstream model names) and IsCodexModelSubstituted (both sides trimmed, lower-cased and stripped of thinking suffixes, dated aliases such as gpt-5.6-terra-2026-05-13 accepted in either direction). - UsageReporter records the served model on the event path and emits the WARN when the attempt publishes its usage record, so no logging work happens before the first event is forwarded. Repeats are throttled per (auth id, requested model, served model) with a 10 minute window, because on an affected credential every request is substituted and an unthrottled warning would mirror the whole request volume into the logs. The credential is labelled auth_index= only: codex credential file names embed the account e-mail, which must not be written to the logs at request rate. Coverage, by entry point. The served model is observed on the HTTP streaming path (both the bootstrap-buffered handshake and the streaming goroutine), the HTTP non-streaming Execute loop, the websocket streaming and non-streaming paths, and the two /responses-shaped image entry points. The remaining codex entry points cannot report it and are therefore left alone: executeCompact (/responses/compact answers with a compaction object that has no event type and no response.model), the two direct image endpoints (/images/generations and /images/edits answer in the Images API shape and stream image_generation.* events), and CountTokens (counts locally with tiktoken, never reaching an upstream). TokenAccountingSchemaVersion is not bumped: it versions the token accounting contract (token breakdown semantics), and this change only adds an optional non-token field that leaves existing consumers and all token math untouched. Note: response_model ships with the usage record and is the counting source; the WARN is a throttled alerting signal and must not be used to count substitutions. Tests: table-driven unit tests for both helpers over real model ids, reporter tests covering the published record, the single throttled warning, the absence of account identifiers in it, concurrent observation and publishing under -race, the throttle window and its entry bound, an executor-level guard for the observeCodexTokenEvent wiring and the per-model records, plus a redisqueue payload assertion for response_model. gofmt, go vet, go test -race on the touched packages and go test ./... are clean. --- internal/redisqueue/plugin.go | 3 + internal/redisqueue/plugin_test.go | 2 + .../executor/codex_executor_execute.go | 1 + .../executor/codex_executor_terminal.go | 4 +- .../runtime/executor/codex_openai_images.go | 2 + .../executor/codex_response_model_test.go | 107 ++++ .../runtime/executor/helps/response_model.go | 181 +++++++ .../executor/helps/response_model_test.go | 479 ++++++++++++++++++ .../runtime/executor/helps/usage_helpers.go | 89 +++- sdk/cliproxy/usage/manager.go | 2 + 10 files changed, 867 insertions(+), 3 deletions(-) create mode 100644 internal/runtime/executor/codex_response_model_test.go create mode 100644 internal/runtime/executor/helps/response_model.go create mode 100644 internal/runtime/executor/helps/response_model_test.go diff --git a/internal/redisqueue/plugin.go b/internal/redisqueue/plugin.go index 0b3bf108b8..4eb8280b6a 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 248f2a505f..78c2285652 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 49ef1a127f..ab54892475 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 a83f01d224..b91b276945 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 59ef7a2618..d20b00a91f 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 0000000000..97719ddbdf --- /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 0000000000..1c2eb1d57e --- /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 0000000000..5d456717d4 --- /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 c2a43db29d..2bd22fb309 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 5fd800999e..5bf39920a7 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.