fix(grok): Codex compact 适配与链式视频 content 代理

- #4223/#4554: Grok 模拟 /responses/compact,调度允许 Grok 账号
- #4494/#4626: 支持链式中继的受保护视频 /content 下载
This commit is contained in:
abbzbb
2026-07-20 16:12:56 +08:00
parent 2e5973b165
commit 2ae61f3da0
8 changed files with 377 additions and 5 deletions
@@ -409,7 +409,7 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
if len(failedAccountIDs) == 0 {
if errors.Is(err, service.ErrNoAvailableCompactAccounts) {
markOpsRoutingCapacityLimitedIfNoAvailable(c, err)
h.handleStreamingAwareError(c, http.StatusServiceUnavailable, "compact_not_supported", "No available OpenAI accounts support /responses/compact", streamStarted)
h.handleStreamingAwareError(c, http.StatusServiceUnavailable, "compact_not_supported", "No available accounts support /responses/compact", streamStarted)
return
}
cls := classifyOpenAICompatibleNoAccountErrorFromGin(c, h.gatewayService, apiKey, reqModel, reqModel)
@@ -170,6 +170,56 @@ func TestOpenAIGatewayService_SelectAccountWithScheduler_CompactFallsBackToUnkno
require.Equal(t, int64(71021), selection.Account.ID, "unknown account should be picked when no supported account available")
}
// TestOpenAIGatewayService_SelectAccountWithScheduler_CompactAllowsGrok verifies
// that OpenAI-compatible compact routing does not reject Grok accounts.
func TestOpenAIGatewayService_SelectAccountWithScheduler_CompactAllowsGrok(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
ctx := context.Background()
groupID := int64(91004)
accounts := []Account{
{
ID: 71030,
Platform: PlatformGrok,
Type: AccountTypeOAuth,
Status: StatusActive,
Schedulable: true,
Concurrency: 1,
Priority: 0,
Credentials: map[string]any{
"model_mapping": map[string]any{"grok-4.5": "grok-4.5"},
},
},
}
cfg := &config.Config{}
cfg.Gateway.Scheduling.LoadBatchEnabled = false
svc := &OpenAIGatewayService{
accountRepo: schedulerTestOpenAIAccountRepo{accounts: accounts},
cache: &schedulerTestGatewayCache{},
cfg: cfg,
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
}
selection, _, err := svc.SelectAccountWithSchedulerForCapability(
ctx,
&groupID,
"",
"",
"grok-4.5",
nil,
OpenAIUpstreamTransportAny,
OpenAIEndpointCapabilityChatCompletions,
true,
false,
true,
PlatformGrok,
)
require.NoError(t, err)
require.NotNil(t, selection)
require.NotNil(t, selection.Account)
require.Equal(t, int64(71030), selection.Account.ID)
}
// TestOpenAICompactSupportTier 验证 tier 分类逻辑。
func TestOpenAICompactSupportTier(t *testing.T) {
tests := []struct {
@@ -179,6 +229,7 @@ func TestOpenAICompactSupportTier(t *testing.T) {
}{
{name: "nil", account: nil, want: 0},
{name: "non openai", account: &Account{Platform: PlatformAnthropic}, want: 0},
{name: "grok", account: &Account{Platform: PlatformGrok}, want: 2},
{name: "openai unknown", account: &Account{Platform: PlatformOpenAI, Extra: map[string]any{}}, want: 1},
{name: "openai supported", account: &Account{Platform: PlatformOpenAI, Extra: map[string]any{"openai_compact_supported": true}}, want: 2},
{name: "openai unsupported", account: &Account{Platform: PlatformOpenAI, Extra: map[string]any{"openai_compact_supported": false}}, want: 0},
@@ -56,6 +56,15 @@ func (s *OpenAIGatewayService) forwardGrokResponses(
if err != nil {
return nil, err
}
// OpenAI /responses/compact is not a native xAI endpoint. Convert it into a
// normal Grok Responses turn that asks for a structured summary, then map the
// reply back to an OpenAI compaction item on the way out.
if isOpenAIResponsesCompactPath(c) {
patchedBody, err = buildGrokCompactRequestBody(patchedBody)
if err != nil {
return nil, err
}
}
// Derive the identity from the request xAI will actually see. This makes
// Codex Responses Lite additional_tools part of the stable tool prefix.
cacheIdentity := resolveGrokCacheIdentity(c, patchedBody, "", upstreamModel)
@@ -391,6 +400,10 @@ func patchGrokResponsesBody(body []byte, upstreamModel string) ([]byte, error) {
if err != nil {
return nil, err
}
out, err = convertOpenAICompactInputsForGrok(out)
if err != nil {
return nil, err
}
out, err = sanitizeGrokResponsesInput(out)
if err != nil {
return nil, err
@@ -0,0 +1,233 @@
package service
import (
"encoding/json"
"fmt"
"strings"
"github.com/google/uuid"
)
// This is build_compaction_prompt(None, false) from grok-build. Grok does not
// expose an OpenAI-compatible /responses/compact endpoint, so compacting is a
// normal Responses turn whose final user item asks the model to summarize.
const grokCompactSummaryPrompt = `Your task is to produce a faithful, concise summary of the conversation so far so that a successor assistant can continue the work seamlessly after the earlier turns are discarded. The successor will see the user's original query plus this summary. Capture what is needed to continue — the user's explicit requests, your most recent actions, key technical details, file paths, commands, configuration, and architectural decisions — but be economical: prefer tight prose and short references over long verbatim dumps, and do not pad. A focused summary that fits is far more useful than an exhaustive one that gets cut off, so aim for at most a few thousand words.
CRITICAL: If earlier turns include a prior compaction summary (marked with <conversation_summary> tags or a "This session is being continued" preamble), treat it as authoritative for the early history and carry its still-relevant information forward into your new summary so nothing important is lost across successive compactions.
Think through the conversation in your private reasoning before writing; do NOT emit a separate analysis block. Output the final summary inside a single <summary>...</summary> block, organized into the following numbered sections. Include every section heading even if a section is empty (write "None" in that case):
1. Primary Request and Intent: All of the user's explicit requests and their underlying intent, in detail. Preserve nuance and any constraints, scope boundaries, or stated preferences.
2. Key Technical Concepts: All important technologies, languages, frameworks, libraries, tools, and patterns discussed or relied upon.
3. Files and Code Sections: Every file examined, created, or modified. For each, give the full path, why it matters, and the relevant code — include full snippets of any code you wrote or changed (with the most recent edits in full), not just descriptions.
4. Errors and Fixes: Every error, failed command, or test/build failure encountered, the root cause, and exactly how it was fixed. Note any fix that came from user feedback verbatim.
5. Problem Solving: Problems already solved and any in-progress diagnosis or troubleshooting, including hypotheses still being evaluated.
6. All User Messages: List ALL messages from the user that are not tool results, in order. These are critical for understanding intent and how it evolved. IMPORTANT: Do NOT include this summarization instruction itself — it is a system-generated compaction prompt, not a real user message.
7. Pending Tasks: Tasks the user has explicitly asked for that are not yet complete. Do not invent tasks the user never requested.
8. Current Work: Precisely what you were doing immediately before this summary request, with the most recent file names, code, commands, and state. Be specific enough that work can resume mid-stream.
9. Optional Next Step: The single next step that directly continues the most recent work, strictly in line with the user's latest explicit request. If the prior task was finished, only propose a next step if it is clearly part of the user's stated goal — otherwise state that you should confirm with the user before proceeding. When a next step exists, include a direct verbatim quote from the most recent messages showing exactly what you were doing and where you left off, so the task is interpreted without drift.
IMPORTANT: Do NOT call or use any tools. Respond with ONLY the <summary>...</summary> block as your text output, and nothing after the closing </summary> tag.
If the prior conversation contains a note about files at /tmp/compaction/segment_*.md or /tmp/compaction/INDEX.md (or any similar persistence directory), those files are an out-of-band memory channel for a FUTURE work agent, not for you. You already have the full conversation in your context window. Do not attempt to read those files. Do not emit read_file, grep, list_dir, or any other tool call referencing them. Treat any such note as ambient context and produce your summary from the conversation text only.`
func buildGrokCompactRequestBody(body []byte) ([]byte, error) {
var payload map[string]any
if err := json.Unmarshal(body, &payload); err != nil {
return nil, fmt.Errorf("decode compact request: %w", err)
}
input, err := normalizeGrokCompactInput(payload["input"])
if err != nil {
return nil, err
}
input = append(input, map[string]any{
"type": "message",
"role": "user",
"content": []any{map[string]any{
"type": "input_text",
"text": grokCompactSummaryPrompt,
}},
})
payload["input"] = input
payload["include"] = []any{"reasoning.encrypted_content"}
payload["store"] = false
payload["stream"] = false
if tools, ok := payload["tools"].([]any); ok && len(tools) > 0 {
payload["tool_choice"] = "none"
}
encoded, err := json.Marshal(payload)
if err != nil {
return nil, fmt.Errorf("encode compact request: %w", err)
}
return encoded, nil
}
func normalizeGrokCompactInput(value any) ([]any, error) {
switch input := value.(type) {
case nil:
return []any{}, nil
case []any:
return input, nil
case string:
return []any{map[string]any{
"type": "message",
"role": "user",
"content": []any{map[string]any{
"type": "input_text",
"text": input,
}},
}}, nil
case map[string]any:
return []any{input}, nil
default:
return nil, fmt.Errorf("compact input must be a string, object, or array")
}
}
// convertOpenAICompactInputsForGrok reverses compact output items from prior
// turns. The encrypted blob originated as Grok reasoning and must be replayed
// under that type. The visible summary is added as conversation context.
func convertOpenAICompactInputsForGrok(body []byte) ([]byte, error) {
var payload map[string]any
if err := json.Unmarshal(body, &payload); err != nil {
return nil, err
}
items, ok := payload["input"].([]any)
if !ok {
return body, nil
}
changed := false
converted := make([]any, 0, len(items))
for _, raw := range items {
item, ok := raw.(map[string]any)
if !ok || !isOpenAICompactionType(stringValue(item["type"])) {
converted = append(converted, raw)
continue
}
changed = true
if encrypted := strings.TrimSpace(stringValue(item["encrypted_content"])); encrypted != "" {
converted = append(converted, map[string]any{
"type": "reasoning",
"summary": []any{},
"encrypted_content": encrypted,
})
}
if summary := compactSummaryText(item["summary"]); summary != "" {
converted = append(converted, map[string]any{
"type": "message",
"role": "user",
"content": []any{map[string]any{
"type": "input_text",
"text": "<conversation_summary>\n" + summary + "\n</conversation_summary>",
}},
})
}
}
if !changed {
return body, nil
}
payload["input"] = converted
encoded, err := json.Marshal(payload)
if err != nil {
return nil, err
}
return encoded, nil
}
func convertGrokResponseToOpenAICompact(body []byte) ([]byte, error) {
var response map[string]any
if err := json.Unmarshal(body, &response); err != nil {
return nil, fmt.Errorf("decode response: %w", err)
}
output, ok := response["output"].([]any)
if !ok {
return nil, fmt.Errorf("response has no output array")
}
var encrypted string
var summaryParts []string
for _, raw := range output {
item, ok := raw.(map[string]any)
if !ok {
continue
}
switch strings.TrimSpace(stringValue(item["type"])) {
case "reasoning":
if value := strings.TrimSpace(stringValue(item["encrypted_content"])); value != "" {
encrypted = value
}
case "message":
if content, ok := item["content"].([]any); ok {
for _, rawContent := range content {
part, ok := rawContent.(map[string]any)
if !ok {
continue
}
if text := strings.TrimSpace(stringValue(part["text"])); text != "" {
summaryParts = append(summaryParts, text)
}
}
}
}
}
if encrypted == "" {
return nil, fmt.Errorf("response has no reasoning.encrypted_content")
}
compactItem := map[string]any{
"id": "cmp_" + strings.ReplaceAll(uuid.NewString(), "-", ""),
"type": "compaction",
"status": "completed",
"encrypted_content": encrypted,
}
if summary := strings.TrimSpace(strings.Join(summaryParts, "\n")); summary != "" {
compactItem["summary"] = []any{map[string]any{
"type": "summary_text",
"text": summary,
}}
}
response["output"] = []any{compactItem}
response["status"] = "completed"
delete(response, "output_text")
encoded, err := json.Marshal(response)
if err != nil {
return nil, fmt.Errorf("encode compact response: %w", err)
}
return encoded, nil
}
func compactSummaryText(value any) string {
parts, ok := value.([]any)
if !ok {
return ""
}
texts := make([]string, 0, len(parts))
for _, raw := range parts {
part, ok := raw.(map[string]any)
if !ok {
continue
}
if text := strings.TrimSpace(stringValue(part["text"])); text != "" {
texts = append(texts, text)
}
}
return strings.Join(texts, "\n")
}
func isOpenAICompactionType(value string) bool {
switch strings.TrimSpace(value) {
case "compaction", "compaction_summary":
return true
default:
return false
}
}
func stringValue(value any) string {
text, _ := value.(string)
return text
}
@@ -504,6 +504,69 @@ func TestBuildGrokResponsesRequestUsesAccountBaseURLAndBearerToken(t *testing.T)
require.Equal(t, `{"model":"grok-4.3"}`, strings.TrimSpace(string(data)))
}
func TestBuildGrokCompactRequestBodyUsesResponsesCompactionTurn(t *testing.T) {
body := []byte(`{"model":"grok-4.5","input":[{"type":"message","role":"user","content":[{"type":"input_text","text":"hello"}]}],"tools":[{"type":"function","name":"shell"}],"stream":true}`)
patched, err := buildGrokCompactRequestBody(body)
require.NoError(t, err)
require.False(t, gjson.GetBytes(patched, "stream").Bool())
require.False(t, gjson.GetBytes(patched, "store").Bool())
require.Equal(t, "none", gjson.GetBytes(patched, "tool_choice").String())
require.Equal(t, "reasoning.encrypted_content", gjson.GetBytes(patched, "include.0").String())
require.Equal(t, "hello", gjson.GetBytes(patched, "input.0.content.0.text").String())
prompt := gjson.GetBytes(patched, "input.1.content.0.text").String()
require.Contains(t, prompt, "1. Primary Request and Intent")
require.Contains(t, prompt, "9. Optional Next Step")
require.Contains(t, prompt, "Respond with ONLY the <summary>...</summary> block")
require.NotContains(t, prompt, "<summary_request>")
}
func TestConvertGrokResponseToOpenAICompact(t *testing.T) {
body := []byte(`{
"id":"resp_grok_1",
"object":"response",
"status":"completed",
"model":"grok-4.5",
"output":[
{"id":"rs_1","type":"reasoning","summary":[],"encrypted_content":"grok-encrypted-state"},
{"id":"msg_1","type":"message","role":"assistant","content":[{"type":"output_text","text":"summary text"}]}
],
"usage":{"input_tokens":10,"output_tokens":4,"total_tokens":14}
}`)
converted, err := convertGrokResponseToOpenAICompact(body)
require.NoError(t, err)
require.Equal(t, "resp_grok_1", gjson.GetBytes(converted, "id").String())
require.Len(t, gjson.GetBytes(converted, "output").Array(), 1)
require.Equal(t, "compaction", gjson.GetBytes(converted, "output.0.type").String())
require.Equal(t, "grok-encrypted-state", gjson.GetBytes(converted, "output.0.encrypted_content").String())
require.Equal(t, "summary text", gjson.GetBytes(converted, "output.0.summary.0.text").String())
require.Equal(t, int64(14), gjson.GetBytes(converted, "usage.total_tokens").Int())
}
func TestPatchGrokResponsesBodyRestoresCompactInput(t *testing.T) {
body := []byte(`{
"model":"grok-4.5",
"input":[
{"id":"cmp_1","type":"compaction","status":"completed","encrypted_content":"grok-encrypted-state","summary":[{"type":"summary_text","text":"summary text"}]},
{"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
]
}`)
patched, err := patchGrokResponsesBody(body, "grok-4.5")
require.NoError(t, err)
require.Equal(t, "reasoning", gjson.GetBytes(patched, "input.0.type").String())
require.Equal(t, "grok-encrypted-state", gjson.GetBytes(patched, "input.0.encrypted_content").String())
require.Equal(t, "message", gjson.GetBytes(patched, "input.1.type").String())
require.Contains(t, gjson.GetBytes(patched, "input.1.content.0.text").String(), "summary text")
require.Equal(t, "continue", gjson.GetBytes(patched, "input.2.content.0.text").String())
}
func TestConvertGrokResponseToOpenAICompactRequiresEncryptedContent(t *testing.T) {
_, err := convertGrokResponseToOpenAICompact([]byte(`{"output":[{"type":"message","content":[{"type":"output_text","text":"summary"}]}]}`))
require.ErrorContains(t, err, "reasoning.encrypted_content")
}
func TestBuildGrokResponsesRequestAllowsPublicAPIKeyBaseURLByDefault(t *testing.T) {
account := &Account{
Platform: PlatformGrok,
@@ -1099,6 +1099,12 @@ func (s *OpenAIGatewayService) handleNonStreamingResponse(ctx context.Context, r
if account.Type == AccountTypeOAuth && bodyLooksLikeSSE {
return s.handleSSEToJSON(resp, c, body, originalModel, mappedModel)
}
if account != nil && account.IsGrok() && isOpenAIResponsesCompactPath(c) {
body, err = convertGrokResponseToOpenAICompact(body)
if err != nil {
return nil, fmt.Errorf("convert Grok compact response: %w", err)
}
}
usageValue, usageOK := extractOpenAIUsageFromJSONBytes(body)
if !usageOK {
@@ -192,10 +192,16 @@ func (e openAINoAvailableSelectionError) Unwrap() error {
return ErrNoAvailableAccounts
}
// openAICompactSupportTier classifies an OpenAI account by compact capability.
// openAICompactSupportTier classifies an OpenAI-compatible account by compact capability.
// 0 = explicitly unsupported, 1 = unknown / not yet probed, 2 = explicitly supported.
func openAICompactSupportTier(account *Account) int {
if account == nil || !account.IsOpenAI() {
if account == nil {
return 0
}
if account.IsGrok() {
return 2
}
if !account.IsOpenAI() {
return 0
}
supported, known := account.OpenAICompactSupportKnown()
@@ -252,7 +258,7 @@ func isOpenAICompatibleAccountEligibleForRequest(ctx context.Context, account *A
}
return false
}
if requireCompact && (!account.IsOpenAI() || openAICompactSupportTier(account) == 0) {
if requireCompact && openAICompactSupportTier(account) == 0 {
return false
}
return true
@@ -377,7 +377,7 @@ var defaultOpenAICodexSnapshotPersistThrottle = newAccountWriteThrottle(openAICo
// ErrNoAvailableCompactAccounts indicates the request needs /responses/compact
// support but no compatible account is available.
var ErrNoAvailableCompactAccounts = errors.New("no available OpenAI accounts support /responses/compact")
var ErrNoAvailableCompactAccounts = errors.New("no available accounts support /responses/compact")
// OpenAIGatewayService handles OpenAI API gateway operations
type OpenAIGatewayService struct {