package chatgpt import ( "aurora/httpclient" "aurora/internal/accounts" "aurora/internal/sseparser" "aurora/typings/chatgpt" "encoding/json" "fmt" "io" "net/http" "net/http/httptest" "net/url" "strings" "sync" "testing" "time" fhttp "github.com/bogdanfinn/fhttp" "github.com/bogdanfinn/websocket" "github.com/gin-gonic/gin" ) type fakeAuroraClient struct { response *http.Response headers httpclient.AuroraHeaders body string } func (f *fakeAuroraClient) Request(method httpclient.HttpMethod, url string, headers httpclient.AuroraHeaders, cookies []*http.Cookie, body io.Reader) (*http.Response, error) { f.headers = headers if body != nil { data, _ := io.ReadAll(body) f.body = string(data) } return f.response, nil } func (f *fakeAuroraClient) SetProxy(url string) error { return nil } func (f *fakeAuroraClient) SetCookies(rawUrl string, cookies []*http.Cookie) { } func (f *fakeAuroraClient) GetCookies(rawUrl string) []*http.Cookie { return nil } func TestHandlerStreamsPatchEvents(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"p":"/conversation_id","o":"replace","v":"conv-test"}`, `data: {"p":"/message/id","o":"replace","v":"msg-test"}`, `data: {"p":"/message/content/parts/0","o":"append","v":"hello"}`, `data: {"p":"/message/content/parts/0","o":"append","v":" world"}`, `data: {"p":"/message/metadata/finish_details","o":"replace","v":{"type":"stop"}}`, `data: {"p":"/message/end_turn","o":"replace","v":true}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) full, continueInfo := Handler(c, response, nil, nil, "request-id", chatGPTRequestForTest(), true, "auto") if continueInfo != nil { t.Fatalf("continueInfo = %#v, want nil", continueInfo) } if full != "hello world" { t.Fatalf("full response = %q, want %q", full, "hello world") } output := writer.Body.String() if !strings.Contains(output, `"content":"hello"`) { t.Fatalf("stream output does not contain first delta: %s", output) } if !strings.Contains(output, `"content":" world"`) { t.Fatalf("stream output does not contain second delta: %s", output) } if !strings.Contains(output, `"finish_reason":"stop"`) { t.Fatalf("stream output does not contain stop chunk: %s", output) } } func TestHandlerStreamsOpenAIChunksAndSentinel(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"choices":[{"delta":{"role":"assistant"}}]}`, `data: {"choices":[{"delta":{"content":"你"}}]}`, `data: {"choices":[{"delta":{"content":"好"}}]}`, `data: {"choices":[{"delta":{}}],"sentinel":{"event":"artifact_pending","kind":"generated_image","title":"正在生成图片..."}}`, `data: {"choices":[{"delta":{}}],"sentinel":{"event":"artifact","kind":"generated_image","slot_index":1,"revision":2,"file_id":"file-yyy","url":"http://example.test/image.png"}}`, `data: {"choices":[{"delta":{},"finish_reason":"stop"}],"conversation_id":"conv-xxx"}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) request := chatGPTRequestForTest() full, continueInfo := Handler(c, response, nil, nil, "request-id", request, true, request.Model) if continueInfo != nil { t.Fatalf("continueInfo = %#v, want nil", continueInfo) } if full != "你好" { t.Fatalf("full response = %q, want %q", full, "你好") } chunks := parseSSEChunks(t, writer.Body.String()) if len(chunks) != 6 { t.Fatalf("chunk count = %d, want 6; output: %s", len(chunks), writer.Body.String()) } if chunks[0]["choices"].([]interface{})[0].(map[string]interface{})["delta"].(map[string]interface{})["role"] != "assistant" { t.Fatalf("first chunk should include assistant role: %#v", chunks[0]) } if chunks[1]["choices"].([]interface{})[0].(map[string]interface{})["delta"].(map[string]interface{})["content"] != "你" { t.Fatalf("second chunk should include first content delta: %#v", chunks[1]) } if chunks[3]["sentinel"].(map[string]interface{})["event"] != "artifact_pending" { t.Fatalf("sentinel pending chunk missing: %#v", chunks[3]) } if chunks[4]["sentinel"].(map[string]interface{})["file_id"] != "file-yyy" { t.Fatalf("sentinel artifact chunk missing file_id: %#v", chunks[4]) } if chunks[5]["conversation_id"] != "conv-xxx" { t.Fatalf("stop chunk conversation_id = %#v, want conv-xxx", chunks[5]["conversation_id"]) } if chunks[5]["choices"].([]interface{})[0].(map[string]interface{})["finish_reason"] != "stop" { t.Fatalf("stop chunk finish_reason missing: %#v", chunks[5]) } } func TestHandlerStreamsConcatenatedOpenAIChunks(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"id":"chunk-1","object":"chat.completion.chunk","created":0,"model":"auto","choices":[{"delta":{"content":"is","role":"assistant"},"index":0,"finish_reason":null}]}`, `data: {"id":"chunk-2","object":"chat.completion.chunk","created":0,"model":"auto","choices":[{"delta":{"content":"This is a test!"},"index":0,"finish_reason":null}]}`, `data: {"id":"chunk-3","object":"chat.completion.chunk","created":0,"model":"auto","choices":[{"delta":{},"index":0,"finish_reason":"stop"}],"conversation_id":"conv-stream"}`, `data: [DONE]`, }, "") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) result := HandlerDetailed(c, response, nil, nil, "request-id", chatGPTRequestForTest(), true, "auto") if result.Text != "isThis is a test!" { t.Fatalf("text = %q, want concatenated deltas", result.Text) } if result.ConversationID != "conv-stream" { t.Fatalf("conversationID = %q, want conv-stream", result.ConversationID) } if !result.StopSent { t.Fatalf("StopSent = false, want true") } chunks := parseSSEChunks(t, writer.Body.String()) if len(chunks) != 3 { t.Fatalf("chunk count = %d, want 3; output: %s", len(chunks), writer.Body.String()) } firstDelta := chunks[0]["choices"].([]interface{})[0].(map[string]interface{})["delta"].(map[string]interface{}) if firstDelta["role"] != "assistant" || firstDelta["content"] != "is" { t.Fatalf("first chunk delta = %#v, want role+content in one chunk", firstDelta) } if chunks[1]["choices"].([]interface{})[0].(map[string]interface{})["delta"].(map[string]interface{})["content"] != "This is a test!" { t.Fatalf("second chunk should include second content delta: %#v", chunks[1]) } if chunks[2]["conversation_id"] != "conv-stream" { t.Fatalf("stop chunk conversation_id = %#v, want conv-stream", chunks[2]["conversation_id"]) } } func TestStreamHandoffTopicFromPayload(t *testing.T) { payload := `{"type":"stream_handoff","options":[{"type":"subscribe_ws_topic","topic_id":"conversation-turn-abc"}]}` topicID, skip := sseparser.HandoffTopicFromPayload(payload, "") if !skip { t.Fatalf("skip = false, want true") } if topicID != "conversation-turn-abc" { t.Fatalf("topicID = %q, want conversation-turn-abc", topicID) } topicID, skip = sseparser.HandoffTopicFromPayload(`{"metadata":{"turn_exchange_id":"xyz"}}`, "server_ste_metadata") if !skip || topicID != "conversation-turn-xyz" { t.Fatalf("server metadata topic = %q skip=%v, want conversation-turn-xyz true", topicID, skip) } } func TestChatWebsocketEncodedItem(t *testing.T) { frame := map[string]interface{}{ "type": "message", "topic_id": "conversation-turn-abc", "payload": map[string]interface{}{ "payload": map[string]interface{}{ "encoded_item": "data: {\"v\":\"hello\"}\n", }, }, } encoded := chatWebsocketEncodedItem(frame, "conversation-turn-abc") if encoded != "data: {\"v\":\"hello\"}\n" { t.Fatalf("encoded = %q, want websocket encoded item", encoded) } if chatWebsocketEncodedItem(frame, "conversation-turn-other") != "" { t.Fatalf("encoded item should be ignored for other topic") } } func TestHandlerStreamsShortDeltaEncoding(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `event: delta_encoding`, `data: "v1"`, `event: delta`, `data: {"v":"hello"}`, `event: delta`, `data: {"v":" world"}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) result := HandlerDetailed(c, response, nil, nil, "request-id", chatGPTRequestForTest(), true, "auto") if result.Text != "hello world" { t.Fatalf("text = %q, want hello world", result.Text) } if !strings.Contains(writer.Body.String(), `"content":"hello"`) { t.Fatalf("stream output missing short delta content: %s", writer.Body.String()) } } func TestHandlerCollectsPatchEventsWithoutStreaming(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"p":"/message/content/parts/0","o":"append","v":"hello"}`, `data: {"p":"/message/content/parts/0","o":"append","v":" world"}`, `data: {"p":"/message/end_turn","o":"replace","v":true}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) full, continueInfo := Handler(c, response, nil, nil, "request-id", chatGPTRequestForTest(), false, "auto") if continueInfo != nil { t.Fatalf("continueInfo = %#v, want nil", continueInfo) } if full != "hello world" { t.Fatalf("full response = %q, want %q", full, "hello world") } if writer.Body.Len() != 0 { t.Fatalf("non-streaming handler wrote body %q, want no direct write", writer.Body.String()) } } func TestHandlerDetailedCollectsOpenAIChunkMetadataWithoutStreaming(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"choices":[{"delta":{"role":"assistant"}}]}`, `data: {"choices":[{"delta":{"content":"你"}}]}`, `data: {"choices":[{"delta":{"content":"好"}}]}`, `data: {"choices":[{"delta":{}}],"sentinel":{"event":"artifact","kind":"generated_image","slot_index":1,"url":"http://example.test/image.png"}}`, `data: {"choices":[{"delta":{},"finish_reason":"stop"}],"conversation_id":"conv-xxx"}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) request := chatGPTRequestForTest() result := HandlerDetailed(c, response, nil, nil, "request-id", request, false, request.Model) if result.Text != "你好" { t.Fatalf("text = %q, want %q", result.Text, "你好") } if result.ConversationID != "conv-xxx" { t.Fatalf("conversationID = %q, want conv-xxx", result.ConversationID) } if !result.StopSent { t.Fatalf("StopSent = false, want true") } if len(result.Sentinel) != 1 { t.Fatalf("sentinel count = %d, want 1", len(result.Sentinel)) } if result.Sentinel[0]["event"] != "artifact" || result.Sentinel[0]["url"] != "http://example.test/image.png" { t.Fatalf("sentinel = %#v, want artifact with url", result.Sentinel[0]) } if writer.Body.Len() != 0 { t.Fatalf("non-streaming handler wrote body %q, want no direct write", writer.Body.String()) } } func TestHandlerDetailedDoesNotMarkStopWhenUpstreamOnlySendsDone(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"choices":[{"delta":{"role":"assistant"},"index":0,"finish_reason":null}]}`, `data: {"choices":[{"delta":{"content":"This"},"index":0,"finish_reason":null}]}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) result := HandlerDetailed(c, response, nil, nil, "request-id", chatGPTRequestForTest(), true, "auto") if result.Text != "This" { t.Fatalf("text = %q, want This", result.Text) } if result.StopSent { t.Fatalf("StopSent = true, want false so caller can emit stop before [DONE]") } output := writer.Body.String() if strings.Contains(output, `"finish_reason":"stop"`) { t.Fatalf("handler should not invent stop internally for upstream DONE-only stream: %s", output) } if !strings.Contains(output, `"content":"This"`) { t.Fatalf("stream output missing content chunk: %s", output) } } func TestGetConduitTokenAllowsNullToken(t *testing.T) { client := &fakeAuroraClient{ response: &http.Response{ StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"status":"ok","conduit_token":null}`)), }, } token, err := getConduitToken(client, chatGPTRequestForTest(), &accounts.Account{}, nil, "trace-id") if err != nil { t.Fatalf("getConduitToken returned error for null token: %v", err) } if token != "" { t.Fatalf("token = %q, want empty token", token) } if client.headers["x-conduit-token"] != "" { t.Fatalf("prepare x-conduit-token = %q, want empty (no-token 字面量已废弃,真实浏览器首次 prepare 送空串)", client.headers["x-conduit-token"]) } } func TestPrepareConversationConduitDoesNotUseSentinelHeaders(t *testing.T) { client := &fakeAuroraClient{ response: &http.Response{ StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"status":"ok","conduit_token":"abc"}`)), }, } token, err := PrepareConversationConduit(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id") if err != nil { t.Fatalf("PrepareConversationConduit returned error: %v", err) } if token != "abc" { t.Fatalf("token = %q, want abc", token) } for _, key := range []string{"openai-sentinel-chat-requirements-token", "openai-sentinel-proof-token", "openai-sentinel-turnstile-token"} { if _, ok := client.headers[key]; ok { t.Fatalf("prepare conduit unexpectedly includes %s", key) } } } // sequentialAuroraClient 是一个按调用顺序返回不同响应的 fake client, // 用于验证 PrepareConversationConduitFull 的三态时序和 token 链路。 type sequentialAuroraClient struct { mu sync.Mutex responses []*http.Response requests []capturedRequest } type capturedRequest struct { url string headers httpclient.AuroraHeaders body string } func (s *sequentialAuroraClient) Request(method httpclient.HttpMethod, url string, headers httpclient.AuroraHeaders, cookies []*http.Cookie, body io.Reader) (*http.Response, error) { s.mu.Lock() defer s.mu.Unlock() idx := len(s.requests) var bodyStr string if body != nil { data, _ := io.ReadAll(body) bodyStr = string(data) } s.requests = append(s.requests, capturedRequest{url: url, headers: headers, body: bodyStr}) if idx >= len(s.responses) { return nil, fmt.Errorf("unexpected request #%d: %s", idx, url) } return s.responses[idx], nil } func (s *sequentialAuroraClient) SetProxy(url string) error { return nil } func (s *sequentialAuroraClient) SetCookies(rawUrl string, cookies []*http.Cookie) {} func (s *sequentialAuroraClient) GetCookies(rawUrl string) []*http.Cookie { // 返回一个 cf_clearance,让 ensureBootstrapped 直接走 fast-path return []*http.Cookie{{Name: "cf_clearance", Value: "test-cf"}} } func (s *sequentialAuroraClient) requestBodies() []string { out := make([]string, len(s.requests)) for i, r := range s.requests { out[i] = r.body } return out } func (s *sequentialAuroraClient) requestHeaders() []httpclient.AuroraHeaders { out := make([]httpclient.AuroraHeaders, len(s.requests)) for i, r := range s.requests { out[i] = r.headers } return out } // TestPrepareConversationConduitFullChainsThreeStates 验证三态 prepare: // 1) 严格发起三次 HTTP 请求,顺序 none -> sent -> success // 2) 第 2/3 次的 X-Conduit-Token 头分别是上一次响应中的 token // 3) client_prepare_state 字段在三次请求中分别为 none / sent / success // 4) partial_query 仅在 sent / success 出现,none 阶段不带 // 5) 返回值是第三次响应的 conduit_token func TestPrepareConversationConduitFullChainsThreeStates(t *testing.T) { client := &sequentialAuroraClient{ responses: []*http.Response{ {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"token-step-1"}`))}, {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"token-step-2"}`))}, {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"token-step-3"}`))}, }, } msg := chatGPTRequestForTest() msg.AddMessage("user", "hello world") // partial_query 取前 5 rune = "hello" token, err := PrepareConversationConduitFull(client, msg, &accounts.Account{}, "", "trace-id", nil) if err != nil { t.Fatalf("PrepareConversationConduitFull returned error: %v", err) } // (1) 调用次数 if got := len(client.requests); got != 3 { t.Fatalf("request count = %d, want 3", got) } // (2) X-Conduit-Token 链路 headers := client.requestHeaders() if got := headers[0]["X-Conduit-Token"]; got != "" { t.Fatalf("step1 X-Conduit-Token = %q, want empty (首次无前驱 token)", got) } if got := headers[1]["X-Conduit-Token"]; got != "token-step-1" { t.Fatalf("step2 X-Conduit-Token = %q, want token-step-1", got) } if got := headers[2]["X-Conduit-Token"]; got != "token-step-2" { t.Fatalf("step3 X-Conduit-Token = %q, want token-step-2", got) } // (3) client_prepare_state 序列 bodies := client.requestBodies() states := make([]string, 3) hasPartial := make([]bool, 3) for i, b := range bodies { var p map[string]interface{} if err := json.Unmarshal([]byte(b), &p); err != nil { t.Fatalf("step%d body is invalid json: %v", i+1, err) } states[i], _ = p["client_prepare_state"].(string) _, hasPartial[i] = p["partial_query"] } wantStates := []string{"none", "sent", "success"} for i, want := range wantStates { if states[i] != want { t.Fatalf("step%d client_prepare_state = %q, want %q", i+1, states[i], want) } } // (4) partial_query: none 阶段不带,sent/success 带 if hasPartial[0] { t.Fatalf("step1 (none) should not carry partial_query, body=%s", bodies[0]) } if !hasPartial[1] || !hasPartial[2] { t.Fatalf("step2/3 should carry partial_query, has=[none=%v sent=%v success=%v]", hasPartial[0], hasPartial[1], hasPartial[2]) } // (5) 返回值 = 第三次响应的 token if token != "token-step-3" { t.Fatalf("returned token = %q, want token-step-3", token) } } // TestPrepareConversationConduitFullFailsFastOnStep1 验证第一步失败时 // 直接返回带 "prepare(none) failed:" 前缀的错误,不会发起后续请求。 func TestPrepareConversationConduitFullFailsFastOnStep1(t *testing.T) { client := &sequentialAuroraClient{ responses: []*http.Response{ {StatusCode: http.StatusInternalServerError, Body: io.NopCloser(strings.NewReader(`{"error":"upstream"}`))}, }, } _, err := PrepareConversationConduitFull(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id", nil) if err == nil { t.Fatalf("expected error, got nil") } if !strings.Contains(err.Error(), "prepare(none) failed") { t.Fatalf("error = %q, want prefix prepare(none) failed", err.Error()) } if got := len(client.requests); got != 1 { t.Fatalf("request count = %d, want 1 (fail-fast on step1)", got) } } // TestPrepareConversationConduitFullFailsFastOnStep2 验证第二步失败时报 // "prepare(sent) failed",且只发起了 2 次请求。 func TestPrepareConversationConduitFullFailsFastOnStep2(t *testing.T) { client := &sequentialAuroraClient{ responses: []*http.Response{ {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"token-step-1"}`))}, {StatusCode: http.StatusForbidden, Body: io.NopCloser(strings.NewReader(`{"error":"forbidden"}`))}, }, } _, err := PrepareConversationConduitFull(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id", nil) if err == nil { t.Fatalf("expected error, got nil") } if !strings.Contains(err.Error(), "prepare(sent) failed") { t.Fatalf("error = %q, want prefix prepare(sent) failed", err.Error()) } if got := len(client.requests); got != 2 { t.Fatalf("request count = %d, want 2 (fail-fast on step2)", got) } } // TestPrepareConversationConduitFullFailsFastOnStep3 验证第三步失败时报 // "prepare(success) failed",且发起了 3 次请求。 func TestPrepareConversationConduitFullFailsFastOnStep3(t *testing.T) { client := &sequentialAuroraClient{ responses: []*http.Response{ {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"t1"}`))}, {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"t2"}`))}, {StatusCode: http.StatusBadGateway, Body: io.NopCloser(strings.NewReader(`bad gateway`))}, }, } _, err := PrepareConversationConduitFull(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id", nil) if err == nil { t.Fatalf("expected error, got nil") } if !strings.Contains(err.Error(), "prepare(success) failed") { t.Fatalf("error = %q, want prefix prepare(success) failed", err.Error()) } if got := len(client.requests); got != 3 { t.Fatalf("request count = %d, want 3 (fail-fast on step3)", got) } } // TestPrepareConversationConduitFullPropagatesState 验证 state.DeviceID / // SessionID 在三次请求中保持一致(整段对话使用同一对 fingerprint ID)。 func TestPrepareConversationConduitFullPropagatesState(t *testing.T) { client := &sequentialAuroraClient{ responses: []*http.Response{ {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"t1"}`))}, {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"t2"}`))}, {StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"conduit_token":"t3"}`))}, }, } state := NewChatClientState() state.DeviceID = "device-3state" state.SessionID = "session-3state" state.ParentMessageID = "parent-3state" _, err := PrepareConversationConduitFull(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id", state) if err != nil { t.Fatalf("PrepareConversationConduitFull returned error: %v", err) } headers := client.requestHeaders() for i, h := range headers { if got := h["Oai-Device-Id"]; got != "device-3state" { t.Fatalf("step%d Oai-Device-Id = %q, want device-3state", i+1, got) } if got := h["Oai-Session-Id"]; got != "session-3state" { t.Fatalf("step%d Oai-Session-Id = %q, want session-3state", i+1, got) } } } func TestConversationHeadersKeepEmptyConduitHeaderForConversation(t *testing.T) { header := conversationHeaders(&accounts.Account{}, nil, "text/event-stream", "/backend-api/f/conversation", "", "trace-id") if _, ok := header["X-Conduit-Token"]; !ok { t.Fatalf("X-Conduit-Token header missing for empty conversation conduit token") } if header["X-Conduit-Token"] != "" { t.Fatalf("X-Conduit-Token = %q, want empty string", header["X-Conduit-Token"]) } } func TestRequiresConversationWebsocket(t *testing.T) { tests := []struct { name string stream bool thinkingEffort string want bool }{ {name: "streaming standard", stream: true, thinkingEffort: "standard", want: true}, {name: "non-streaming standard", thinkingEffort: "standard", want: false}, {name: "non-streaming low normalizes to standard", thinkingEffort: "low", want: false}, {name: "non-streaming extended", thinkingEffort: "extended", want: true}, {name: "non-streaming medium normalizes to extended", thinkingEffort: "medium", want: true}, {name: "non-streaming max", thinkingEffort: "max", want: true}, {name: "non-streaming high normalizes to max", thinkingEffort: "high", want: true}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { if got := RequiresConversationWebsocket(tt.stream, tt.thinkingEffort); got != tt.want { t.Fatalf("RequiresConversationWebsocket(%v, %q) = %v, want %v", tt.stream, tt.thinkingEffort, got, tt.want) } }) } } func TestShouldUseWebsocketHandoffSkipsCatchupAfterHTTPBody(t *testing.T) { if shouldUseWebsocketHandoff(false, "conversation-turn-abc", &websocket.Conn{}, "hello", nil) { t.Fatalf("websocket handoff should be skipped when HTTP SSE already emitted text") } if shouldUseWebsocketHandoff(false, "conversation-turn-abc", &websocket.Conn{}, "", []string{"![image](url)"}) { t.Fatalf("websocket handoff should be skipped when HTTP SSE already emitted image content") } if !shouldUseWebsocketHandoff(false, "conversation-turn-abc", &websocket.Conn{}, "", nil) { t.Fatalf("websocket handoff should be used when HTTP SSE has no body yet") } if shouldUseWebsocketHandoff(true, "conversation-turn-abc", &websocket.Conn{}, "", nil) { t.Fatalf("websocket handoff should not be used while already reading websocket") } } func TestCreateBaseHeaderMatchesWebClientShape(t *testing.T) { first := createBaseHeader() second := createBaseHeader() // 对齐 conversation.txt 2026-06 抓包:英文浏览器 (en-US,en) if first["Oai-Language"] != "en-US" { t.Fatalf("Oai-Language = %q, want en-US", first["Oai-Language"]) } if first["Accept-Language"] != "en-US,en;q=0.9" { t.Fatalf("Accept-Language = %q, want en-US,en;q=0.9", first["Accept-Language"]) } // UA must be the Chrome 148 variant to match sec-ch-ua="Google Chrome";"v="148" ua := first["User-Agent"] if !strings.Contains(ua, "Chrome/148.") { t.Fatalf("User-Agent = %q, want Chrome 148 to match sec-ch-ua=Chrome 148", ua) } if strings.Contains(ua, "Edg/") { t.Fatalf("User-Agent = %q, must not be Edge variant (conversation.txt uses Chrome)", ua) } if first["Oai-Device-Id"] == "" || first["Oai-Device-Id"] != second["Oai-Device-Id"] { t.Fatalf("Oai-Device-Id should be stable across headers: first=%q second=%q", first["Oai-Device-Id"], second["Oai-Device-Id"]) } if first["Oai-Session-Id"] == "" || first["Oai-Session-Id"] != second["Oai-Session-Id"] { t.Fatalf("Oai-Session-Id should be stable across headers: first=%q second=%q", first["Oai-Session-Id"], second["Oai-Session-Id"]) } // 对齐 2026-06-26 浏览器抓包 if first["Oai-Client-Version"] != "prod-dbbd612ddb47498515c3eecf8579bcafa0066e07" { t.Fatalf("Oai-Client-Version = %q, want prod-dbbd612ddb47498515c3eecf8579bcafa0066e07", first["Oai-Client-Version"]) } if first["Oai-Client-Build-Number"] != "7823760" { t.Fatalf("Oai-Client-Build-Number = %q, want 7823760", first["Oai-Client-Build-Number"]) } // sec-ch-ua 必须跟 UA 一致(都是 Chrome 148) if first["Sec-Ch-Ua"] != `"Chromium";v="148", "Google Chrome";v="148", "Not/A)Brand";v="99"` { t.Fatalf("Sec-Ch-Ua = %q, want Chrome 148", first["Sec-Ch-Ua"]) } } func TestPrepareConversationConduitUsesClientState(t *testing.T) { client := &fakeAuroraClient{ response: &http.Response{ StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"status":"ok","conduit_token":"abc"}`)), }, } state := NewChatClientState() state.DeviceID = "device-state" state.SessionID = "session-state" state.ParentMessageID = "parent-state" state.StartTime = time.Now().Add(-2500 * time.Millisecond) token, err := PrepareConversationConduitWithState(client, chatGPTRequestForTest(), &accounts.Account{}, "", "trace-id", state) if err != nil { t.Fatalf("PrepareConversationConduitWithState returned error: %v", err) } if token != "abc" { t.Fatalf("token = %q, want abc", token) } if client.headers["Oai-Device-Id"] != "device-state" { t.Fatalf("Oai-Device-Id = %q, want state device", client.headers["Oai-Device-Id"]) } if client.headers["Oai-Session-Id"] != "session-state" { t.Fatalf("Oai-Session-Id = %q, want state session", client.headers["Oai-Session-Id"]) } var payload map[string]interface{} if err := json.Unmarshal([]byte(client.body), &payload); err != nil { t.Fatalf("prepare body is invalid json: %v", err) } if payload["parent_message_id"] != "parent-state" { t.Fatalf("parent_message_id = %#v, want parent-state", payload["parent_message_id"]) } if payload["supports_buffering"] != true { t.Fatalf("supports_buffering = %#v, want prepare capability", payload["supports_buffering"]) } if _, ok := payload["supported_encodings"].([]interface{}); !ok { t.Fatalf("supported_encodings missing: %#v", payload) } contextInfo, ok := payload["client_contextual_info"].(map[string]interface{}) if !ok { t.Fatalf("client_contextual_info missing: %#v", payload) } loaded, ok := contextInfo["time_since_loaded"].(float64) if !ok || loaded < 2 || loaded > 5 { t.Fatalf("time_since_loaded = %#v, want dynamic seconds around 3", contextInfo["time_since_loaded"]) } } func TestConversationCompletionOmitsRiskyCapabilityFields(t *testing.T) { client := &fakeAuroraClient{ response: &http.Response{ StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader("")), }, } request := chatGPTRequestForTest() request.EnableMessageFollowups = true request.SupportsBuffering = true request.SupportedEncodings = []string{"v1"} request.ClientContextualInfo = map[string]interface{}{"app_name": "chatgpt.com"} request.ThinkingEffort = "extended" response, err := POSTconversationPreparedWithState(client, request, &accounts.Account{}, nil, "", "conduit", "trace", NewChatClientState()) if err != nil { t.Fatalf("POSTconversationPreparedWithState returned error: %v", err) } if response == nil || response.StatusCode != http.StatusOK { t.Fatalf("response = %#v, want HTTP 200", response) } var payload map[string]interface{} if err := json.Unmarshal([]byte(client.body), &payload); err != nil { t.Fatalf("completion body is invalid json: %v", err) } for _, key := range []string{"enable_message_followups", "supports_buffering", "supported_encodings", "client_contextual_info"} { if _, ok := payload[key]; ok { t.Fatalf("completion request unexpectedly includes %s: %#v", key, payload[key]) } } if payload["client_prepare_state"] != "success" { t.Fatalf("client_prepare_state = %#v, want success", payload["client_prepare_state"]) } if payload["thinking_effort"] != "extended" { t.Fatalf("thinking_effort = %#v, want extended", payload["thinking_effort"]) } } func TestConversationCompletionNormalizesThinkingEffort(t *testing.T) { for _, tt := range []struct { input string want string }{ {input: "low", want: "standard"}, {input: "medium", want: "extended"}, {input: "high", want: "max"}, {input: "turbo", want: "standard"}, } { t.Run(tt.input, func(t *testing.T) { client := &fakeAuroraClient{ response: &http.Response{ StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader("")), }, } request := chatGPTRequestForTest() request.ThinkingEffort = tt.input _, err := POSTconversationPreparedWithState(client, request, &accounts.Account{}, nil, "", "conduit", "trace", NewChatClientState()) if err != nil { t.Fatalf("POSTconversationPreparedWithState returned error: %v", err) } var payload map[string]interface{} if err := json.Unmarshal([]byte(client.body), &payload); err != nil { t.Fatalf("completion body is invalid json: %v", err) } if payload["thinking_effort"] != tt.want { t.Fatalf("thinking_effort = %#v, want %q", payload["thinking_effort"], tt.want) } }) } } func TestChatWebsocketConversationUpdateItem(t *testing.T) { frame := map[string]interface{}{ "type": "message", "topic_id": "conversations", "payload": map[string]interface{}{ "type": "conversation-update", "conversation_id": "conv-img", "payload": map[string]interface{}{ "asset_pointer": "sediment://file-img123", }, }, } items := chatWebsocketSSEItems(frame, "conversation-turn-abc") if len(items) != 1 { t.Fatalf("items = %#v, want one conversation-update SSE item", items) } payloads := sseparser.DataPayloads(items[0]) if len(payloads) != 1 || !strings.Contains(payloads[0], "conversation-update") { t.Fatalf("payloads = %#v, want conversation-update data", payloads) } var raw map[string]interface{} if err := json.Unmarshal([]byte(payloads[0]), &raw); err != nil { t.Fatalf("conversation-update item is invalid json: %v", err) } signals := ExtractSignalsFromJSON(raw) if len(ImageFileIDsFromSignals(signals)) != 1 || ImageFileIDsFromSignals(signals)[0] != "file-img123" { t.Fatalf("signals = %#v, want image file id", signals) } } func TestWebsocketProxyFuncUsesExplicitProxy(t *testing.T) { proxyFunc, err := websocketProxyFunc("http://127.0.0.1:10808") if err != nil { t.Fatalf("websocketProxyFunc returned error: %v", err) } reqURL, _ := url.Parse("wss://example.test/ws") req := &fhttp.Request{URL: reqURL} proxyURL, err := proxyFunc(req) if err != nil { t.Fatalf("proxy func returned error: %v", err) } if proxyURL.String() != "http://127.0.0.1:10808" { t.Fatalf("proxy URL = %q, want explicit proxy", proxyURL.String()) } } func TestArtifactAccumulatorGeneratedImageRevision(t *testing.T) { acc := newArtifactAccumulator() raw := map[string]interface{}{ "type": "conversation-update", "conversation_id": "conv-img", "message": map[string]interface{}{ "id": "msg-img", "author": map[string]interface{}{"role": "assistant"}, "metadata": map[string]interface{}{ "image_gen_task_id": "task-1", }, "content": map[string]interface{}{ "content_type": "image_asset_pointer", "parts": []interface{}{ map[string]interface{}{ "asset_pointer": "sediment://file-img123", "metadata": map[string]interface{}{ "dalle": map[string]interface{}{ "gen_id": "gen-1", "slot_index": float64(1), }, }, }, }, }, }, } events := acc.ObserveRaw(raw, "conv-img") finalEvents := acc.Finalize() if len(events) < 2 { t.Fatalf("events = %#v, want pending and artifact", events) } if events[0].Event != StreamEventArtifactPending { t.Fatalf("first event = %#v, want artifact_pending", events[0]) } last := events[len(events)-1] if last.Event != StreamEventArtifact || last.Kind != "generated_image" || last.FileID != "file-img123" || last.GenID != "gen-1" || last.SlotIndex != 1 || last.Revision != 1 { t.Fatalf("image artifact event = %#v", last) } if len(finalEvents) != 1 || finalEvents[0].Event != StreamEventArtifactSlotFinal { t.Fatalf("finalEvents = %#v, want slot final", finalEvents) } } func TestHandlerSeparatesAnalysisAndFinalChannels(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"v":{"conversation_id":"conv-thinking","message":{"id":"msg-thinking","author":{"role":"assistant"},"channel":"analysis","content":{"content_type":"text","parts":["think"]},"metadata":{"message_type":"next"},"recipient":"all"}}}`, `data: {"v":{"conversation_id":"conv-thinking","message":{"id":"msg-final","author":{"role":"assistant"},"channel":"final","content":{"content_type":"text","parts":["answer"]},"metadata":{"message_type":"next"},"recipient":"all","end_turn":true}}}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) result := HandlerDetailedWithOptions(c, response, nil, nil, "request-id", chatGPTRequestForTest(), false, "auto", HandlerDetailedOptions{}) if result.Text != "answer" { t.Fatalf("text = %q, want final answer only", result.Text) } if result.ThinkingText != "think" { t.Fatalf("thinking = %q, want analysis text", result.ThinkingText) } if len(result.Sentinel) == 0 || result.Sentinel[0]["event"] != "thinking" { t.Fatalf("sentinel = %#v, want thinking event", result.Sentinel) } if result.ParentMessageID != "msg-final" { t.Fatalf("parent message id = %q, want msg-final", result.ParentMessageID) } } func TestHandlerAppliesBatchPatch(t *testing.T) { gin.SetMode(gin.TestMode) body := strings.Join([]string{ `data: {"p":"/conversation_id","o":"replace","v":"conv-batch"}`, `data: {"p":"/message/id","o":"replace","v":"msg-batch"}`, `data: {"p":"/message/author/role","o":"replace","v":"assistant"}`, `data: {"p":"/message/content/parts/0","o":"append","v":"hello "}`, `data: {"p":"","o":"patch","v":[{"p":"/message/content/parts/0","o":"append","v":"world"}]}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} writer := httptest.NewRecorder() c, _ := gin.CreateTestContext(writer) result := HandlerDetailedWithOptions(c, response, nil, nil, "request-id", chatGPTRequestForTest(), false, "auto", HandlerDetailedOptions{}) if result.Text != "hello world" { t.Fatalf("text = %q, want hello world", result.Text) } if result.ConversationID != "conv-batch" { t.Fatalf("conversation id = %q, want conv-batch", result.ConversationID) } if result.ParentMessageID != "msg-batch" { t.Fatalf("parent message id = %q, want msg-batch", result.ParentMessageID) } } func TestHandlerTTSParsesPatchStream(t *testing.T) { body := strings.Join([]string{ `data: {"p":"/conversation_id","o":"replace","v":"conv-tts"}`, `data: {"p":"/message/id","o":"replace","v":"msg-tts"}`, `data: {"p":"/message/author/role","o":"replace","v":"assistant"}`, `data: {"p":"/message/content/parts/0","o":"append","v":"hello tts"}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} msgID, convID := HandlerTTS(response, "hello tts") if msgID != "msg-tts" || convID != "conv-tts" { t.Fatalf("msgID=%q convID=%q, want msg-tts conv-tts", msgID, convID) } } func TestHandlerTTSFallsBackToAssistantMessageID(t *testing.T) { body := strings.Join([]string{ `data: {"conversation_id":"conv-tts","message":{"id":"msg-tts","author":{"role":"assistant"},"content":{"content_type":"text","parts":["different text"]},"metadata":{"message_type":"next"},"recipient":"all"}}`, `data: [DONE]`, ``, }, "\n") response := &http.Response{Body: io.NopCloser(strings.NewReader(body))} msgID, convID := HandlerTTS(response, "requested text") if msgID != "msg-tts" || convID != "conv-tts" { t.Fatalf("msgID=%q convID=%q, want fallback assistant message", msgID, convID) } } func chatGPTRequestForTest() chatgpt.ChatGPTRequest { return chatgpt.NewChatGPTRequest() } func parseSSEChunks(t *testing.T, output string) []map[string]interface{} { t.Helper() var chunks []map[string]interface{} for _, line := range strings.Split(output, "\n") { line = strings.TrimSpace(line) if !strings.HasPrefix(line, "data: ") || line == "data: [DONE]" { continue } var chunk map[string]interface{} if err := json.Unmarshal([]byte(strings.TrimPrefix(line, "data: ")), &chunk); err != nil { t.Fatalf("invalid chunk %q: %v", line, err) } chunks = append(chunks, chunk) } return chunks }