Revert "修复 Grok WebSearch SSE action 兼容"

This reverts commit 726de30101.
This commit is contained in:
IanShaw
2026-08-20 06:36:42 -07:00
parent 5ae254f771
commit ab9cb69e7e
2 changed files with 1 additions and 70 deletions
@@ -439,9 +439,7 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context.
sendErrorEvent("stream_read_error")
return resultWithUsage(), fmt.Errorf("stream read error: %w", scanErr), true
}
var pendingGrokWebSearchCompleted string
var processSSELine func(line string, queueDrained bool)
processSSELine = func(line string, queueDrained bool) {
processSSELine := func(line string, queueDrained bool) {
if streamEarlyErr != nil {
return
}
@@ -450,26 +448,6 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context.
dataBytes := []byte(data)
eventTypeRaw := gjson.GetBytes(dataBytes, "type").String()
eventType := strings.TrimSpace(eventTypeRaw)
// Grok Build's xAI decoder requires action on the completed search
// event, while OpenAI-compatible upstreams often provide it only on
// the following output_item.done event. Hold that one event until the
// matching item arrives, then replay it with action injected.
if account != nil && account.Platform == PlatformGrok && account.Type == AccountTypeAPIKey &&
eventType == "response.web_search_call.completed" &&
!gjson.GetBytes(dataBytes, "action").Exists() {
pendingGrokWebSearchCompleted = line
return
}
if pendingGrokWebSearchCompleted != "" && eventType == "response.output_item.done" {
if adapted, adaptedOK := adaptGrokWebSearchCompletedAction(
[]byte(strings.TrimSpace(strings.TrimPrefix(pendingGrokWebSearchCompleted, "data:"))), dataBytes,
); adaptedOK {
pendingGrokWebSearchCompleted = "data: " + string(adapted)
}
pending := pendingGrokWebSearchCompleted
pendingGrokWebSearchCompleted = ""
processSSELine(pending, queueDrained)
}
observer.ObserveOpenAI(dataBytes, eventTypeRaw)
// 初始上游 data 的 type 只解析一次:原始值保持终止事件的精确匹配,规范化值供后续分支复用。
if openAIStreamEventIsTerminalWithType(data, eventTypeRaw) {
@@ -662,11 +640,6 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context.
// A blank line dispatches a guarded event from the attempt-local stage.
if stageFirstOutput && line == "" {
if pendingGrokWebSearchCompleted != "" {
pending := pendingGrokWebSearchCompleted
pendingGrokWebSearchCompleted = ""
processSSELine(pending, queueDrained)
}
if !clientDisconnected {
if _, err := writePendingString("\n"); err != nil {
handlePendingWriteError(err)
@@ -681,11 +654,6 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context.
// or queue-drain flush must never split an open SSE event.
shouldFlush := false
if line == "" {
if pendingGrokWebSearchCompleted != "" {
pending := pendingGrokWebSearchCompleted
pendingGrokWebSearchCompleted = ""
processSSELine(pending, queueDrained)
}
shouldFlush = eventShouldFlush || (queueDrained && clientOutputStarted)
eventShouldFlush = false
}
@@ -1205,32 +1173,6 @@ func (s *OpenAIGatewayService) bindHTTPResponseAccount(ctx context.Context, c *g
logOpenAIWSBindResponseAccountWarn(groupID, account.ID, responseID, store.BindResponseAccount(ctx, groupID, responseID, account.ID, ttl))
}
// adaptGrokWebSearchCompletedAction fills the xAI-specific action field that
// Grok Build expects on response.web_search_call.completed. OpenAI-compatible
// upstreams commonly emit the action only on the later output_item.done event.
func adaptGrokWebSearchCompletedAction(completed, outputItem []byte) ([]byte, bool) {
if strings.TrimSpace(gjson.GetBytes(completed, "type").String()) != "response.web_search_call.completed" {
return completed, false
}
itemID := strings.TrimSpace(gjson.GetBytes(completed, "item_id").String())
item := gjson.GetBytes(outputItem, "item")
if itemID == "" || !item.Exists() || strings.TrimSpace(item.Get("type").String()) != "web_search_call" {
return completed, false
}
if itemID != strings.TrimSpace(item.Get("id").String()) {
return completed, false
}
action := item.Get("action")
if !action.Exists() || strings.TrimSpace(action.Raw) == "" || action.Raw == "null" {
return completed, false
}
updated, err := sjson.SetRawBytes(completed, "action", []byte(action.Raw))
if err != nil {
return completed, false
}
return updated, true
}
func openAIUsageFromGJSON(value gjson.Result) (OpenAIUsage, bool) {
if !value.Exists() || !value.IsObject() {
return OpenAIUsage{}, false
@@ -3463,17 +3463,6 @@ func TestExtractOpenAIUsageFromJSONBytes_AcceptsResponseAndChatUsageShapes(t *te
require.Equal(t, 500, usage.OutputTokens)
}
func TestAdaptGrokWebSearchCompletedAction(t *testing.T) {
completed := []byte(`{"type":"response.web_search_call.completed","item_id":"call_1"}`)
done := []byte(`{"type":"response.output_item.done","item":{"type":"web_search_call","id":"call_1","action":{"type":"search","query":"latest news"}}}`)
got, ok := adaptGrokWebSearchCompletedAction(completed, done)
require.True(t, ok)
require.Equal(t, "search", gjson.GetBytes(got, "action.type").String())
require.Equal(t, "latest news", gjson.GetBytes(got, "action.query").String())
_, ok = adaptGrokWebSearchCompletedAction(completed, []byte(`{"type":"response.output_item.done","item":{"type":"web_search_call","id":"other"}}`))
require.False(t, ok)
}
func TestExtractOpenAIUsageFromJSONBytes_IncludesGrokReasoningTokens(t *testing.T) {
usage, ok := extractOpenAIUsageFromJSONBytes([]byte(`{"usage":{"prompt_tokens":32,"completion_tokens":9,"total_tokens":135,"completion_tokens_details":{"reasoning_tokens":94}}}`))
require.True(t, ok)