fix(openai): adapt service tier observation to upstream constraints

- cc_pipeline: observe Chat Completions chunks/bodies as untyped payloads
  (empty event type) so the upstream-echoed service_tier is trusted, matching
  the upstream constraint that only terminal events and untyped bodies report
  the actual processing tier.
- tests: add model field to terminal SSE frames (observation only triggers on
  model-bearing frames) and assert response.created tier echo is ignored.
This commit is contained in:
alfadb
2026-08-24 11:52:48 +08:00
parent c0c3e1cb47
commit e457f0fa22
4 changed files with 22 additions and 12 deletions
@@ -757,17 +757,24 @@ func TestUpstreamResponseModelObserver_ObservesServiceTier(t *testing.T) {
t.Parallel()
observer := &upstreamResponseModelObserver{}
// 上游约束:非终止且有类型的事件(response.created)回显的是请求档位而非
// 实际处理档位,忽略。
observer.ObserveOpenAI([]byte(`{"type":"response.created","response":{"model":"gpt-5.5","service_tier":"flex"}}`), "response.created")
require.Equal(t, "flex", observer.ServiceTier())
require.Empty(t, observer.ServiceTier())
// terminal 声明优先。
// terminal 声明(带 model 帧)优先。
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"model":"gpt-5.5","service_tier":"default"}}`), "response.completed")
require.Equal(t, "default", observer.ServiceTier())
// Chat Completions 顶层 service_tier 同样可观察。
// Chat Completions 顶层 service_tier 按 untyped payload 观察(无 type 字段)。
ccObserver := &upstreamResponseModelObserver{}
ccObserver.ObserveOpenAI([]byte(`{"id":"chatcmpl-1","model":"gpt-5.5","service_tier":"priority","choices":[]}`), "chat.completion")
ccObserver.ObserveOpenAI([]byte(`{"id":"chatcmpl-1","model":"gpt-5.5","service_tier":"priority","choices":[]}`), "")
require.Equal(t, "priority", ccObserver.ServiceTier())
// 无 model 的帧不触发 tier 观察(上游约束:tier 声明必带 model)。
modelFree := &upstreamResponseModelObserver{}
modelFree.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"service_tier":"default"}}`), "response.completed")
require.Empty(t, modelFree.ServiceTier())
}
func TestResolvedOpenAIUpstreamServiceTier(t *testing.T) {
@@ -779,7 +786,7 @@ func TestResolvedOpenAIUpstreamServiceTier(t *testing.T) {
gin.SetMode(gin.TestMode)
c, _ := gin.CreateTestContext(nil)
observer := beginUpstreamResponseModelObservation(c)
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"service_tier":"default"}}`), "response.completed")
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"model":"gpt-5.5","service_tier":"default"}}`), "response.completed")
got := resolvedOpenAIUpstreamServiceTier(c, priority)
require.NotNil(t, got)
@@ -800,7 +807,7 @@ func TestResolvedOpenAIUpstreamServiceTier(t *testing.T) {
gin.SetMode(gin.TestMode)
c, _ := gin.CreateTestContext(nil)
observer := beginUpstreamResponseModelObservation(c)
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"service_tier":"fast"}}`), "response.completed")
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"model":"gpt-5.5","service_tier":"fast"}}`), "response.completed")
got := resolvedOpenAIUpstreamServiceTier(c, nil)
require.NotNil(t, got)
@@ -819,7 +826,7 @@ func TestResolvedOpenAIUpstreamServiceTier(t *testing.T) {
t.Run("local observer wins without gin context", func(t *testing.T) {
observer := &upstreamResponseModelObserver{}
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"service_tier":"default"}}`), "response.completed")
observer.ObserveOpenAI([]byte(`{"type":"response.completed","response":{"model":"gpt-5.5","service_tier":"default"}}`), "response.completed")
got := resolvedOpenAIUpstreamServiceTierFromObserver(observer, priority)
require.NotNil(t, got)
@@ -274,8 +274,10 @@ func (s *OpenAIGatewayService) scanCCStream(
break
}
// 观察上游 CC chunk 回显的 model / service_tier(计费以回显为准)。
// CC chunk 无 type 字段,按 untyped payload 观察(上游约束:只有终止
// 事件与无类型 body 报告实际处理档位)。
if observer := upstreamResponseModelObserverFromContext(c); observer != nil {
observer.ObserveOpenAI([]byte(payload), "chat.completion.chunk")
observer.ObserveOpenAI([]byte(payload), "")
}
if u := extractCCStreamUsage(payload); u != nil {
@@ -337,8 +339,9 @@ func (s *OpenAIGatewayService) readCCUpstreamJSONResponse(
return nil, OpenAIUsage{}, fmt.Errorf("parse chat completions response: %w", err)
}
// 观察上游 CC JSON 回显的 model / service_tier(计费以回显为准)。
// CC JSON 无 type 字段,按 untyped payload 观察(上游约束)。
if observer := upstreamResponseModelObserverFromContext(c); observer != nil {
observer.ObserveOpenAI(respBody, "chat.completion")
observer.ObserveOpenAI(respBody, "")
}
usage := OpenAIUsage{}
@@ -50,7 +50,7 @@ func TestForwardOpenAIWSV2_UpstreamDefaultServiceTierWinsOverRequest(t *testing.
captureConn := &openAIWSCaptureConn{
events: [][]byte{
[]byte(`{"type":"response.completed","response":{"id":"resp_tier_v2","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`),
[]byte(`{"type":"response.completed","response":{"id":"resp_tier_v2","model":"gpt-5.5","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`),
},
}
captureDialer := &openAIWSCaptureDialer{conn: captureConn}
@@ -52,7 +52,7 @@ func TestProxyOpenAIWSHTTPBridgeTurn_UpstreamDefaultServiceTierWinsOverRequest(t
// policy。本测试只覆盖局部 observer:canonical 请求 priority 被上游
// response.completed service_tier=default 覆盖。
sse := strings.Join([]string{
`data: {"type":"response.completed","response":{"id":"resp_tier","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`,
`data: {"type":"response.completed","response":{"id":"resp_tier","model":"gpt-5.5","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`,
``,
}, "\n")
upstream := &httpUpstreamRecorder{resp: &http.Response{
@@ -89,7 +89,7 @@ func TestProxyOpenAIWSHTTPBridgeTurn_UpstreamDefaultWinsOverFastAlias(t *testing
// 客户端别名 fast 同样被上游回显的 default 覆盖:局部 observer 的
// ServiceTier() 是唯一计费依据,绝不回退到请求侧 fast。
sse := strings.Join([]string{
`data: {"type":"response.completed","response":{"id":"resp_tier2","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`,
`data: {"type":"response.completed","response":{"id":"resp_tier2","model":"gpt-5.5","status":"completed","service_tier":"default","usage":{"input_tokens":1,"output_tokens":1}}}`,
``,
}, "\n")
upstream := &httpUpstreamRecorder{resp: &http.Response{