From f702bc1ac2631953b3eea2d99eb6761b283a99ef Mon Sep 17 00:00:00 2001 From: Luis Pater Date: Sun, 13 Sep 2026 15:05:03 +0800 Subject: [PATCH] feat(codex): preserve native fidelity for responses-lite requests - Detect native responses-lite requests via headers and client metadata. - Skip instructions normalization and synthetic session cloaking for native requests. - Preserve upstream completion output during websocket response forwarding. Closes: #5780 --- .../executor/codex_executor_execute.go | 4 +- .../executor/codex_executor_request.go | 22 +-- .../runtime/executor/codex_executor_stream.go | 11 +- .../runtime/executor/codex_executor_tokens.go | 2 +- .../executor/codex_native_fidelity_test.go | 156 ++++++++++++++++++ .../executor/codex_websockets_connection.go | 3 +- .../executor/codex_websockets_execute.go | 5 +- .../codex_websockets_executor_test.go | 129 +++++++++++++-- .../executor/codex_websockets_request.go | 16 +- .../executor/codex_websockets_stream.go | 13 +- .../runtime/executor/helps/codex_native.go | 20 +++ internal/util/codex.go | 17 ++ .../openai/openai_responses_websocket.go | 10 +- .../openai_responses_websocket_forward.go | 9 +- .../openai/openai_responses_websocket_test.go | 45 ++++- 15 files changed, 403 insertions(+), 59 deletions(-) create mode 100644 internal/runtime/executor/codex_native_fidelity_test.go create mode 100644 internal/runtime/executor/helps/codex_native.go create mode 100644 internal/util/codex.go diff --git a/internal/runtime/executor/codex_executor_execute.go b/internal/runtime/executor/codex_executor_execute.go index c68757d63..49ef1a127 100644 --- a/internal/runtime/executor/codex_executor_execute.go +++ b/internal/runtime/executor/codex_executor_execute.go @@ -61,7 +61,7 @@ func (e *CodexExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, re body, _ = sjson.DeleteBytes(body, "prompt_cache_retention") body, _ = sjson.DeleteBytes(body, "safety_identifier") body, _ = sjson.DeleteBytes(body, "stream_options") - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, helps.IsNativeCodexRequest(req.Payload, opts)) if e.cfg == nil || e.cfg.DisableImageGeneration == config.DisableImageGenerationOff { body = ensureImageGenerationTool(body, baseModel, auth, opts.Headers) } @@ -239,7 +239,7 @@ func (e *CodexExecutor) executeCompact(ctx context.Context, auth *cliproxyauth.A body = helps.ApplyPayloadConfigWithRequest(e.cfg, baseModel, to.String(), from.String(), "", body, originalTranslated, requestedModel, requestPath, opts.Headers) body = helps.SetStringIfDifferent(body, "model", baseModel) body, _ = sjson.DeleteBytes(body, "stream") - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, helps.IsNativeCodexRequest(req.Payload, opts)) body = sanitizeOpenAIResponsesReasoningEncryptedContent(ctx, "codex executor", body) body = normalizeCodexParallelToolCalls(body, opts.Headers) body = helps.NormalizeCodexToolSchemas(body) diff --git a/internal/runtime/executor/codex_executor_request.go b/internal/runtime/executor/codex_executor_request.go index ceb2fc597..b67eda88c 100644 --- a/internal/runtime/executor/codex_executor_request.go +++ b/internal/runtime/executor/codex_executor_request.go @@ -27,7 +27,6 @@ const ( codexOriginator = "codex-tui" codexDefaultImageToolModel = "gpt-image-2" codexResponsesLiteHeader = "X-OpenAI-Internal-Codex-Responses-Lite" - codexResponsesLiteMetadata = "client_metadata.ws_request_header_x_openai_internal_codex_responses_lite" ) var dataTag = []byte("data:") @@ -378,7 +377,10 @@ func applyCodexCloakingHeaders(headers http.Header, cfg *config.Config) { headers.Set("Originator", codexOriginator) } -func normalizeCodexInstructions(body []byte) []byte { +func normalizeCodexInstructions(body []byte, nativeRequest ...bool) []byte { + if len(nativeRequest) > 0 && nativeRequest[0] { + return body + } instructions := gjson.GetBytes(body, "instructions") if !instructions.Exists() || instructions.Type == gjson.Null { body, _ = sjson.SetBytes(body, "instructions", "") @@ -420,20 +422,8 @@ func isImageGenerationFunctionTool(tool gjson.Result) bool { return false } -func isCodexResponsesLiteRequest(body []byte, headers http.Header) bool { - if strings.EqualFold(strings.TrimSpace(headers.Get(codexResponsesLiteHeader)), "true") { - return true - } - // Codex Desktop mirrors websocket-only request headers into client_metadata. - value := gjson.GetBytes(body, codexResponsesLiteMetadata) - if !value.Exists() { - return false - } - return value.Type == gjson.True || value.Type == gjson.String && strings.EqualFold(strings.TrimSpace(value.String()), "true") -} - func ensureImageGenerationTool(body []byte, baseModel string, auth *cliproxyauth.Auth, headers http.Header) []byte { - if isCodexResponsesLiteRequest(body, headers) { + if util.IsCodexResponsesLiteRequest(body, headers) { return body } if strings.HasSuffix(baseModel, "spark") { @@ -458,7 +448,7 @@ func ensureImageGenerationTool(body []byte, baseModel string, auth *cliproxyauth } func normalizeCodexParallelToolCalls(body []byte, headers http.Header) []byte { - if isCodexResponsesLiteRequest(body, headers) { + if util.IsCodexResponsesLiteRequest(body, headers) { body = helps.SetBoolIfDifferent(body, "parallel_tool_calls", false) return body } diff --git a/internal/runtime/executor/codex_executor_stream.go b/internal/runtime/executor/codex_executor_stream.go index 1ea663e0b..54a6c4416 100644 --- a/internal/runtime/executor/codex_executor_stream.go +++ b/internal/runtime/executor/codex_executor_stream.go @@ -40,6 +40,7 @@ func (e *CodexExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Au from := opts.SourceFormat responseFormat := cliproxyexecutor.ResponseFormatOrSource(opts) + preserveNativeOutput := helps.IsNativeCodexRequest(req.Payload, opts) isGrokClient := grokbuild.IsGrokClientContext(ctx, opts.Headers) to := sdktranslator.FromString("codex") originalPayloadSource := req.Payload @@ -67,7 +68,7 @@ func (e *CodexExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Au body, _ = sjson.SetBytes(body, "stream_options.reasoning_summary_delivery", reasoningSummaryDelivery.Value()) } body = helps.SetStringIfDifferent(body, "model", baseModel) - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, preserveNativeOutput) if e.cfg == nil || e.cfg.DisableImageGeneration == config.DisableImageGenerationOff { body = ensureImageGenerationTool(body, baseModel, auth, opts.Headers) } @@ -222,7 +223,9 @@ func (e *CodexExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Au reporter.EnsurePublished(ctx) } publishCodexImageToolUsage(ctx, reporter, body, data) - data = patchCodexCompletedOutput(data, outputItemsByIndex, outputItemsFallback) + if !preserveNativeOutput { + data = patchCodexCompletedOutput(data, outputItemsByIndex, outputItemsFallback) + } if eventType == "response.completed" || eventType == "response.done" { cacheCodexReasoningReplayFromCompleted(replayScope, data) } @@ -360,7 +363,9 @@ func (e *CodexExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Au reporter.EnsurePublished(ctx) } publishCodexImageToolUsage(ctx, reporter, body, data) - data = patchCodexCompletedOutput(data, outputItemsByIndex, outputItemsFallback) + if !preserveNativeOutput { + data = patchCodexCompletedOutput(data, outputItemsByIndex, outputItemsFallback) + } if eventType == "response.completed" || eventType == "response.done" { cacheCodexReasoningReplayFromCompleted(replayScope, data) } diff --git a/internal/runtime/executor/codex_executor_tokens.go b/internal/runtime/executor/codex_executor_tokens.go index a72dcb39a..43f165ad5 100644 --- a/internal/runtime/executor/codex_executor_tokens.go +++ b/internal/runtime/executor/codex_executor_tokens.go @@ -35,7 +35,7 @@ func (e *CodexExecutor) CountTokens(ctx context.Context, auth *cliproxyauth.Auth body, _ = sjson.DeleteBytes(body, "safety_identifier") body, _ = sjson.DeleteBytes(body, "stream_options") body = helps.SetBoolIfDifferent(body, "stream", false) - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, helps.IsNativeCodexRequest(req.Payload, opts)) enc, err := tokenizerForCodexModel(baseModel) if err != nil { diff --git a/internal/runtime/executor/codex_native_fidelity_test.go b/internal/runtime/executor/codex_native_fidelity_test.go new file mode 100644 index 000000000..7e802fb95 --- /dev/null +++ b/internal/runtime/executor/codex_native_fidelity_test.go @@ -0,0 +1,156 @@ +package executor + +import ( + "bytes" + "context" + "fmt" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/gorilla/websocket" + "github.com/router-for-me/CLIProxyAPI/v7/internal/config" + cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" + sdktranslator "github.com/router-for-me/CLIProxyAPI/v7/sdk/translator" + "github.com/tidwall/gjson" +) + +func TestCodexNativeStreamFidelity(t *testing.T) { + for _, source := range []sdktranslator.Format{sdktranslator.FormatCodex, sdktranslator.FormatOpenAIResponse, sdktranslator.FormatClaude, sdktranslator.FormatOpenAI} { + t.Run(source.String(), func(t *testing.T) { testCodexNativeStreamFidelity(t, source) }) + } +} + +func testCodexNativeStreamFidelity(t *testing.T, source sdktranslator.Format) { + t.Helper() + for _, transport := range []string{"http", "websocket"} { + for _, lite := range []string{"", "header", "metadata"} { + for _, buffering := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/lite=%s/buffering=%t", transport, lite, buffering), func(t *testing.T) { + metadata := `{"type":"codex.response.metadata","headers":{"x-models-etag":"models-v1","x-codex-turn-state":"turn-1","x-codex-safety-buffering-enabled":"true","x-codex-safety-buffering-faster-model":"fixture-model"},"future":{"ok":true}}` + completed := `{"type":"response.completed","response":{"id":"resp_1","status":"completed","output":[],"future":{"ok":true},"usage":{"input_tokens":1,"output_tokens":1,"total_tokens":2}}}` + events := []string{metadata, `{"type":"response.output_item.done","output_index":0,"item":{"id":"msg_1","type":"message","role":"assistant","content":[{"type":"output_text","text":"ok"}]}}`, completed} + captured := make(chan []byte, 1) + capturedHeaders := make(chan http.Header, 1) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + capturedHeaders <- r.Header.Clone() + if transport == "websocket" { + upgrader := websocket.Upgrader{} + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + t.Error(err) + return + } + defer func() { _ = conn.Close() }() + _, body, errRead := conn.ReadMessage() + if errRead != nil { + t.Error(errRead) + return + } + captured <- body + for _, event := range events { + if errWrite := conn.WriteMessage(websocket.TextMessage, []byte(event)); errWrite != nil { + t.Error(errWrite) + } + } + return + } + body, errRead := io.ReadAll(r.Body) + if errRead != nil { + t.Error(errRead) + } + captured <- body + w.Header().Set("Content-Type", "text/event-stream") + for _, event := range events { + _, _ = fmt.Fprintf(w, "data: %s\n\n", event) + } + })) + defer server.Close() + payload := []byte(`{"model":"gpt-5.6-sol","input":[],"parallel_tool_calls":false}`) + headers := http.Header{"Session-Id": {"session-1"}, "Thread-Id": {"thread-1"}} + if lite == "header" { + headers.Set(codexResponsesLiteHeader, "true") + } else if lite == "metadata" { + payload = []byte(`{"model":"gpt-5.6-sol","input":[],"parallel_tool_calls":false,"client_metadata":{"ws_request_header_x_openai_internal_codex_responses_lite":"true"}}`) + } + cfg := &config.Config{Codex: config.CodexConfig{StreamBootstrapBuffering: buffering, DisableCodexCloaking: true}} + executeStream := NewCodexExecutor(cfg).ExecuteStream + if transport == "websocket" { + executeStream = NewCodexWebsocketsExecutor(cfg).ExecuteStream + } + result, err := executeStream(context.Background(), codexTestAuth(server.URL), cliproxyexecutor.Request{Model: "gpt-5.6-sol", Payload: payload}, cliproxyexecutor.Options{SourceFormat: source, ResponseFormat: sdktranslator.FormatCodex, Headers: headers, Stream: true}) + if err != nil { + t.Fatal(err) + } + var terminal []byte + var metadataEvents []string + for chunk := range result.Chunks { + if chunk.Err != nil { + t.Fatal(chunk.Err) + } + for _, line := range bytes.Split(chunk.Payload, []byte("\n")) { + data := bytes.TrimSpace(bytes.TrimPrefix(line, []byte("data:"))) + if gjson.GetBytes(data, "type").String() == "codex.response.metadata" { + metadataEvents = append(metadataEvents, string(data)) + } + if gjson.GetBytes(data, "type").String() == "response.completed" { + terminal = bytes.Clone(data) + } + } + } + body := <-captured + upstreamHeaders := <-capturedHeaders + native := lite != "" && (source == sdktranslator.FormatCodex || source == sdktranslator.FormatOpenAIResponse) + if transport == "websocket" { + wantLiteHeader := "" + if native && lite == "header" { + wantLiteHeader = "true" + } + if got := upstreamHeaders.Get(codexResponsesLiteHeader); got != wantLiteHeader { + t.Errorf("upstream Lite header = %q, want %q", got, wantLiteHeader) + } + alias := headerValueCaseInsensitive(upstreamHeaders, "session_id") + t.Logf("upstream session alias: %q", alias) + if (alias == "") != native { + t.Errorf("session alias = %q, native = %t", alias, native) + } + } + t.Logf("downstream metadata: %q", metadataEvents) + if len(metadataEvents) != 1 || metadataEvents[0] != metadata { + t.Errorf("metadata event changed or duplicated: %q", metadataEvents) + } + t.Logf("upstream request: %s; downstream completion: %s", body, terminal) + if native { + if gjson.GetBytes(body, "instructions").Exists() { + t.Errorf("native request gained instructions: %s", body) + } + if string(terminal) != completed { + t.Errorf("native completion changed: %s", terminal) + } + } else if gjson.GetBytes(body, "instructions").Type != gjson.String || gjson.GetBytes(terminal, "response.output.0.id").String() != "msg_1" { + t.Errorf("compatibility normalization/backfill lost: %s; %s", body, terminal) + } + }) + } + } + } +} + +func TestCodexWebsocketLiteHeaderWithoutSessionHeaders(t *testing.T) { + for _, native := range []bool{false, true} { + for _, disableCloaking := range []bool{false, true} { + cfg := &config.Config{Codex: config.CodexConfig{DisableCodexCloaking: disableCloaking}} + headers := http.Header{} + headers.Set(codexResponsesLiteHeader, "true") + got := applyCodexWebsocketHeaders(context.Background(), nil, nil, "fixture-token", cfg, native, headers) + want := "" + if native { + want = "true" + } + if value := got.Get(codexResponsesLiteHeader); value != want { + t.Errorf("native=%t disableCloaking=%t: Lite header = %q, want %q", native, disableCloaking, value, want) + } + } + } +} diff --git a/internal/runtime/executor/codex_websockets_connection.go b/internal/runtime/executor/codex_websockets_connection.go index 7f4856c92..54b3695c3 100644 --- a/internal/runtime/executor/codex_websockets_connection.go +++ b/internal/runtime/executor/codex_websockets_connection.go @@ -13,6 +13,7 @@ import ( "github.com/gorilla/websocket" "github.com/router-for-me/CLIProxyAPI/v7/internal/config" "github.com/router-for-me/CLIProxyAPI/v7/internal/runtime/executor/helps" + "github.com/router-for-me/CLIProxyAPI/v7/internal/util" cliproxyauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" "github.com/router-for-me/CLIProxyAPI/v7/sdk/proxyutil" @@ -123,7 +124,7 @@ func mapCodexWebsocketReadError(err error) error { } func normalizeCodexWebsocketParallelToolCalls(body []byte, headers http.Header) []byte { - if !isCodexResponsesLiteRequest(body, headers) { + if !util.IsCodexResponsesLiteRequest(body, headers) { return body } body = helps.SetBoolIfDifferent(body, "parallel_tool_calls", false) diff --git a/internal/runtime/executor/codex_websockets_execute.go b/internal/runtime/executor/codex_websockets_execute.go index e2c8026db..70a214f24 100644 --- a/internal/runtime/executor/codex_websockets_execute.go +++ b/internal/runtime/executor/codex_websockets_execute.go @@ -37,6 +37,7 @@ func (e *CodexWebsocketsExecutor) Execute(ctx context.Context, auth *cliproxyaut from := opts.SourceFormat responseFormat := cliproxyexecutor.ResponseFormatOrSource(opts) + nativeRequest := helps.IsNativeCodexRequest(req.Payload, opts) to := sdktranslator.FromString("codex") originalPayloadSource := req.Payload if len(opts.OriginalRequest) > 0 { @@ -57,7 +58,7 @@ func (e *CodexWebsocketsExecutor) Execute(ctx context.Context, auth *cliproxyaut body = helps.SetBoolIfDifferent(body, "stream", true) body, _ = sjson.DeleteBytes(body, "prompt_cache_retention") body, _ = sjson.DeleteBytes(body, "safety_identifier") - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, nativeRequest) if e.cfg == nil || e.cfg.DisableImageGeneration == config.DisableImageGenerationOff { body = ensureImageGenerationTool(body, baseModel, auth, opts.Headers) } @@ -85,7 +86,7 @@ func (e *CodexWebsocketsExecutor) Execute(ctx context.Context, auth *cliproxyaut var identityState codexIdentityConfuseState upstreamBody, identityState := applyCodexIdentityConfuseBody(e.cfg, auth, originalPayloadSource, body) reporter.SetTranslatedReasoningEffort(clientBody, to.String()) - wsHeaders = applyCodexWebsocketHeaders(ctx, wsHeaders, auth, apiKey, e.cfg, opts.Headers) + wsHeaders = applyCodexWebsocketHeaders(ctx, wsHeaders, auth, apiKey, e.cfg, nativeRequest, opts.Headers) applyModelHeaderOverrides(wsHeaders, baseModel) applyCodexIdentityConfuseHeaders(wsHeaders, &identityState) diff --git a/internal/runtime/executor/codex_websockets_executor_test.go b/internal/runtime/executor/codex_websockets_executor_test.go index cf12cfc90..2dc594d59 100644 --- a/internal/runtime/executor/codex_websockets_executor_test.go +++ b/internal/runtime/executor/codex_websockets_executor_test.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "errors" + "fmt" "io" "net/http" "net/http/httptest" @@ -233,6 +234,9 @@ func TestCodexWebsocketsExecuteResponsesLiteDoesNotInjectImageGenerationTool(t * select { case payload := <-capturedPayload: + if instructions := gjson.GetBytes(payload, "instructions"); instructions.Exists() { + t.Errorf("unexpected instructions in responses-lite upstream payload: %s", payload) + } if tools := gjson.GetBytes(payload, "tools"); tools.Exists() { t.Fatalf("unexpected tools in responses-lite upstream payload: %s", tools.Raw) } @@ -1080,7 +1084,7 @@ func TestCodexWebsocketsUpstreamDisconnectChanSignalsOnInvalidate(t *testing.T) } func TestApplyCodexWebsocketHeadersDefaultsToCurrentResponsesBeta(t *testing.T) { - headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, nil, "", nil) + headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, nil, "", nil, false) if got := headers.Get("OpenAI-Beta"); got != codexResponsesWebsocketBetaHeaderValue { t.Fatalf("OpenAI-Beta = %s, want %s", got, codexResponsesWebsocketBetaHeaderValue) @@ -1157,7 +1161,7 @@ func TestApplyCodexWebsocketHeadersDefaultsToCodexCloaking(t *testing.T) { headers.Set("User-Agent", "existing-ua") headers.Set("Originator", "existing-origin") - headers = applyCodexWebsocketHeaders(ctx, headers, tt.auth, tt.token, cfg) + headers = applyCodexWebsocketHeaders(ctx, headers, tt.auth, tt.token, cfg, false) if got := headers.Get("User-Agent"); got != codexUserAgent { t.Fatalf("User-Agent = %q, want %q", got, codexUserAgent) @@ -1182,9 +1186,12 @@ func TestApplyCodexWebsocketHeadersPassesThroughClientIdentityHeadersWhenCloakin "X-Codex-Turn-Metadata": `{"turn_id":"turn-1"}`, "X-Client-Request-Id": "019d2233-e240-7162-992d-38df0a2a0e0d", "session-id": "legacy-session", + "Thread-Id": "thread-1", + "X-Codex-Routing-Hint": "route-1", + "X-Codex-Window-Id": "window-1", }) - headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "", cfg) + headers := applyCodexWebsocketHeaders(ctx, http.Header{"session_id": {"cache-key"}, "Conversation_id": {"cache-key"}}, auth, "", cfg, true) if got := headers.Get("Originator"); got != "Codex Desktop" { t.Fatalf("Originator = %s, want %s", got, "Codex Desktop") @@ -1201,11 +1208,101 @@ func TestApplyCodexWebsocketHeadersPassesThroughClientIdentityHeadersWhenCloakin if got := headers.Get("X-Client-Request-Id"); got != "019d2233-e240-7162-992d-38df0a2a0e0d" { t.Fatalf("X-Client-Request-Id = %s, want %s", got, "019d2233-e240-7162-992d-38df0a2a0e0d") } - if got := headers["session_id"]; len(got) != 1 || got[0] != "legacy-session" { - t.Fatalf("session_id = %#v, want [legacy-session]", got) + if got := headerValueCaseInsensitive(headers, "session_id"); got != "" { + t.Fatalf("unexpected session_id = %q", got) } - if got := headers.Get("Session-Id"); got != "" { - t.Fatalf("Session-Id = %s, want empty", got) + if got := headerValueCaseInsensitive(headers, "conversation_id"); got != "" { + t.Fatalf("unexpected conversation_id = %q", got) + } + for key, want := range map[string]string{"Session-Id": "legacy-session", "Thread-Id": "thread-1", "X-Codex-Routing-Hint": "route-1", "X-Codex-Window-Id": "window-1"} { + if got := headers.Get(key); got != want { + t.Errorf("%s = %q, want %q", key, got, want) + } + } +} + +func TestApplyCodexWebsocketHeadersNativeSessionCombinations(t *testing.T) { + cfg := &config.Config{ + Codex: config.CodexConfig{DisableCodexCloaking: true}, + } + auth := &cliproxyauth.Auth{ + Provider: "codex", + Metadata: map[string]any{"email": "user@example.com"}, + } + + tests := []struct { + name string + clientHeaders map[string]string + wantSessionID string + wantThreadID string + }{ + { + name: "both session and thread present", + clientHeaders: map[string]string{ + "Session-Id": "sess-both", + "Thread-Id": "thread-both", + }, + wantSessionID: "sess-both", + wantThreadID: "thread-both", + }, + { + name: "only session present", + clientHeaders: map[string]string{ + "Session-Id": "sess-only", + }, + wantSessionID: "sess-only", + wantThreadID: "", + }, + { + name: "only thread present", + clientHeaders: map[string]string{ + "Thread-Id": "thread-only", + }, + wantSessionID: "", + wantThreadID: "thread-only", + }, + { + name: "neither present", + clientHeaders: map[string]string{}, + wantSessionID: "", + wantThreadID: "", + }, + } + + for _, tt := range tests { + for _, withCacheAliases := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/cache_aliases=%t", tt.name, withCacheAliases), func(t *testing.T) { + ctx := contextWithGinHeaders(tt.clientHeaders) + initialHeaders := http.Header{} + if withCacheAliases { + initialHeaders = http.Header{"session_id": {"cache-alias"}, "Conversation_id": {"cache-alias"}} + } + got := applyCodexWebsocketHeaders(ctx, initialHeaders, auth, "", cfg, true) + + if tt.wantSessionID != "" { + if val := got.Get("Session-Id"); val != tt.wantSessionID { + t.Errorf("Session-Id = %q, want %q", val, tt.wantSessionID) + } + } else if val := got.Get("Session-Id"); val != "" { + t.Errorf("unexpected Session-Id = %q", val) + } + + if tt.wantThreadID != "" { + if val := got.Get("Thread-Id"); val != tt.wantThreadID { + t.Errorf("Thread-Id = %q, want %q", val, tt.wantThreadID) + } + } else if val := got.Get("Thread-Id"); val != "" { + t.Errorf("unexpected Thread-Id = %q", val) + } + + if hasSessionAlias := headerValueCaseInsensitive(got, "session_id"); hasSessionAlias != "" { + t.Errorf("unexpected synthesized session_id alias = %q", hasSessionAlias) + } + if hasConversationAlias := headerValueCaseInsensitive(got, "conversation_id"); hasConversationAlias != "" { + t.Errorf("unexpected synthesized conversation_id alias = %q", hasConversationAlias) + } + }) + } } } @@ -1220,7 +1317,7 @@ func TestApplyCodexWebsocketHeadersCanonicalizesLegacyUnderscoreSessionHeader(t "Session_id": "legacy-underscore-session", }) - headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "", nil) + headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "", nil, false) if got := headers["session_id"]; len(got) != 1 || got[0] != "legacy-underscore-session" { t.Fatalf("session_id = %#v, want [legacy-underscore-session]", got) @@ -1243,7 +1340,7 @@ func TestApplyCodexWebsocketHeadersUsesConfigDefaultsForOAuth(t *testing.T) { Metadata: map[string]any{"email": "user@example.com"}, } - headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "", cfg) + headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "", cfg, false) if got := headers.Get("User-Agent"); got != "my-codex-client/1.0" { t.Fatalf("User-Agent = %s, want %s", got, "my-codex-client/1.0") @@ -1276,7 +1373,7 @@ func TestApplyCodexWebsocketHeadersPrefersExistingHeadersOverClientAndConfig(t * headers.Set("User-Agent", "existing-ua") headers.Set("X-Codex-Beta-Features", "existing-beta") - got := applyCodexWebsocketHeaders(ctx, headers, auth, "", cfg) + got := applyCodexWebsocketHeaders(ctx, headers, auth, "", cfg, false) if gotVal := got.Get("User-Agent"); gotVal != "existing-ua" { t.Fatalf("User-Agent = %s, want %s", gotVal, "existing-ua") @@ -1303,7 +1400,7 @@ func TestApplyCodexWebsocketHeadersConfigUserAgentOverridesClientHeader(t *testi "X-Codex-Beta-Features": "client-beta", }) - headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "", cfg) + headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "", cfg, false) if got := headers.Get("User-Agent"); got != "config-ua" { t.Fatalf("User-Agent = %s, want %s", got, "config-ua") @@ -1326,7 +1423,7 @@ func TestApplyCodexWebsocketHeadersIgnoresConfigForAPIKeyAuth(t *testing.T) { Attributes: map[string]string{"api_key": "sk-test"}, } - headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "sk-test", cfg) + headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "sk-test", cfg, false) if got := headers.Get("User-Agent"); got != "" { t.Fatalf("User-Agent = %s, want empty", got) @@ -1343,7 +1440,7 @@ func TestApplyCodexWebsocketHeadersPreservesExplicitAPIKeyUserAgent(t *testing.T auth := &cliproxyauth.Auth{Provider: "codex", Attributes: map[string]string{"api_key": "sk-test"}} ctx := contextWithGinHeaders(map[string]string{"User-Agent": "api-key-client/1.0", "Originator": "explicit-origin"}) - headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "sk-test", nil) + headers := applyCodexWebsocketHeaders(ctx, http.Header{}, auth, "sk-test", nil, false) if got := headers.Get("User-Agent"); got != "api-key-client/1.0" { t.Fatalf("User-Agent = %s, want api-key-client/1.0", got) @@ -1356,7 +1453,7 @@ func TestApplyCodexWebsocketHeadersPreservesExplicitAPIKeyUserAgent(t *testing.T func TestApplyCodexWebsocketHeadersUsesCanonicalAccountHeader(t *testing.T) { auth := &cliproxyauth.Auth{Provider: "codex", Metadata: map[string]any{"account_id": "acct-1"}} - headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "", nil) + headers := applyCodexWebsocketHeaders(context.Background(), http.Header{}, auth, "", nil, false) if got := headerValueCaseInsensitive(headers, "ChatGPT-Account-ID"); got != "acct-1" { t.Fatalf("ChatGPT-Account-ID = %s, want acct-1", got) @@ -1509,7 +1606,7 @@ func TestApplyCodexWebsocketHeadersIdentityConfuseRemapsPromptCacheKey(t *testin "X-Codex-Turn-Metadata": `{"prompt_cache_key":"cache-ws-1","turn_id":"turn-ws-1","window_id":"cache-ws-1:0"}`, "X-Client-Request-Id": "client-request-1", }) - headers = applyCodexWebsocketHeaders(ctx, headers, auth, "oauth-token", cfg) + headers = applyCodexWebsocketHeaders(ctx, headers, auth, "oauth-token", cfg, false) applyCodexIdentityConfuseHeaders(headers, &identityState) expectedPromptCacheKey := codexIdentityConfuseUUID("auth-ws-1", "prompt-cache", "cache-ws-1") @@ -1805,7 +1902,7 @@ func TestApplyCodexWebsocketHeaders_EmptyAPIKey_OmitsAuthorizationAndOAuthHeader BetaFeatures: "oauth-beta", }, } - headers := applyCodexWebsocketHeaders(context.Background(), nil, auth, "", cfg) + headers := applyCodexWebsocketHeaders(context.Background(), nil, auth, "", cfg, false) if got := headers.Get("Authorization"); got != "" { t.Fatalf("Authorization = %q, want empty for empty API key", got) } diff --git a/internal/runtime/executor/codex_websockets_request.go b/internal/runtime/executor/codex_websockets_request.go index f9c027b7c..2c75bff62 100644 --- a/internal/runtime/executor/codex_websockets_request.go +++ b/internal/runtime/executor/codex_websockets_request.go @@ -64,7 +64,7 @@ func applyCodexPromptCacheHeadersWithContext(ctx context.Context, from sdktransl return rawJSON, headers, nil } -func applyCodexWebsocketHeaders(ctx context.Context, headers http.Header, auth *cliproxyauth.Auth, token string, cfg *config.Config, clientHeaders ...http.Header) http.Header { +func applyCodexWebsocketHeaders(ctx context.Context, headers http.Header, auth *cliproxyauth.Auth, token string, cfg *config.Config, nativeRequest bool, clientHeaders ...http.Header) http.Header { if headers == nil { headers = http.Header{} } @@ -89,6 +89,9 @@ func applyCodexWebsocketHeaders(ctx context.Context, headers http.Header, auth * misc.EnsureHeader(headers, ginHeaders, "x-client-request-id", "") misc.EnsureHeader(headers, ginHeaders, "x-responsesapi-include-timing-metrics", "") misc.EnsureHeader(headers, ginHeaders, "Version", "") + if nativeRequest { + misc.EnsureHeader(headers, ginHeaders, codexResponsesLiteHeader, "") + } if isAPIKey { ensureHeaderWithPriority(headers, ginHeaders, "User-Agent", "", "") } else { @@ -108,6 +111,17 @@ func applyCodexWebsocketHeaders(ctx context.Context, headers http.Header, auth * sessionFallback = uuid.NewString() } ensureCodexWebsocketSessionHeader(headers, ginHeaders, sessionFallback) + if nativeRequest && cfg != nil && cfg.Codex.DisableCodexCloaking { + deleteHeaderCaseInsensitive(headers, "session_id") + deleteHeaderCaseInsensitive(headers, "conversation_id") + for key, values := range ginHeaders { + switch strings.ToLower(key) { + case "session-id", "session_id", "conversation_id", "thread-id", "x-codex-routing-hint", "x-codex-window-id": + deleteHeaderCaseInsensitive(headers, key) + headers[key] = append([]string(nil), values...) + } + } + } if originator := strings.TrimSpace(ginHeaders.Get("Originator")); originator != "" { headers.Set("Originator", originator) } else if !isAPIKey { diff --git a/internal/runtime/executor/codex_websockets_stream.go b/internal/runtime/executor/codex_websockets_stream.go index 0247c4f3f..19232a4e8 100644 --- a/internal/runtime/executor/codex_websockets_stream.go +++ b/internal/runtime/executor/codex_websockets_stream.go @@ -38,6 +38,7 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr from := opts.SourceFormat responseFormat := cliproxyexecutor.ResponseFormatOrSource(opts) + preserveNativeOutput := helps.IsNativeCodexRequest(req.Payload, opts) to := sdktranslator.FromString("codex") originalPayloadSource := req.Payload if len(opts.OriginalRequest) > 0 { @@ -55,7 +56,7 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr requestPath := helps.PayloadRequestPath(opts) body = helps.ApplyPayloadConfigWithRequest(e.cfg, baseModel, to.String(), from.String(), "", body, originalTranslated, requestedModel, requestPath, opts.Headers) body = helps.SetStringIfDifferent(body, "model", baseModel) - body = normalizeCodexInstructions(body) + body = normalizeCodexInstructions(body, preserveNativeOutput) if e.cfg == nil || e.cfg.DisableImageGeneration == config.DisableImageGenerationOff { body = ensureImageGenerationTool(body, baseModel, auth, opts.Headers) } @@ -83,7 +84,7 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr var identityState codexIdentityConfuseState upstreamBody, identityState := applyCodexIdentityConfuseBody(e.cfg, auth, originalPayloadSource, body) reporter.SetTranslatedReasoningEffort(clientBody, to.String()) - wsHeaders = applyCodexWebsocketHeaders(ctx, wsHeaders, auth, apiKey, e.cfg, opts.Headers) + wsHeaders = applyCodexWebsocketHeaders(ctx, wsHeaders, auth, apiKey, e.cfg, preserveNativeOutput, opts.Headers) applyModelHeaderOverrides(wsHeaders, baseModel) applyCodexIdentityConfuseHeaders(wsHeaders, &identityState) @@ -437,7 +438,9 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr completedPayload := payload if eventType == "response.completed" || eventType == "response.done" || eventType == "response.incomplete" { completedPayload = normalizeCodexWebsocketCompletion(completedPayload) - completedPayload = patchCodexCompletedOutput(completedPayload, outputItemsByIndex, outputItemsFallback) + if !preserveNativeOutput { + completedPayload = patchCodexCompletedOutput(completedPayload, outputItemsByIndex, outputItemsFallback) + } if eventType != "response.incomplete" { cacheCodexReasoningReplayFromCompleted(replayScope, completedPayload) } @@ -664,7 +667,9 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr completedPayload := payload if eventType == "response.completed" || eventType == "response.done" || eventType == "response.incomplete" { completedPayload = normalizeCodexWebsocketCompletion(completedPayload) - completedPayload = patchCodexCompletedOutput(completedPayload, outputItemsByIndex, outputItemsFallback) + if !preserveNativeOutput { + completedPayload = patchCodexCompletedOutput(completedPayload, outputItemsByIndex, outputItemsFallback) + } if eventType != "response.incomplete" { cacheCodexReasoningReplayFromCompleted(replayScope, completedPayload) } diff --git a/internal/runtime/executor/helps/codex_native.go b/internal/runtime/executor/helps/codex_native.go new file mode 100644 index 000000000..bb02ea737 --- /dev/null +++ b/internal/runtime/executor/helps/codex_native.go @@ -0,0 +1,20 @@ +package helps + +import ( + "strings" + + "github.com/router-for-me/CLIProxyAPI/v7/internal/util" + cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" + sdktranslator "github.com/router-for-me/CLIProxyAPI/v7/sdk/translator" +) + +// IsNativeCodexRequest checks the client dialect for use inside Codex executors. +func IsNativeCodexRequest(body []byte, opts cliproxyexecutor.Options) bool { + for _, format := range []sdktranslator.Format{opts.SourceFormat, cliproxyexecutor.ResponseFormatOrSource(opts)} { + name := strings.TrimSpace(format.String()) + if !strings.EqualFold(name, sdktranslator.FormatCodex.String()) && !strings.EqualFold(name, sdktranslator.FormatOpenAIResponse.String()) { + return false + } + } + return util.IsCodexResponsesLiteRequest(body, opts.Headers) +} diff --git a/internal/util/codex.go b/internal/util/codex.go new file mode 100644 index 000000000..8d3bdcac1 --- /dev/null +++ b/internal/util/codex.go @@ -0,0 +1,17 @@ +package util + +import ( + "net/http" + "strings" + + "github.com/tidwall/gjson" +) + +// IsCodexResponsesLiteRequest recognizes the native header and its websocket metadata mirror. +func IsCodexResponsesLiteRequest(body []byte, headers http.Header) bool { + if strings.EqualFold(strings.TrimSpace(headers.Get("X-OpenAI-Internal-Codex-Responses-Lite")), "true") { + return true + } + value := gjson.GetBytes(body, "client_metadata.ws_request_header_x_openai_internal_codex_responses_lite") + return value.Type == gjson.True || value.Type == gjson.String && strings.EqualFold(strings.TrimSpace(value.String()), "true") +} diff --git a/sdk/api/handlers/openai/openai_responses_websocket.go b/sdk/api/handlers/openai/openai_responses_websocket.go index a265db068..6207cf054 100644 --- a/sdk/api/handlers/openai/openai_responses_websocket.go +++ b/sdk/api/handlers/openai/openai_responses_websocket.go @@ -16,6 +16,7 @@ import ( "github.com/google/uuid" "github.com/gorilla/websocket" "github.com/router-for-me/CLIProxyAPI/v7/internal/interfaces" + "github.com/router-for-me/CLIProxyAPI/v7/internal/util" "github.com/router-for-me/CLIProxyAPI/v7/sdk/api/handlers" coreauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" @@ -592,6 +593,8 @@ func (h *OpenAIResponsesAPIHandler) ResponsesWebsocket(c *gin.Context) { lastAttemptedAuthID := pinnedAuthID attemptedUpstreamMode := responsesWebsocketUpstreamModeUnknown selectedAuthObserved := false + nativeRequest := util.IsCodexResponsesLiteRequest(payload, c.Request.Header) + var preserveNativeOutput atomic.Bool pinnedAuthAttempted := false cliCtx, cliCancel := h.GetContextWithCancel(h, c, executionParent) cliCtx = cliproxyexecutor.WithDownstreamWebsocket(cliCtx) @@ -600,6 +603,7 @@ func (h *OpenAIResponsesAPIHandler) ResponsesWebsocket(c *gin.Context) { } cliCtx = handlers.WithExecutionSessionID(cliCtx, passthroughSessionID) cliCtx = handlers.WithSelectedAuthIDCallback(cliCtx, func(authID string) { + preserveNativeOutput.Store(false) authID = strings.TrimSpace(authID) if authID == "" || h == nil || h.AuthManager == nil { return @@ -612,6 +616,7 @@ func (h *OpenAIResponsesAPIHandler) ResponsesWebsocket(c *gin.Context) { return } attemptedUpstreamMode = upstreamModeForAuth(selectedAuth) + preserveNativeOutput.Store(nativeRequest && strings.EqualFold(strings.TrimSpace(selectedAuth.Provider), "codex")) }) if pinnedAuthID != "" && !routeOverridesModelResolution { cliCtx = handlers.WithPinnedAuthID(cliCtx, pinnedAuthID) @@ -638,8 +643,9 @@ func (h *OpenAIResponsesAPIHandler) ResponsesWebsocket(c *gin.Context) { wsTimelineLog, passthroughSessionID, responsesWebsocketForwardOptions{ - toolCacheTurn: toolCacheTurn, - suppressError: replayPinnedAuthFailure, + preserveCompletionOutput: preserveNativeOutput.Load, + toolCacheTurn: toolCacheTurn, + suppressError: replayPinnedAuthFailure, }, ) if errForward != nil { diff --git a/sdk/api/handlers/openai/openai_responses_websocket_forward.go b/sdk/api/handlers/openai/openai_responses_websocket_forward.go index da3666633..fa937b099 100644 --- a/sdk/api/handlers/openai/openai_responses_websocket_forward.go +++ b/sdk/api/handlers/openai/openai_responses_websocket_forward.go @@ -22,9 +22,10 @@ import ( ) type responsesWebsocketForwardOptions struct { - toolCacheTurn *responsesWebsocketToolCacheTurn - suppressError func(*interfaces.ErrorMessage) bool - keepAliveInterval *time.Duration + preserveCompletionOutput func() bool + toolCacheTurn *responsesWebsocketToolCacheTurn + suppressError func(*interfaces.ErrorMessage) bool + keepAliveInterval *time.Duration } func (h *OpenAIResponsesAPIHandler) forwardResponsesWebsocket( @@ -141,7 +142,7 @@ func (h *OpenAIResponsesAPIHandler) forwardResponsesWebsocket( for i := range payloads { collectResponsesWebsocketOutputItem(payloads[i], outputItemsByIndex, &outputItemsFallback) eventType := gjson.GetBytes(payloads[i], "type").String() - if isResponsesWebsocketCompletionEvent(eventType) { + if isResponsesWebsocketCompletionEvent(eventType) && (opts.preserveCompletionOutput == nil || !opts.preserveCompletionOutput()) { payloads[i] = restoreResponsesWebsocketCompletionOutput(payloads[i], outputItemsByIndex, outputItemsFallback) } if toolCacheTurn != nil { diff --git a/sdk/api/handlers/openai/openai_responses_websocket_test.go b/sdk/api/handlers/openai/openai_responses_websocket_test.go index c3a936054..7efcd2f5c 100644 --- a/sdk/api/handlers/openai/openai_responses_websocket_test.go +++ b/sdk/api/handlers/openai/openai_responses_websocket_test.go @@ -691,6 +691,7 @@ type websocketBootstrapFallbackExecutor struct { type websocketDirectCaptureExecutor struct { mu sync.Mutex provider string + streamItems bool failStatus int authIDs []string models []string @@ -806,7 +807,7 @@ func (e *websocketDirectCaptureExecutor) ExecuteStream(ctx context.Context, auth failStatus := e.failStatus e.mu.Unlock() - chunks := make(chan coreexecutor.StreamChunk, 1) + chunks := make(chan coreexecutor.StreamChunk, 2) if failStatus > 0 { chunks <- coreexecutor.StreamChunk{Err: websocketPinnedFailoverStatusError{ status: failStatus, @@ -816,7 +817,12 @@ func (e *websocketDirectCaptureExecutor) ExecuteStream(ctx context.Context, auth return &coreexecutor.StreamResult{Chunks: chunks}, nil } responseID := fmt.Sprintf("resp-%d", count) - chunks <- coreexecutor.StreamChunk{Payload: []byte(fmt.Sprintf(`{"type":"response.completed","response":{"id":%q,"output":[{"type":"message","id":"out-%d"}]}}`, responseID, count))} + output := fmt.Sprintf(`[{"type":"message","id":"out-%d"}]`, count) + if e.streamItems { + chunks <- coreexecutor.StreamChunk{Payload: []byte(fmt.Sprintf(`{"type":"response.output_item.done","output_index":0,"item":{"type":"message","id":"out-%d"}}`, count))} + output = "[]" + } + chunks <- coreexecutor.StreamChunk{Payload: []byte(fmt.Sprintf(`{"type":"response.completed","response":{"id":%q,"output":%s}}`, responseID, output))} close(chunks) if count >= 2 && e.done != nil { e.doneOnce.Do(func() { @@ -2456,6 +2462,15 @@ func TestRecordResponsesWebsocketCustomToolCallsFromOutputItemDoneWithCache(t *t } func TestForwardResponsesWebsocketRestoresAndForwardsCompletedOutput(t *testing.T) { + for _, preserve := range []bool{false, true} { + t.Run(fmt.Sprintf("preserve=%t", preserve), func(t *testing.T) { + testForwardResponsesWebsocketCompletedOutput(t, preserve) + }) + } +} + +func testForwardResponsesWebsocketCompletedOutput(t *testing.T, preserve bool) { + t.Helper() gin.SetMode(gin.TestMode) serverErrCh := make(chan error, 1) @@ -2491,6 +2506,7 @@ func TestForwardResponsesWebsocketRestoresAndForwardsCompletedOutput(t *testing. errCh, timelineLog, "session-1", + responsesWebsocketForwardOptions{preserveCompletionOutput: func() bool { return preserve }}, ) if err != nil { serverErrCh <- err @@ -2550,7 +2566,11 @@ func TestForwardResponsesWebsocketRestoresAndForwardsCompletedOutput(t *testing. if strings.Contains(string(payload), "response.done") { t.Fatalf("payload unexpectedly rewrote completed event: %s", payload) } - if got := gjson.GetBytes(payload, "response.output.0.id").String(); got != "call-1" { + if preserve { + if string(payload) != `{"type":"response.completed","response":{"id":"resp-1","output":[]}}` { + t.Fatalf("native completion changed: %s", payload) + } + } else if got := gjson.GetBytes(payload, "response.output.0.id").String(); got != "call-1" { t.Fatalf("downstream completion output id = %q, want call-1; payload=%s", got, payload) } @@ -3494,8 +3514,8 @@ func TestResponsesWebsocketFullRequestCanRouteFromNativeWebsocketToBuiltInProvid const sourceModel = "codex-provider-route-source" const targetModel = "claude-provider-route-target" - codexExecutor := &websocketDirectCaptureExecutor{provider: "codex"} - claudeExecutor := &websocketDirectCaptureExecutor{provider: "claude"} + codexExecutor := &websocketDirectCaptureExecutor{provider: "codex", streamItems: true} + claudeExecutor := &websocketDirectCaptureExecutor{provider: "claude", streamItems: true} manager := coreauth.NewManager(nil, nil, nil) manager.RegisterExecutor(codexExecutor) manager.RegisterExecutor(claudeExecutor) @@ -3537,18 +3557,25 @@ func TestResponsesWebsocketFullRequestCanRouteFromNativeWebsocketToBuiltInProvid } defer func() { _ = conn.Close() }() - firstRequest := []byte(fmt.Sprintf(`{"type":"response.create","model":%q,"input":[{"type":"message","id":"msg-1"}]}`, sourceModel)) + firstRequest := []byte(fmt.Sprintf(`{"type":"response.create","model":%q,"input":[{"type":"message","id":"msg-1"}],"client_metadata":{"ws_request_header_x_openai_internal_codex_responses_lite":"true"}}`, sourceModel)) if errWrite := conn.WriteMessage(websocket.TextMessage, firstRequest); errWrite != nil { t.Fatalf("write first websocket message: %v", errWrite) } if _, _, errRead := conn.ReadMessage(); errRead != nil { t.Fatalf("read first websocket response: %v", errRead) } + _, nativeResponse, errRead := conn.ReadMessage() + if errRead != nil || gjson.GetBytes(nativeResponse, "response.output").Raw != "[]" { + t.Fatalf("native completion = %s, error = %v", nativeResponse, errRead) + } - routedRequest := []byte(fmt.Sprintf(`{"type":"response.create","model":%q,"route_to_claude":true,"input":[{"type":"message","id":"msg-routed"}]}`, sourceModel)) + routedRequest := []byte(fmt.Sprintf(`{"type":"response.create","model":%q,"route_to_claude":true,"input":[{"type":"message","id":"msg-routed"}],"client_metadata":{"ws_request_header_x_openai_internal_codex_responses_lite":"true"}}`, sourceModel)) if errWrite := conn.WriteMessage(websocket.TextMessage, routedRequest); errWrite != nil { t.Fatalf("write routed websocket message: %v", errWrite) } + if _, _, errRead := conn.ReadMessage(); errRead != nil { + t.Fatalf("read routed output item: %v", errRead) + } _, response, errRead := conn.ReadMessage() if errRead != nil { t.Fatalf("read routed websocket response: %v", errRead) @@ -3556,6 +3583,10 @@ func TestResponsesWebsocketFullRequestCanRouteFromNativeWebsocketToBuiltInProvid if got := gjson.GetBytes(response, "type").String(); got != wsEventTypeCompleted { t.Fatalf("routed response type = %q, want %q: %s", got, wsEventTypeCompleted, response) } + t.Logf("native completion: %s; routed completion: %s", nativeResponse, response) + if got := gjson.GetBytes(response, "response.output.0.id").String(); got != "out-1" { + t.Fatalf("cross-provider output repair lost: %s", response) + } if got := len(codexExecutor.Payloads()); got != 1 { t.Fatalf("codex payload count = %d, want 1", got) }