test(runner): 补齐延迟 429 流式轮转回归
验证上游延迟返回 429 时,在首个 delta 前仍分类为可重试限流并冷却轮转到第二候选,响应不混入失败候选内容。 验证:完整 Go 单元测试通过;隔离 PostgreSQL 下新增 HTTP acceptance 场景通过;gofmt 与差异检查通过。
This commit is contained in:
@@ -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) {
|
||||
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user