From ef51690b5a71501a2ce76296d6979f4ea09f6779 Mon Sep 17 00:00:00 2001 From: Achord Chan Date: Sun, 20 Sep 2026 10:46:09 +0800 Subject: [PATCH 01/78] fix(antigravity): preserve string const constraints in tool schemas --- .../pkg/antigravity/schema_cleaner.go | 21 +++++++++ .../pkg/antigravity/schema_const_test.go | 44 +++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 backend/internal/pkg/antigravity/schema_const_test.go diff --git a/backend/internal/pkg/antigravity/schema_cleaner.go b/backend/internal/pkg/antigravity/schema_cleaner.go index 0ee746aa3..3cbd0fcae 100644 --- a/backend/internal/pkg/antigravity/schema_cleaner.go +++ b/backend/internal/pkg/antigravity/schema_cleaner.go @@ -119,6 +119,27 @@ func cleanJSONSchemaRecursive(value any) any { // 0. [NEW] 合并 allOf mergeAllOf(schemaMap) + // Gemini's enum representation supports strings. Preserve string constants + // before the allowlist removes const, including schemas with no explicit type. + if constant, ok := schemaMap["const"].(string); ok { + values := []any{constant} + if existing, ok := schemaMap["enum"].([]any); ok { + // Both constraints apply: retain their intersection, not a wider enum. + values = []any{} + for _, value := range existing { + if text, ok := value.(string); ok && text == constant { + values = append(values, constant) + break + } + } + } + schemaMap["enum"] = values + if _, exists := schemaMap["type"]; !exists { + schemaMap["type"] = "string" + } + delete(schemaMap, "const") + } + // 1. [CRITICAL] 深度递归处理子项 if props, ok := schemaMap["properties"].(map[string]any); ok { for _, v := range props { diff --git a/backend/internal/pkg/antigravity/schema_const_test.go b/backend/internal/pkg/antigravity/schema_const_test.go new file mode 100644 index 000000000..ac83a4026 --- /dev/null +++ b/backend/internal/pkg/antigravity/schema_const_test.go @@ -0,0 +1,44 @@ +package antigravity + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestBuildToolsPreservesStringConst(t *testing.T) { + for _, tc := range []struct { + name string + schema string + want string + }{ + {"typed", `{"type":"string","const":"browser"}`, `{"type":"string","enum":["browser"]}`}, + {"inferred", `{"const":"browser"}`, `{"type":"string","enum":["browser"]}`}, + {"empty", `{"type":"string","const":""}`, `{"type":"string","enum":[""]}`}, + {"existing enum", `{"type":"string","const":"browser","enum":["browser","shell"]}`, `{"type":"string","enum":["browser"]}`}, + {"conflicting enum", `{"type":"string","const":"browser","enum":["shell"]}`, `{"type":"string","enum":[]}`}, + {"ordinary enum", `{"type":"string","enum":["browser","shell"]}`, `{"type":"string","enum":["browser","shell"]}`}, + {"array items", `{"type":"array","items":{"type":"string","const":"browser"}}`, `{"type":"array","items":{"type":"string","enum":["browser"]}}`}, + {"nested property named const", `{"type":"object","properties":{"const":{"const":"browser"}}}`, `{"type":"object","properties":{"const":{"type":"string","enum":["browser"]}}}`}, + {"schema metadata", `{"type":"string","const":"browser","$schema":"https://json-schema.org/draft/2020-12/schema","description":"Tool kind"}`, `{"type":"string","enum":["browser"],"description":"Tool kind"}`}, + } { + t.Run(tc.name, func(t *testing.T) { + var property map[string]any + require.NoError(t, json.Unmarshal([]byte(tc.schema), &property)) + tools := buildTools([]ClaudeTool{{ + Name: "dispatch", + InputSchema: map[string]any{ + "type": "object", + "properties": map[string]any{"action": property}, + }, + }}) + require.Len(t, tools, 1) + require.Len(t, tools[0].FunctionDeclarations, 1) + properties := tools[0].FunctionDeclarations[0].Parameters["properties"].(map[string]any) + got, err := json.Marshal(properties["action"]) + require.NoError(t, err) + require.JSONEq(t, tc.want, string(got)) + }) + } +} From b13200d7a18486797f77b7833f4b2e35c00be7ff Mon Sep 17 00:00:00 2001 From: wucm667 Date: Sat, 19 Sep 2026 21:26:18 +0800 Subject: [PATCH 02/78] fix(gateway): forward received chat stream usage --- .../anthropic_chat_stream_usage_test.go | 111 ++++++++++++++++++ .../gateway_forward_as_chat_completions.go | 22 +++- ...ateway_forward_as_chat_completions_test.go | 4 +- ...enai_gateway_anthropic_native_pump_test.go | 4 +- ...teway_chat_completions_anthropic_native.go | 9 +- 5 files changed, 138 insertions(+), 12 deletions(-) create mode 100644 backend/internal/service/anthropic_chat_stream_usage_test.go diff --git a/backend/internal/service/anthropic_chat_stream_usage_test.go b/backend/internal/service/anthropic_chat_stream_usage_test.go new file mode 100644 index 000000000..99ce4d278 --- /dev/null +++ b/backend/internal/service/anthropic_chat_stream_usage_test.go @@ -0,0 +1,111 @@ +//go:build unit + +package service + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/Wei-Shaw/sub2api/internal/pkg/apicompat" + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/require" +) + +func TestAnthropicChatStreamAuthoritativeUsage(t *testing.T) { + gin.SetMode(gin.TestMode) + for _, adapter := range []string{"anthropic", "native"} { + for _, options := range []string{"", `,"stream_options":{"include_usage":true}`, `,"stream_options":{"include_usage":false}`} { + for _, tc := range []struct { + name, start, delta string + want bool + prompt, output, cached, created int + }{ + {"cache", `,"usage":{"input_tokens":308,"cache_read_input_tokens":241,"cache_creation_input_tokens":17}`, `,"usage":{"output_tokens":49}`, true, 566, 49, 241, 17}, + {"zero_start", `,"usage":{"input_tokens":0,"output_tokens":0}`, "", true, 0, 0, 0, 0}, + {"zero_delta", "", `,"usage":{"input_tokens":0,"output_tokens":0}`, true, 0, 0, 0, 0}, + {"absent", "", "", false, 0, 0, 0, 0}, + {"null", `,"usage":null`, `,"usage":null`, false, 0, 0, 0, 0}, + } { + for _, terminal := range []string{"stop", "eof"} { + t.Run(adapter+"/"+options+"/"+tc.name+"/"+terminal, func(t *testing.T) { + sse := fmt.Sprintf("event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\",\"type\":\"message\",\"role\":\"assistant\",\"model\":\"k3\",\"content\":[]%s}}\n\nevent: message_delta\ndata: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}%s}\n\n", tc.start, tc.delta) + if terminal == "stop" { + // Repeated terminal events and EOF must not duplicate usage. + sse += "event: message_stop\ndata: {\"type\":\"message_stop\"}\n\nevent: message_stop\ndata: {\"type\":\"message_stop\"}\n\n" + } + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader(sse)), + }} + body := []byte(`{"model":"k3","messages":[{"role":"user","content":"hi"}],"stream":true` + options + "}") + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader(string(body))) + if adapter == "native" { + svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream} + result, err := svc.forwardChatCompletionsViaNativeAnthropic(context.Background(), c, nativeAnthropicTestAccount(), body, "") + require.NoError(t, err) + require.Equal(t, tc.output, result.Usage.OutputTokens) + require.Equal(t, tc.cached, result.Usage.CacheReadInputTokens) + require.Equal(t, tc.created, result.Usage.CacheCreationInputTokens) + } else { + svc := &GatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream} + account := &Account{ID: 1, Platform: PlatformAnthropic, Type: AccountTypeAPIKey, Credentials: map[string]any{"api_key": "sk-test"}} + result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, nil) + require.NoError(t, err) + require.Equal(t, tc.output, result.Usage.OutputTokens) + require.Equal(t, tc.cached, result.Usage.CacheReadInputTokens) + require.Equal(t, tc.created, result.Usage.CacheCreationInputTokens) + } + count, done := 0, 0 + var last *apicompat.ChatCompletionsChunk + for _, line := range strings.Split(rec.Body.String(), "\n") { + if !strings.HasPrefix(line, "data: ") { + continue + } + payload := strings.TrimPrefix(line, "data: ") + if payload == "[DONE]" { + done++ + if tc.want { + require.NotNil(t, last) + require.NotNil(t, last.Usage) + } + continue + } + require.Zero(t, done, "chunk after DONE") + var chunk apicompat.ChatCompletionsChunk + require.NoError(t, json.Unmarshal([]byte(payload), &chunk)) + last = &chunk + if chunk.Usage == nil { + continue + } + count++ + require.Empty(t, chunk.Choices) + require.Equal(t, tc.prompt, chunk.Usage.PromptTokens) + require.Equal(t, tc.output, chunk.Usage.CompletionTokens) + require.Equal(t, tc.prompt+tc.output, chunk.Usage.TotalTokens) + if tc.cached > 0 { + require.NotNil(t, chunk.Usage.PromptTokensDetails) + require.Equal(t, tc.cached, chunk.Usage.PromptTokensDetails.CachedTokens) + require.Equal(t, tc.created, chunk.Usage.PromptTokensDetails.CacheCreationTokens) + } + } + require.Equal(t, 1, done) + if tc.want { + require.Equal(t, 1, count) + } else { + require.Zero(t, count) + } + }) + } + } + } + } +} diff --git a/backend/internal/service/gateway_forward_as_chat_completions.go b/backend/internal/service/gateway_forward_as_chat_completions.go index 32d38cf40..687e31cc6 100644 --- a/backend/internal/service/gateway_forward_as_chat_completions.go +++ b/backend/internal/service/gateway_forward_as_chat_completions.go @@ -42,7 +42,6 @@ func (s *GatewayService) ForwardAsChatCompletions( } originalModel := ccReq.Model clientStream := ccReq.Stream - includeUsage := ccReq.StreamOptions != nil && ccReq.StreamOptions.IncludeUsage // 2. Convert CC → Responses → Anthropic (chained conversion) responsesReq, err := apicompat.ChatCompletionsToResponses(&ccReq) @@ -185,7 +184,7 @@ func (s *GatewayService) ForwardAsChatCompletions( var result *ForwardResult var handleErr error if clientStream { - result, handleErr = s.handleCCStreamingFromAnthropic(resp, c, originalModel, mappedModel, reasoningEffort, startTime, includeUsage) + result, handleErr = s.handleCCStreamingFromAnthropic(resp, c, originalModel, mappedModel, reasoningEffort, startTime) } else { result, handleErr = s.handleCCBufferedFromAnthropic(resp, c, originalModel, mappedModel, reasoningEffort, startTime) } @@ -358,7 +357,6 @@ func (s *GatewayService) handleCCStreamingFromAnthropic( mappedModel string, reasoningEffort *string, startTime time.Time, - includeUsage bool, ) (*ForwardResult, error) { requestID := resp.Header.Get("x-request-id") @@ -376,7 +374,6 @@ func (s *GatewayService) handleCCStreamingFromAnthropic( anthState.Model = originalModel ccState := apicompat.NewResponsesEventToChatState() ccState.Model = originalModel - ccState.IncludeUsage = includeUsage var usage ClaudeUsage var firstTokenMs *int @@ -467,6 +464,10 @@ func (s *GatewayService) handleCCStreamingFromAnthropic( continue } + // Forward received usage regardless of the client stream_options. + // The intermediate Responses converter synthesizes usage even when absent. + ccState.IncludeUsage = ccState.IncludeUsage || anthropicChatStreamHasUsage(&event, payload) + if processAnthropicEvent(&event) { return resultWithUsage(), nil } @@ -512,3 +513,16 @@ func writeGatewayCCError(c *gin.Context, statusCode int, errType, message string }, }) } + +// anthropicChatStreamHasUsage distinguishes an explicit zero-valued usage object +// from omitted/null usage before the Anthropic→Responses conversion loses that distinction. +func anthropicChatStreamHasUsage(event *apicompat.AnthropicStreamEvent, payload string) bool { + switch event.Type { + case "message_start": + return event.Message != nil && gjson.Get(payload, "message.usage").IsObject() + case "message_delta": + return event.Usage != nil + default: + return false + } +} diff --git a/backend/internal/service/gateway_forward_as_chat_completions_test.go b/backend/internal/service/gateway_forward_as_chat_completions_test.go index f7cb42c69..163b20449 100644 --- a/backend/internal/service/gateway_forward_as_chat_completions_test.go +++ b/backend/internal/service/gateway_forward_as_chat_completions_test.go @@ -199,7 +199,7 @@ func TestHandleCCStreamingFromAnthropic_CompactSSEFormat(t *testing.T) { } svc := &GatewayService{} - result, err := svc.handleCCStreamingFromAnthropic(resp, c, "k3", "k3", nil, time.Now(), true) + result, err := svc.handleCCStreamingFromAnthropic(resp, c, "k3", "k3", nil, time.Now()) require.NoError(t, err) require.NotNil(t, result) require.Equal(t, 21, result.Usage.InputTokens) @@ -236,7 +236,7 @@ func TestHandleCCStreamingFromAnthropic_PreservesMessageStartCacheUsageAndReason } svc := &GatewayService{} - result, err := svc.handleCCStreamingFromAnthropic(resp, c, "gpt-5", "claude-sonnet-4.5", &reasoningEffort, time.Now(), true) + result, err := svc.handleCCStreamingFromAnthropic(resp, c, "gpt-5", "claude-sonnet-4.5", &reasoningEffort, time.Now()) require.NoError(t, err) require.NotNil(t, result) require.Equal(t, 20, result.Usage.InputTokens) diff --git a/backend/internal/service/openai_gateway_anthropic_native_pump_test.go b/backend/internal/service/openai_gateway_anthropic_native_pump_test.go index 83bbd4eef..e5b9d0dd8 100644 --- a/backend/internal/service/openai_gateway_anthropic_native_pump_test.go +++ b/backend/internal/service/openai_gateway_anthropic_native_pump_test.go @@ -140,7 +140,7 @@ func TestCCStreamingFromNativeAnthropic_HangTimesOut(t *testing.T) { resp, pr, pw := newHangingUpstreamResponse() start := time.Now() - res, err := svc.handleCCStreamingFromNativeAnthropic(resp, c, "glm-4.7", "glm-4.7", "glm-4.7", nil, start, true) + res, err := svc.handleCCStreamingFromNativeAnthropic(resp, c, "glm-4.7", "glm-4.7", "glm-4.7", nil, start) _ = pw.Close() _ = pr.Close() @@ -220,7 +220,7 @@ func TestCCStreamingFromNativeAnthropic_HappyPathStillConverts(t *testing.T) { }() defer func() { _ = pr.Close() }() - res, err := svc.handleCCStreamingFromNativeAnthropic(resp, c, "glm-4.7", "glm-4.7", "glm-4.7", nil, time.Now(), true) + res, err := svc.handleCCStreamingFromNativeAnthropic(resp, c, "glm-4.7", "glm-4.7", "glm-4.7", nil, time.Now()) if err != nil { t.Fatalf("unexpected error: %v", err) } diff --git a/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go b/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go index bc865af81..40b00c81e 100644 --- a/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go +++ b/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go @@ -54,7 +54,6 @@ func (s *OpenAIGatewayService) forwardChatCompletionsViaNativeAnthropic( return nil, fmt.Errorf("missing model in request") } clientStream := ccReq.Stream - includeUsage := ccReq.StreamOptions != nil && ccReq.StreamOptions.IncludeUsage // 2. Convert CC → Responses → Anthropic (chained conversion) responsesReq, err := apicompat.ChatCompletionsToResponses(&ccReq) @@ -136,7 +135,7 @@ func (s *OpenAIGatewayService) forwardChatCompletionsViaNativeAnthropic( reasoningEffort = ApplyThinkingEnabledFallback(reasoningEffort, body, billingModel) if clientStream { - return s.handleCCStreamingFromNativeAnthropic(resp, c, originalModel, billingModel, upstreamModel, reasoningEffort, startTime, includeUsage) + return s.handleCCStreamingFromNativeAnthropic(resp, c, originalModel, billingModel, upstreamModel, reasoningEffort, startTime) } return s.handleCCBufferedFromNativeAnthropic(resp, c, originalModel, billingModel, upstreamModel, reasoningEffort, startTime) } @@ -303,7 +302,6 @@ func (s *OpenAIGatewayService) handleCCStreamingFromNativeAnthropic( upstreamModel string, reasoningEffort *string, startTime time.Time, - includeUsage bool, ) (*OpenAIForwardResult, error) { requestID := resp.Header.Get("x-request-id") @@ -320,7 +318,6 @@ func (s *OpenAIGatewayService) handleCCStreamingFromNativeAnthropic( anthState.Model = originalModel ccState := apicompat.NewResponsesEventToChatState() ccState.Model = originalModel - ccState.IncludeUsage = includeUsage var usage ClaudeUsage var firstTokenMs *int @@ -461,6 +458,10 @@ func (s *OpenAIGatewayService) handleCCStreamingFromNativeAnthropic( continue } + // Forward received usage regardless of the client stream_options. + // The intermediate Responses converter synthesizes usage even when absent. + ccState.IncludeUsage = ccState.IncludeUsage || anthropicChatStreamHasUsage(&event, payload) + if processAnthropicEvent(&event) { return resultWithUsage(), nil } From ece1ecccca842585c99bc738bb93115bd1455e3e Mon Sep 17 00:00:00 2001 From: wucm667 Date: Mon, 21 Sep 2026 12:36:32 +0800 Subject: [PATCH 03/78] fix(gateway): normalize streamed Anthropic usage --- .../anthropic_chat_stream_usage_test.go | 27 +++-- .../gateway_forward_as_chat_completions.go | 5 + .../service/gateway_forward_as_responses.go | 67 ++++++++--- .../gateway_forward_as_responses_test.go | 110 ++++++++++++++++++ ...teway_chat_completions_anthropic_native.go | 5 + .../service/openai_gateway_cn_fixes_test.go | 80 +++++++++++++ ...enai_gateway_responses_anthropic_native.go | 6 + 7 files changed, 277 insertions(+), 23 deletions(-) diff --git a/backend/internal/service/anthropic_chat_stream_usage_test.go b/backend/internal/service/anthropic_chat_stream_usage_test.go index 99ce4d278..ff1558b07 100644 --- a/backend/internal/service/anthropic_chat_stream_usage_test.go +++ b/backend/internal/service/anthropic_chat_stream_usage_test.go @@ -22,19 +22,28 @@ func TestAnthropicChatStreamAuthoritativeUsage(t *testing.T) { for _, adapter := range []string{"anthropic", "native"} { for _, options := range []string{"", `,"stream_options":{"include_usage":true}`, `,"stream_options":{"include_usage":false}`} { for _, tc := range []struct { - name, start, delta string - want bool - prompt, output, cached, created int + name, start, delta, repeatDelta string + want bool + prompt, billableInput, totalInput, output, cached, created int }{ - {"cache", `,"usage":{"input_tokens":308,"cache_read_input_tokens":241,"cache_creation_input_tokens":17}`, `,"usage":{"output_tokens":49}`, true, 566, 49, 241, 17}, - {"zero_start", `,"usage":{"input_tokens":0,"output_tokens":0}`, "", true, 0, 0, 0, 0}, - {"zero_delta", "", `,"usage":{"input_tokens":0,"output_tokens":0}`, true, 0, 0, 0, 0}, - {"absent", "", "", false, 0, 0, 0, 0}, - {"null", `,"usage":null`, `,"usage":null`, false, 0, 0, 0, 0}, + {"cache", `,"usage":{"input_tokens":308,"cache_read_input_tokens":241,"cache_creation_input_tokens":17}`, `,"usage":{"output_tokens":49}`, "", true, 566, 308, 566, 49, 241, 17}, + {"deepseek_partial_cache", "", `,"usage":{"input_tokens":1200,"output_tokens":30,"prompt_cache_hit_tokens":800,"prompt_cache_miss_tokens":400}`, "", true, 1200, 400, 1200, 30, 800, 0}, + {"openai_partial_cache", "", `,"usage":{"prompt_tokens":1200,"output_tokens":30,"prompt_tokens_details":{"cached_tokens":800}}`, "", true, 1200, 400, 1200, 30, 800, 0}, + {"repeated_cache_bucket", `,"usage":{"input_tokens":1200}`, `,"usage":{"output_tokens":30,"cache_read_input_tokens":800}`, `,"usage":{"output_tokens":30,"cache_read_input_tokens":800}`, true, 1200, 400, 1200, 30, 800, 0}, + {"kimi_full_cache", `,"usage":{"input_tokens":173306}`, `,"usage":{"output_tokens":49,"cache_read_input_tokens":173306}`, "", true, 173306, 0, 173306, 49, 173306, 0}, + {"kimi_full_cache_cached_tokens", `,"usage":{"input_tokens":173306}`, `,"usage":{"output_tokens":49,"cached_tokens":173306}`, "", true, 173306, 0, 173306, 49, 173306, 0}, + {"kimi_full_cache_prompt_details", `,"usage":{"input_tokens":173306}`, `,"usage":{"output_tokens":49,"prompt_tokens_details":{"cached_tokens":173306}}`, "", true, 173306, 0, 173306, 49, 173306, 0}, + {"zero_start", `,"usage":{"input_tokens":0,"output_tokens":0}`, "", "", true, 0, 0, 0, 0, 0, 0}, + {"zero_delta", "", `,"usage":{"input_tokens":0,"output_tokens":0}`, "", true, 0, 0, 0, 0, 0, 0}, + {"absent", "", "", "", false, 0, 0, 0, 0, 0, 0}, + {"null", `,"usage":null`, `,"usage":null`, "", false, 0, 0, 0, 0, 0, 0}, } { for _, terminal := range []string{"stop", "eof"} { t.Run(adapter+"/"+options+"/"+tc.name+"/"+terminal, func(t *testing.T) { sse := fmt.Sprintf("event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\",\"type\":\"message\",\"role\":\"assistant\",\"model\":\"k3\",\"content\":[]%s}}\n\nevent: message_delta\ndata: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}%s}\n\n", tc.start, tc.delta) + if tc.repeatDelta != "" { + sse += fmt.Sprintf("event: message_delta\ndata: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}%s}\n\n", tc.repeatDelta) + } if terminal == "stop" { // Repeated terminal events and EOF must not duplicate usage. sse += "event: message_stop\ndata: {\"type\":\"message_stop\"}\n\nevent: message_stop\ndata: {\"type\":\"message_stop\"}\n\n" @@ -52,6 +61,7 @@ func TestAnthropicChatStreamAuthoritativeUsage(t *testing.T) { svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream} result, err := svc.forwardChatCompletionsViaNativeAnthropic(context.Background(), c, nativeAnthropicTestAccount(), body, "") require.NoError(t, err) + require.Equal(t, tc.totalInput, result.Usage.InputTokens) require.Equal(t, tc.output, result.Usage.OutputTokens) require.Equal(t, tc.cached, result.Usage.CacheReadInputTokens) require.Equal(t, tc.created, result.Usage.CacheCreationInputTokens) @@ -60,6 +70,7 @@ func TestAnthropicChatStreamAuthoritativeUsage(t *testing.T) { account := &Account{ID: 1, Platform: PlatformAnthropic, Type: AccountTypeAPIKey, Credentials: map[string]any{"api_key": "sk-test"}} result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, nil) require.NoError(t, err) + require.Equal(t, tc.billableInput, result.Usage.InputTokens) require.Equal(t, tc.output, result.Usage.OutputTokens) require.Equal(t, tc.cached, result.Usage.CacheReadInputTokens) require.Equal(t, tc.created, result.Usage.CacheCreationInputTokens) diff --git a/backend/internal/service/gateway_forward_as_chat_completions.go b/backend/internal/service/gateway_forward_as_chat_completions.go index 687e31cc6..0936af2b4 100644 --- a/backend/internal/service/gateway_forward_as_chat_completions.go +++ b/backend/internal/service/gateway_forward_as_chat_completions.go @@ -430,6 +430,11 @@ func (s *GatewayService) handleCCStreamingFromAnthropic( mergeAnthropicUsage(&usage, event.Message.Usage) } + // Keep the outward Responses/Chat usage on the same normalized buckets used + // for billing, including converter handlers that consume event usage. + syncAnthropicResponsesUsage(anthState, usage) + normalizeAnthropicEventUsageForResponses(event, usage) + // Chain: Anthropic event → Responses events → CC chunks responsesEvents := apicompat.AnthropicEventToResponsesEvents(event, anthState) for _, resEvt := range responsesEvents { diff --git a/backend/internal/service/gateway_forward_as_responses.go b/backend/internal/service/gateway_forward_as_responses.go index 1b97a7396..6903e659a 100644 --- a/backend/internal/service/gateway_forward_as_responses.go +++ b/backend/internal/service/gateway_forward_as_responses.go @@ -279,22 +279,22 @@ func mergeAnthropicUsage(dst *ClaudeUsage, src apicompat.AnthropicUsage) { return } + cacheReadTokens := src.CacheReadInputTokens + if cacheReadTokens == 0 && src.CachedTokens > 0 { + cacheReadTokens = src.CachedTokens + } + if cacheReadTokens == 0 && src.PromptTokensDetails != nil && src.PromptTokensDetails.CachedTokens > 0 { + cacheReadTokens = src.PromptTokensDetails.CachedTokens + } + if cacheReadTokens == 0 && src.PromptCacheHitTokens != nil { + cacheReadTokens = max(*src.PromptCacheHitTokens, 0) + } + // Some Anthropic-compatible providers retain OpenAI-style prompt/cache // fields. Prefer those authoritative totals or hit/miss buckets over the // overloaded input_tokens field. This covers Kimi's changing stream // semantics as well as GLM/DeepSeek cache aliases. if src.PromptTokens > 0 || src.PromptCacheHitTokens != nil || src.PromptCacheMissTokens != nil { - cacheReadTokens := src.CacheReadInputTokens - if cacheReadTokens == 0 && src.CachedTokens > 0 { - cacheReadTokens = src.CachedTokens - } - if cacheReadTokens == 0 && src.PromptTokensDetails != nil && src.PromptTokensDetails.CachedTokens > 0 { - cacheReadTokens = src.PromptTokensDetails.CachedTokens - } - if cacheReadTokens == 0 && src.PromptCacheHitTokens != nil { - cacheReadTokens = max(*src.PromptCacheHitTokens, 0) - } - if src.PromptCacheMissTokens != nil { dst.InputTokens = max(*src.PromptCacheMissTokens, 0) } else { @@ -303,13 +303,21 @@ func mergeAnthropicUsage(dst *ClaudeUsage, src apicompat.AnthropicUsage) { dst.CacheReadInputTokens = cacheReadTokens dst.CacheCreationInputTokens = src.CacheCreationInputTokens } else { + previousCacheReadTokens := dst.CacheReadInputTokens + previousCacheCreationTokens := dst.CacheCreationInputTokens if src.InputTokens > 0 { dst.InputTokens = src.InputTokens } - if src.CacheReadInputTokens > 0 { - dst.CacheReadInputTokens = src.CacheReadInputTokens - } else if src.CachedTokens > 0 { - dst.CacheReadInputTokens = src.CachedTokens + if src.InputTokens == 0 && dst.InputTokens > 0 && (cacheReadTokens > 0 || src.CacheCreationInputTokens > 0) { + // Some compatible streams put the total prompt count in message_start, + // then provide only cumulative cache buckets in message_delta. Subtract + // only newly observed buckets so repeated deltas are idempotent. + newCacheReadTokens := max(cacheReadTokens-previousCacheReadTokens, 0) + newCacheCreationTokens := max(src.CacheCreationInputTokens-previousCacheCreationTokens, 0) + dst.InputTokens = max(dst.InputTokens-newCacheReadTokens-newCacheCreationTokens, 0) + } + if cacheReadTokens > 0 { + dst.CacheReadInputTokens = cacheReadTokens } if src.CacheCreationInputTokens > 0 { dst.CacheCreationInputTokens = src.CacheCreationInputTokens @@ -320,6 +328,29 @@ func mergeAnthropicUsage(dst *ClaudeUsage, src apicompat.AnthropicUsage) { } } +func syncAnthropicResponsesUsage(state *apicompat.AnthropicEventToResponsesState, usage ClaudeUsage) { + state.InputTokens = usage.InputTokens + state.OutputTokens = usage.OutputTokens + state.CacheReadInputTokens = usage.CacheReadInputTokens + state.CacheCreationInputTokens = usage.CacheCreationInputTokens +} + +func normalizeAnthropicEventUsageForResponses(event *apicompat.AnthropicStreamEvent, usage ClaudeUsage) { + normalize := func(dst *apicompat.AnthropicUsage) { + if dst == nil { + return + } + dst.InputTokens = usage.InputTokens + dst.OutputTokens = usage.OutputTokens + dst.CacheReadInputTokens = usage.CacheReadInputTokens + dst.CacheCreationInputTokens = usage.CacheCreationInputTokens + } + normalize(event.Usage) + if event.Message != nil { + normalize(&event.Message.Usage) + } +} + // parseAnthropicSSEField parses an SSE field line in the form "field:value" or "field: value". // According to the SSE spec (https://html.spec.whatwg.org/multipage/server-sent-events.html#event-stream-interpretation), // the space after the colon is optional. This function handles both formats. @@ -545,6 +576,12 @@ func (s *GatewayService) handleResponsesStreamingResponse( mergeAnthropicUsage(&usage, event.Message.Usage) } + // Keep the terminal Responses usage aligned with the normalized billing + // buckets. Normalize the converter input too, so message handlers cannot + // restore the provider's overlapping raw input total. + syncAnthropicResponsesUsage(state, usage) + normalizeAnthropicEventUsageForResponses(event, usage) + // Convert to Responses events events := apicompat.AnthropicEventToResponsesEvents(event, state) for _, evt := range events { diff --git a/backend/internal/service/gateway_forward_as_responses_test.go b/backend/internal/service/gateway_forward_as_responses_test.go index 972c6bc1a..353eeb884 100644 --- a/backend/internal/service/gateway_forward_as_responses_test.go +++ b/backend/internal/service/gateway_forward_as_responses_test.go @@ -307,6 +307,116 @@ func TestHandleResponsesStreamingResponse_PreservesMessageStartCacheUsage(t *tes require.Contains(t, rec.Body.String(), `response.completed`) } +func TestHandleResponsesStreamingResponse_NormalizesTerminalUsage(t *testing.T) { + gin.SetMode(gin.TestMode) + + tests := []struct { + name string + startUsage string + deltaUsage string + wantInput int + wantOutput int + wantCached int + wantCacheCreation int + }{ + { + name: "cache_read_input_tokens full cache", + startUsage: `"input_tokens":173306`, + deltaUsage: `"output_tokens":8,"cache_read_input_tokens":173306`, + wantOutput: 8, + wantCached: 173306, + }, + { + name: "cached_tokens full cache", + startUsage: `"input_tokens":173306`, + deltaUsage: `"output_tokens":8,"cached_tokens":173306`, + wantOutput: 8, + wantCached: 173306, + }, + { + name: "prompt_tokens_details full cache", + startUsage: `"input_tokens":173306`, + deltaUsage: `"output_tokens":8,"prompt_tokens_details":{"cached_tokens":173306}`, + wantOutput: 8, + wantCached: 173306, + }, + { + name: "DeepSeek input total with hit and miss buckets", + startUsage: `"input_tokens":0`, + deltaUsage: `"input_tokens":1200,"output_tokens":30,"prompt_cache_hit_tokens":800,"prompt_cache_miss_tokens":400`, + wantInput: 400, + wantOutput: 30, + wantCached: 800, + }, + { + name: "OpenAI prompt total with cached details", + startUsage: `"input_tokens":0`, + deltaUsage: `"prompt_tokens":1200,"output_tokens":30,"prompt_tokens_details":{"cached_tokens":800}`, + wantInput: 400, + wantOutput: 30, + wantCached: 800, + }, + { + name: "cache creation only without prompt total or miss bucket", + startUsage: `"input_tokens":1200`, + deltaUsage: `"output_tokens":30,"cache_creation_input_tokens":800`, + wantInput: 400, + wantOutput: 30, + wantCacheCreation: 800, + }, + } + + for _, tt := range tests { + for _, terminal := range []string{"message_stop", "eof"} { + t.Run(tt.name+"/"+terminal, func(t *testing.T) { + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + lines := []string{ + `event: message_start`, + `data: {"type":"message_start","message":{"id":"msg_usage","type":"message","role":"assistant","content":[],"model":"k3","stop_reason":"","usage":{` + tt.startUsage + `}}}`, + ``, + `event: message_delta`, + `data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{` + tt.deltaUsage + `}}`, + ``, + } + if terminal == "message_stop" { + lines = append(lines, `event: message_stop`, `data: {"type":"message_stop"}`, ``) + } + resp := &http.Response{Body: io.NopCloser(strings.NewReader(strings.Join(lines, "\n")))} + + result, err := (&GatewayService{}).handleResponsesStreamingResponse(resp, c, "k3", "k3", nil, time.Now(), apicompat.ResponsesClientToolMapping{}) + require.NoError(t, err) + require.Equal(t, tt.wantInput, result.Usage.InputTokens) + require.Equal(t, tt.wantCached, result.Usage.CacheReadInputTokens) + require.Equal(t, tt.wantCacheCreation, result.Usage.CacheCreationInputTokens) + + var completed apicompat.ResponsesStreamEvent + for _, line := range strings.Split(rec.Body.String(), "\n") { + if !strings.HasPrefix(line, "data: ") { + continue + } + var event apicompat.ResponsesStreamEvent + require.NoError(t, json.Unmarshal([]byte(strings.TrimPrefix(line, "data: ")), &event)) + if event.Type == "response.completed" { + completed = event + } + } + require.NotNil(t, completed.Response) + require.NotNil(t, completed.Response.Usage) + require.Equal(t, tt.wantInput+tt.wantCached+tt.wantCacheCreation, completed.Response.Usage.InputTokens) + require.Equal(t, tt.wantOutput, completed.Response.Usage.OutputTokens) + require.Equal(t, completed.Response.Usage.InputTokens+tt.wantOutput, completed.Response.Usage.TotalTokens) + require.Equal(t, tt.wantCacheCreation, completed.Response.Usage.CacheCreationInputTokens) + if tt.wantCached == 0 { + require.Nil(t, completed.Response.Usage.InputTokensDetails) + } else { + require.Equal(t, tt.wantCached, completed.Response.Usage.InputTokensDetails.CachedTokens) + } + }) + } + } +} + func TestParseAnthropicSSEField(t *testing.T) { t.Parallel() diff --git a/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go b/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go index 40b00c81e..03c8b8e72 100644 --- a/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go +++ b/backend/internal/service/openai_gateway_chat_completions_anthropic_native.go @@ -407,6 +407,11 @@ func (s *OpenAIGatewayService) handleCCStreamingFromNativeAnthropic( mergeAnthropicUsage(&usage, event.Message.Usage) } + // Keep the outward Responses/Chat usage on the same normalized buckets used + // for billing, including converter handlers that consume event usage. + syncAnthropicResponsesUsage(anthState, usage) + normalizeAnthropicEventUsageForResponses(event, usage) + // 客户端已断开:跳过转换与写出,继续读上游直到流结束(usage 完整、 // 连接及时归还),不再提前 return。 if clientDisconnected { diff --git a/backend/internal/service/openai_gateway_cn_fixes_test.go b/backend/internal/service/openai_gateway_cn_fixes_test.go index 0dfb7c35b..470c789c7 100644 --- a/backend/internal/service/openai_gateway_cn_fixes_test.go +++ b/backend/internal/service/openai_gateway_cn_fixes_test.go @@ -11,9 +11,12 @@ package service import ( "context" + "encoding/json" "errors" + "io" "net/http" "net/http/httptest" + "strings" "testing" "time" @@ -109,6 +112,83 @@ func TestResponsesStreamingFromNativeAnthropic_ClientDisconnectDrainsUsage(t *te "output_tokens 必须来自排水读到的末尾 message_delta(断开即弃时会是 1)") } +func TestResponsesStreamingFromNativeAnthropic_NormalizesTerminalUsage(t *testing.T) { + gin.SetMode(gin.TestMode) + + tests := []struct { + name string + startUsage string + deltaUsage string + repeatDelta bool + wantInput int + wantOutput int + wantCached int + wantCacheCreation int + }{ + {name: "full cache", startUsage: `"input_tokens":1200`, deltaUsage: `"output_tokens":30,"cache_read_input_tokens":1200`, wantOutput: 30, wantCached: 1200}, + {name: "partial cache", startUsage: `"input_tokens":1200`, deltaUsage: `"output_tokens":30,"cache_read_input_tokens":800`, wantInput: 400, wantOutput: 30, wantCached: 800}, + {name: "cache creation", startUsage: `"input_tokens":1200`, deltaUsage: `"output_tokens":30,"cache_creation_input_tokens":800`, wantInput: 400, wantOutput: 30, wantCacheCreation: 800}, + {name: "repeated cumulative cache bucket", startUsage: `"input_tokens":1200`, deltaUsage: `"output_tokens":30,"cache_read_input_tokens":800`, repeatDelta: true, wantInput: 400, wantOutput: 30, wantCached: 800}, + } + + for _, tt := range tests { + for _, terminal := range []string{"message_stop", "eof"} { + t.Run(tt.name+"/"+terminal, func(t *testing.T) { + lines := []string{ + `event: message_start`, + `data: {"type":"message_start","message":{"id":"msg_usage","type":"message","role":"assistant","content":[],"model":"k3","stop_reason":"","usage":{` + tt.startUsage + `}}}`, + ``, + `event: message_delta`, + `data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{` + tt.deltaUsage + `}}`, + ``, + } + if tt.repeatDelta { + lines = append(lines, `event: message_delta`, `data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{`+tt.deltaUsage+`}}`, ``) + } + if terminal == "message_stop" { + lines = append(lines, `event: message_stop`, `data: {"type":"message_stop"}`, ``) + } + + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/", nil) + resp := &http.Response{StatusCode: http.StatusOK, Header: http.Header{}, Body: io.NopCloser(strings.NewReader(strings.Join(lines, "\n")))} + + result, err := (&OpenAIGatewayService{}).handleResponsesStreamingFromNativeAnthropic( + resp, c, "k3", "k3", "k3", nil, time.Now(), apicompat.ResponsesClientToolMapping{}) + require.NoError(t, err) + require.Equal(t, tt.wantInput+tt.wantCached+tt.wantCacheCreation, result.Usage.InputTokens) + require.Equal(t, tt.wantOutput, result.Usage.OutputTokens) + require.Equal(t, tt.wantCached, result.Usage.CacheReadInputTokens) + require.Equal(t, tt.wantCacheCreation, result.Usage.CacheCreationInputTokens) + + var completed apicompat.ResponsesStreamEvent + for _, line := range strings.Split(rec.Body.String(), "\n") { + if !strings.HasPrefix(line, "data: ") { + continue + } + var event apicompat.ResponsesStreamEvent + require.NoError(t, json.Unmarshal([]byte(strings.TrimPrefix(line, "data: ")), &event)) + if event.Type == "response.completed" { + completed = event + } + } + require.NotNil(t, completed.Response) + require.NotNil(t, completed.Response.Usage) + require.Equal(t, tt.wantInput+tt.wantCached+tt.wantCacheCreation, completed.Response.Usage.InputTokens) + require.Equal(t, tt.wantOutput, completed.Response.Usage.OutputTokens) + require.Equal(t, completed.Response.Usage.InputTokens+tt.wantOutput, completed.Response.Usage.TotalTokens) + require.Equal(t, tt.wantCacheCreation, completed.Response.Usage.CacheCreationInputTokens) + if tt.wantCached == 0 { + require.Nil(t, completed.Response.Usage.InputTokensDetails) + } else { + require.Equal(t, tt.wantCached, completed.Response.Usage.InputTokensDetails.CachedTokens) + } + }) + } + } +} + func TestHandle403_CNProviderHTMLBodySkipsAccountPenalty(t *testing.T) { for _, platform := range []string{PlatformKimi, PlatformZhipu, PlatformDeepseek, PlatformMiniMax} { repo := &rateLimitAccountRepoStub{} diff --git a/backend/internal/service/openai_gateway_responses_anthropic_native.go b/backend/internal/service/openai_gateway_responses_anthropic_native.go index fa2e190b9..c5f8ef688 100644 --- a/backend/internal/service/openai_gateway_responses_anthropic_native.go +++ b/backend/internal/service/openai_gateway_responses_anthropic_native.go @@ -397,6 +397,12 @@ func (s *OpenAIGatewayService) handleResponsesStreamingFromNativeAnthropic( mergeAnthropicUsage(&usage, event.Message.Usage) } + // Keep terminal Responses usage aligned with the normalized billing + // buckets. Normalize converter input too so raw overlapping totals cannot + // overwrite the state when message_start/message_delta handlers run. + syncAnthropicResponsesUsage(state, usage) + normalizeAnthropicEventUsageForResponses(event, usage) + events := apicompat.AnthropicEventToResponsesEvents(event, state) if clientDisconnected { return From f867e3e2ece67781d23cff532c64fcd8bbc90c73 Mon Sep 17 00:00:00 2001 From: Achord Chan Date: Tue, 22 Sep 2026 18:34:36 +0800 Subject: [PATCH 04/78] test(antigravity): check schema properties type assertion --- backend/internal/pkg/antigravity/schema_const_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/backend/internal/pkg/antigravity/schema_const_test.go b/backend/internal/pkg/antigravity/schema_const_test.go index ac83a4026..05551fbd8 100644 --- a/backend/internal/pkg/antigravity/schema_const_test.go +++ b/backend/internal/pkg/antigravity/schema_const_test.go @@ -35,7 +35,8 @@ func TestBuildToolsPreservesStringConst(t *testing.T) { }}) require.Len(t, tools, 1) require.Len(t, tools[0].FunctionDeclarations, 1) - properties := tools[0].FunctionDeclarations[0].Parameters["properties"].(map[string]any) + properties, ok := tools[0].FunctionDeclarations[0].Parameters["properties"].(map[string]any) + require.True(t, ok, "tool parameters must contain a properties object") got, err := json.Marshal(properties["action"]) require.NoError(t, err) require.JSONEq(t, tc.want, string(got)) From 1af3269adb1bceea4e8d362ed5e55438f081ece8 Mon Sep 17 00:00:00 2001 From: "chihao.ou" Date: Tue, 22 Sep 2026 21:43:13 +0800 Subject: [PATCH 05/78] fix(groups): clean up modal listeners and pending searches --- .../admin/group/GroupRPMOverridesModal.vue | 6 +++- .../admin/group/GroupRateMultipliersModal.vue | 6 +++- .../__tests__/GroupModal.cleanup.spec.ts | 30 +++++++++++++++++++ 3 files changed, 40 insertions(+), 2 deletions(-) create mode 100644 frontend/src/components/admin/group/__tests__/GroupModal.cleanup.spec.ts diff --git a/frontend/src/components/admin/group/GroupRPMOverridesModal.vue b/frontend/src/components/admin/group/GroupRPMOverridesModal.vue index 025054162..439878014 100644 --- a/frontend/src/components/admin/group/GroupRPMOverridesModal.vue +++ b/frontend/src/components/admin/group/GroupRPMOverridesModal.vue @@ -206,7 +206,7 @@