fix(gateway): 记录国产厂商原生 Anthropic 直通路径的 reasoning_effort

kimi/zhipu/deepseek 平台分组的 /v1/messages 请求在账号 api_protocol=anthropic
时走零转换直通路径(forwardAnthropicViaNativeAnthropicEndpoint)。该路径的
两个结果构造函数(非流式/流式)从不设置 ReasoningEffort,导致客户端显式
携带 output_config.effort 或启用 thinking 时,usage_log.reasoning_effort
仍恒为 NULL,用量页推理强度列长期显示 '-'。

修复:转发前从入站 body 提取推理强度——优先 Claude 协议 output_config.effort
(NormalizeClaudeOutputEffort),缺失且 thinking 已启用时经
ApplyThinkingEnabledFallback 按 passback-required 白名单兜底为 high,与
Anthropic 网关(gateway_handler)的记录语义对齐;两个结果构造函数补上
ReasoningEffort 字段透传。

测试:新增 4 个单元测试覆盖显式 effort(非流式/流式)、thinking 兜底、
无字段保持 nil 的守卫场景。
This commit is contained in:
clearmann.liu
2026-08-20 08:16:17 -04:00
parent 2bc139ab52
commit 2074fe3ba8
2 changed files with 164 additions and 11 deletions
@@ -58,6 +58,15 @@ func (s *OpenAIGatewayService) forwardAnthropicViaNativeAnthropicEndpoint(
body = rewritten
}
// 记录客户端请求的推理强度:优先 Claude 协议的 output_config.effort;
// 缺失且 thinking 已启用时,按国产 passback-required 模型兜底为 high
// (对齐 Anthropic 网关 gateway_handler 的记录语义,避免该路径长期落 NULL)。
reasoningEffort := ApplyThinkingEnabledFallback(
NormalizeClaudeOutputEffort(gjson.GetBytes(body, "output_config.effort").String()),
body,
billingModel,
)
// 与 Anthropic 平台 passthrough 相同的 pre-filter:剥离空文本块与上游
// 无法接受的 web-search 历史块(GLM/Kimi/DeepSeek 对 server_tool_use 400)。
body = StripEmptyTextBlocks(body)
@@ -104,9 +113,9 @@ func (s *OpenAIGatewayService) forwardAnthropicViaNativeAnthropicEndpoint(
}
if clientStream {
return s.handleNativeAnthropicStreamingResponse(ctx, resp, c, account, originalModel, billingModel, upstreamModel, startTime)
return s.handleNativeAnthropicStreamingResponse(ctx, resp, c, account, originalModel, billingModel, upstreamModel, reasoningEffort, startTime)
}
return s.handleNativeAnthropicBufferedResponse(ctx, resp, c, account, originalModel, billingModel, upstreamModel, startTime)
return s.handleNativeAnthropicBufferedResponse(ctx, resp, c, account, originalModel, billingModel, upstreamModel, reasoningEffort, startTime)
}
// nativeAnthropicTargetURL 组装国产供应商原生 Anthropic messages 端点。
@@ -194,6 +203,7 @@ func (s *OpenAIGatewayService) handleNativeAnthropicBufferedResponse(
originalModel string,
billingModel string,
upstreamModel string,
reasoningEffort *string,
startTime time.Time,
) (*OpenAIForwardResult, error) {
if s.rateLimitService != nil {
@@ -239,6 +249,7 @@ func (s *OpenAIGatewayService) handleNativeAnthropicBufferedResponse(
UpstreamModel: upstreamModel,
UpstreamEndpoint: "/v1/messages",
Stream: false,
ReasoningEffort: reasoningEffort,
Duration: time.Since(startTime),
}, nil
}
@@ -254,6 +265,7 @@ func (s *OpenAIGatewayService) handleNativeAnthropicStreamingResponse(
originalModel string,
billingModel string,
upstreamModel string,
reasoningEffort *string,
startTime time.Time,
) (*OpenAIForwardResult, error) {
observer := upstreamResponseModelObserverFromContext(c)
@@ -382,28 +394,28 @@ func (s *OpenAIGatewayService) handleNativeAnthropicStreamingResponse(
flusher.Flush()
}
if !sawTerminalEvent {
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream usage incomplete: missing terminal event")
}
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime), nil
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime), nil
}
if ev.err != nil {
if sawTerminalEvent {
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime), nil
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime), nil
}
if clientDisconnected {
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream usage incomplete after disconnect: %w", ev.err)
}
if errors.Is(ev.err, context.Canceled) || errors.Is(ev.err, context.DeadlineExceeded) {
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream usage incomplete: %w", ev.err)
}
if errors.Is(ev.err, bufio.ErrTooLong) {
logger.LegacyPrintf("service.gateway", "[CN Anthropic 直通] SSE line too long: account=%d max_size=%d error=%v", account.ID, maxLineSize, ev.err)
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime), ev.err
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime), ev.err
}
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream read error: %w", ev.err)
}
@@ -451,14 +463,14 @@ func (s *OpenAIGatewayService) handleNativeAnthropicStreamingResponse(
continue
}
if clientDisconnected {
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream usage incomplete after timeout")
}
logger.LegacyPrintf("service.gateway", "[CN Anthropic 直通] Stream data interval timeout: account=%d model=%s interval=%s", account.ID, upstreamModel, streamInterval)
if s.rateLimitService != nil {
s.rateLimitService.HandleStreamTimeout(ctx, account, upstreamModel)
}
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, startTime),
return s.nativeAnthropicStreamResult(c, resp, usage, firstTokenMs, clientDisconnected, originalModel, billingModel, upstreamModel, reasoningEffort, startTime),
fmt.Errorf("stream data interval timeout")
case <-keepaliveCh:
@@ -496,6 +508,7 @@ func (s *OpenAIGatewayService) nativeAnthropicStreamResult(
originalModel string,
billingModel string,
upstreamModel string,
reasoningEffort *string,
startTime time.Time,
) *OpenAIForwardResult {
if usage == nil {
@@ -509,6 +522,7 @@ func (s *OpenAIGatewayService) nativeAnthropicStreamResult(
UpstreamModel: upstreamModel,
UpstreamEndpoint: "/v1/messages",
Stream: true,
ReasoningEffort: reasoningEffort,
Duration: time.Since(startTime),
FirstTokenMs: firstTokenMs,
ClientDisconnect: clientDisconnect,
@@ -0,0 +1,139 @@
//go:build unit
package service
import (
"context"
"io"
"net/http"
"strings"
"testing"
"github.com/gin-gonic/gin"
"github.com/stretchr/testify/require"
)
// 国产供应商原生 Anthropic 直通路径(api_protocol=anthropic)的 reasoning_effort 记录。
// 回归背景:该路径此前从不提取 Claude 协议的 output_config.effort,也不做
// thinking-enabled 兜底,导致 kimi/zhipu/deepseek 平台分组的 /v1/messages 请求
// usage_log.reasoning_effort 恒为 NULL。
func nativeAnthropicTestAccount() *Account {
return &Account{
ID: 702,
Name: "kimi-native",
Platform: PlatformKimi,
Type: AccountTypeAPIKey,
Concurrency: 1,
Credentials: map[string]any{
"api_key": "sk-test",
"api_protocol": APIProtocolAnthropic,
"api_base_urls": map[string]any{
APIProtocolAnthropic: "http://anthropic.example",
},
},
}
}
func nativeAnthropicBufferedResponse() *http.Response {
return &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(
`{"id":"msg_1","type":"message","role":"assistant","model":"k3",` +
`"content":[{"type":"text","text":"pong"}],"stop_reason":"end_turn",` +
`"usage":{"input_tokens":93,"output_tokens":16}}`,
)),
}
}
func nativeAnthropicStreamResponse() *http.Response {
sse := `event: message_start
data: {"type":"message_start","message":{"id":"msg_1","type":"message","role":"assistant","model":"k3","content":[],"stop_reason":null,"usage":{"input_tokens":93,"output_tokens":1}}}
event: content_block_start
data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"pong"}}
event: content_block_stop
data: {"type":"content_block_stop","index":0}
event: message_delta
data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":16}}
event: message_stop
data: {"type":"message_stop"}
`
return &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"text/event-stream"}},
Body: io.NopCloser(strings.NewReader(sse)),
}
}
func TestNativeAnthropicPassthroughRecordsOutputConfigEffort(t *testing.T) {
gin.SetMode(gin.TestMode)
body := []byte(`{"model":"k3","max_tokens":32,"stream":false,` +
`"output_config":{"effort":"low"},` +
`"messages":[{"role":"user","content":"hi"}]}`)
upstream := &httpUpstreamRecorder{resp: nativeAnthropicBufferedResponse()}
svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream}
result, err := svc.ForwardAsAnthropic(context.Background(),
adaptiveProtocolTestContext("/v1/messages", body), nativeAnthropicTestAccount(), body, "", "")
require.NoError(t, err)
require.NotNil(t, result)
require.NotNil(t, result.ReasoningEffort)
require.Equal(t, "low", *result.ReasoningEffort)
}
func TestNativeAnthropicPassthroughThinkingEnabledFallback(t *testing.T) {
gin.SetMode(gin.TestMode)
// 未显式传 effort,但 thinking 已启用:k3 属于 passback-required 白名单,应兜底记为 high。
body := []byte(`{"model":"k3","max_tokens":32,"stream":false,` +
`"thinking":{"type":"enabled","budget_tokens":1024},` +
`"messages":[{"role":"user","content":"hi"}]}`)
upstream := &httpUpstreamRecorder{resp: nativeAnthropicBufferedResponse()}
svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream}
result, err := svc.ForwardAsAnthropic(context.Background(),
adaptiveProtocolTestContext("/v1/messages", body), nativeAnthropicTestAccount(), body, "", "")
require.NoError(t, err)
require.NotNil(t, result)
require.NotNil(t, result.ReasoningEffort)
require.Equal(t, "high", *result.ReasoningEffort)
}
func TestNativeAnthropicPassthroughStreamRecordsEffort(t *testing.T) {
gin.SetMode(gin.TestMode)
body := []byte(`{"model":"k3","max_tokens":32,"stream":true,` +
`"output_config":{"effort":"max"},` +
`"messages":[{"role":"user","content":"hi"}]}`)
upstream := &httpUpstreamRecorder{resp: nativeAnthropicStreamResponse()}
svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream}
result, err := svc.ForwardAsAnthropic(context.Background(),
adaptiveProtocolTestContext("/v1/messages", body), nativeAnthropicTestAccount(), body, "", "")
require.NoError(t, err)
require.NotNil(t, result)
require.NotNil(t, result.ReasoningEffort)
require.Equal(t, "max", *result.ReasoningEffort)
}
func TestNativeAnthropicPassthroughNoEffortStaysNil(t *testing.T) {
gin.SetMode(gin.TestMode)
// 既无 output_config.effort 也未启用 thinking:保持 nil,不做语义注入。
body := []byte(`{"model":"k3","max_tokens":32,"stream":false,` +
`"messages":[{"role":"user","content":"hi"}]}`)
upstream := &httpUpstreamRecorder{resp: nativeAnthropicBufferedResponse()}
svc := &OpenAIGatewayService{cfg: rawChatCompletionsTestConfig(), httpUpstream: upstream}
result, err := svc.ForwardAsAnthropic(context.Background(),
adaptiveProtocolTestContext("/v1/messages", body), nativeAnthropicTestAccount(), body, "", "")
require.NoError(t, err)
require.NotNil(t, result)
require.Nil(t, result.ReasoningEffort)
}