From b2c9b4f6d90db0ab39f8aae939274f1a65d92dc4 Mon Sep 17 00:00:00 2001 From: wangbo Date: Tue, 4 Aug 2026 13:20:14 +0800 Subject: [PATCH] =?UTF-8?q?test(runner):=20=E8=A1=A5=E9=BD=90=E5=BB=B6?= =?UTF-8?q?=E8=BF=9F=20429=20=E6=B5=81=E5=BC=8F=E8=BD=AE=E8=BD=AC=E5=9B=9E?= =?UTF-8?q?=E5=BD=92?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 验证上游延迟返回 429 时,在首个 delta 前仍分类为可重试限流并冷却轮转到第二候选,响应不混入失败候选内容。 验证:完整 Go 单元测试通过;隔离 PostgreSQL 下新增 HTTP acceptance 场景通过;gofmt 与差异检查通过。 --- apps/api/internal/clients/clients_test.go | 78 +++++++++++++++++++ ...policy_http_acceptance_integration_test.go | 62 +++++++++++++++ 2 files changed, 140 insertions(+) diff --git a/apps/api/internal/clients/clients_test.go b/apps/api/internal/clients/clients_test.go index 7396291..d6a9aa4 100644 --- a/apps/api/internal/clients/clients_test.go +++ b/apps/api/internal/clients/clients_test.go @@ -344,6 +344,84 @@ func TestOpenAIClientChatContract(t *testing.T) { } } +func TestOpenAIClientDelayedStream429ReturnsBeforeAnyDelta(t *testing.T) { + requestStarted := make(chan struct{}) + releaseResponse := make(chan struct{}) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + close(requestStarted) + <-releaseResponse + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusTooManyRequests) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": map[string]any{"message": "controlled delayed 429"}, + }) + })) + defer server.Close() + release := func() { + select { + case <-releaseResponse: + default: + close(releaseResponse) + } + } + defer release() + + type runResult struct { + response Response + err error + } + deltaCount := 0 + done := make(chan runResult, 1) + go func() { + response, err := (OpenAIClient{HTTPClient: server.Client()}).Run(context.Background(), Request{ + Kind: "chat.completions", + Model: "delayed-429-model", + Stream: true, + Body: map[string]any{ + "model": "delayed-429-model", + "stream": true, + "messages": []any{map[string]any{"role": "user", "content": "ping"}}, + }, + Candidate: store.RuntimeModelCandidate{ + BaseURL: server.URL, + ProviderModelName: "delayed-429-model", + Credentials: map[string]any{"apiKey": "test-key"}, + }, + StreamDelta: func(StreamDeltaEvent) error { + deltaCount++ + return nil + }, + }) + done <- runResult{response: response, err: err} + }() + + select { + case <-requestStarted: + case <-time.After(time.Second): + t.Fatal("upstream request did not start") + } + select { + case result := <-done: + t.Fatalf("stream returned before upstream status arrived: response=%+v err=%v", result.response, result.err) + case <-time.After(50 * time.Millisecond): + } + release() + + var result runResult + select { + case result = <-done: + case <-time.After(time.Second): + t.Fatal("stream did not return after delayed 429") + } + clientErr, ok := result.err.(*ClientError) + if !ok || clientErr.Code != "rate_limit" || clientErr.StatusCode != http.StatusTooManyRequests || !clientErr.Retryable { + t.Fatalf("unexpected delayed 429 error: %T %+v", result.err, result.err) + } + if deltaCount != 0 { + t.Fatalf("delayed 429 emitted %d stream deltas, want 0", deltaCount) + } +} + func TestOpenAIClientAliyunChatSendsStandardInputAudioAndNormalizesLegacyAudioURL(t *testing.T) { var captured map[string]any server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { diff --git a/apps/api/internal/httpapi/failure_policy_http_acceptance_integration_test.go b/apps/api/internal/httpapi/failure_policy_http_acceptance_integration_test.go index 1c74048..27e5991 100644 --- a/apps/api/internal/httpapi/failure_policy_http_acceptance_integration_test.go +++ b/apps/api/internal/httpapi/failure_policy_http_acceptance_integration_test.go @@ -871,6 +871,68 @@ func (f *failurePolicyAcceptanceFixture) testAllCoolingAsyncRecovery(t *testing. } func (f *failurePolicyAcceptanceFixture) testStreamingBoundary(t *testing.T) { + t.Run("延迟 429 在首个 delta 前仍冷却并轮转", func(t *testing.T) { + model := "acceptance-stream-delayed-429-" + f.suffix + baseModelID := f.createBaseModel(t, model, "text_generate") + failed := newControlledProviderStub(func(_ int, w http.ResponseWriter, _ *http.Request) { + time.Sleep(125 * time.Millisecond) + writeStubJSON(w, http.StatusTooManyRequests, map[string]any{ + "error": map[string]any{"message": "controlled delayed 429"}, + }) + }) + defer failed.close() + success := newControlledProviderStub(func(_ int, w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + _, _ = io.WriteString(w, "data: {\"id\":\"chunk-delayed-429\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{\"role\":\"assistant\",\"content\":\"delayed 429 fallback\"},\"finish_reason\":null}]}\n\n") + _, _ = io.WriteString(w, "data: {\"id\":\"chunk-delayed-429\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n") + _, _ = io.WriteString(w, "data: [DONE]\n\n") + }) + defer success.close() + platformA, modelA := f.createPlatformModel(t, "stream-delayed-429-a", "openai", failed.server.URL+"/v1", 10, baseModelID, "stream-delayed-429", "text_generate", 1, "", nil) + platformB, _ := f.createPlatformModel(t, "stream-delayed-429-b", "openai", success.server.URL+"/v1", 20, baseModelID, "stream-delayed-success", "text_generate", 1, "", nil) + + startedAt := time.Now() + status, _, body := f.postJSON(t, "/api/v1/chat/completions", map[string]any{ + "model": model, "stream": true, "messages": []any{map[string]any{"role": "user", "content": "stream delayed 429"}}, + }, nil) + bodyText := string(body) + if status != http.StatusOK || failed.calls.Load() != 1 || success.calls.Load() != 1 { + t.Fatalf("unexpected delayed 429 failover status=%d calls=%d/%d body=%s", status, failed.calls.Load(), success.calls.Load(), body) + } + if strings.Count(bodyText, "delayed 429 fallback") != 1 || strings.Contains(bodyText, "controlled delayed 429") || !strings.Contains(bodyText, "[DONE]") { + t.Fatalf("delayed 429 response leaked or duplicated candidate output: %s", bodyText) + } + + detail := f.loadTask(t, f.taskIDForModel(t, model, startedAt)) + if len(detail.Attempts) != 2 || detail.Attempts[0].PlatformID != platformA.ID || detail.Attempts[0].ErrorCode != "upstream_rate_limited" || detail.Attempts[1].PlatformID != platformB.ID || detail.Attempts[1].Status != "succeeded" || countFailureEffects(detail, "cooldown") != 1 { + t.Fatalf("unexpected delayed 429 attempt chain: %+v", detail) + } + decisionCount := 0 + for _, attempt := range detail.Attempts { + trace := append([]map[string]any(nil), attempt.Trace...) + for _, raw := range anySlice(attempt.DynamicMetrics["trace"]) { + if entry, ok := raw.(map[string]any); ok { + trace = append(trace, entry) + } + } + for _, entry := range trace { + if entry["event"] == "failure_decision" && entry["route"] == "next" && entry["effect"] == "cooldown" && entry["category"] == "rate_limit" { + decisionCount++ + } + } + } + if decisionCount != 1 { + t.Fatalf("delayed 429 should record one cooldown-and-next decision, got %d: %+v", decisionCount, detail) + } + var modelCooling bool + if err := f.pool.QueryRow(f.ctx, `SELECT COALESCE(cooldown_until > now(), false) FROM platform_models WHERE id=$1::uuid`, modelA.ID).Scan(&modelCooling); err != nil { + t.Fatalf("read delayed 429 cooldown state: %v", err) + } + if !modelCooling { + t.Fatal("delayed 429 failed platform model was not cooled") + } + }) + t.Run("首个 delta 前允许轮转", func(t *testing.T) { model := "acceptance-stream-before-" + f.suffix baseModelID := f.createBaseModel(t, model, "text_generate")