diff --git a/research-wiki/designs/2026-07-21-m25-resilience-design.md b/research-wiki/designs/2026-07-21-m25-resilience-design.md index 04f0426..9b9f47f 100644 --- a/research-wiki/designs/2026-07-21-m25-resilience-design.md +++ b/research-wiki/designs/2026-07-21-m25-resilience-design.md @@ -105,6 +105,11 @@ **动因**: 第五轮(正午高峰,外部条件比首跑更严苛)剩余失败三分: 阶梯耗尽 ~1.1%(结构化档死亡率 3.2%,重问仅 1 次)、retry_exhausted ~0.9%(源1 尝试失败率 20.5% = 空补全 14.4% + **429 残漏 5.8%**,后者说明 AIMD 0.7 削减在临界点上方震荡)、400 拒绝 ~0.5%(语料固有毒负载,环境常量)。 **修正**: ① `PGW_STRUCTURED_MAX_RETRIES` 库缺省 1→2(instructor 库缺省 3 的保守版;成本只在解析失败时新增一跳);② AIMD `_CUT` 0.7→0.5(更快收敛到网关水位之下,Netflix 建议区间 0.5-0.9 内取激进端)。预期: 阶梯死亡 ~3.2%→~0.6%、429 出链后 retry_exhausted ~0.9%→~0.3%,合计残余 ≈ 1.0-1.4%(含 0.5% 环境常量)。 +### 3.38 迭代 5 补遗: 429 免重试预算(gRPC pushback 语义;2026-07-21 第六轮数据驱动) + +**动因**: 第六轮实测网关限速是**按账号速率**而非并发——外部负载高峰时我方仅 26 req/min 仍吃 16% 429,AIMD(并发型)压不住速率型限流;429 每次消耗 1/3 重试预算,饱和窗口里调用被"429 链"干净杀死(16 retry_exhausted/336)。 +**修正**: 携 Retry-After 的 429 是**服务端调度指令而非失败**(gRPC A6 pushback / SRE 语义): 照常退避(与 Retry-After 取大)、照常喂 AIMD/健康分,但**不消耗重试预算**;新增**调用级时间上限**(复用 `stall_window_s`,缺省 300s)在重试循环顶部兜底——持续饱和超窗即按既有 `stalled` 语义抛出,防无限循环。对下游的承诺变化: 饱和期调用最长等待 stall_window 而非快速失败(生产语义,迁移文档标注)。 + ### 3.4 不做与预留(方案 C 组件的接入点) - 账号级 429 共享退避: 不做;预留 = 冷却备忘 key 从 source_name 换 account_key 即可接入(LiteLLM 先例,治理粒度=配额粒度原则记入 ARCHITECTURE)。 diff --git a/research-wiki/migrations/chsanalyzer.md b/research-wiki/migrations/chsanalyzer.md index a1b6320..1bca441 100644 --- a/research-wiki/migrations/chsanalyzer.md +++ b/research-wiki/migrations/chsanalyzer.md @@ -188,3 +188,5 @@ stack = ExtractionProviderStack( | 开路时长指数递增,`retry_after_s` 上限由 cooldown_s(60s)变为 max_cooldown_s(缺省 300s) | **G1 消费点注意**: tracking.py arq defer 延期时长最多放大 5 倍;属期望行为(反复坏的源就该等更久) | | 阈值自动抬升由全局并发改为源级并发 | CHS .env 注释约定的本意(防并发误熔)保留;全局并发抬升是 M2 错误移植,已废除 | | 缺省选源 round_robin → health_aware(EWMA×在途 P2C) | 显式配置 `LLM__SELECTOR=round_robin` 者行为不变;迁移时建议直接吃新缺省 | +| 429 不消耗重试预算(pushback 语义),调用级时间上限 = stall_window_s | 饱和期调用延迟上限从 ~3×backoff 变为 stall_window(300s);CHS 快失败偏好者可调小 `LLM__BACKPRESSURE__STALL_WINDOW_S` | +| 结构化重问缺省 1→2 | 解析失败时最多多一跳成本;`PGW_STRUCTURED_MAX_RETRIES=1` 可显式还原 | diff --git a/src/polygateway/middleware/retry.py b/src/polygateway/middleware/retry.py index ef0af52..1c03eea 100644 --- a/src/polygateway/middleware/retry.py +++ b/src/polygateway/middleware/retry.py @@ -210,6 +210,16 @@ class RetryMW: attempt_fails: dict[str, int] = {} entered_at = self._now() # 调用级累计计时,循环内不重置(CHS governance.py:207) while True: + # 调用级时间上限(迭代 5): 429 免预算后的兜底,防饱和期无限循环。 + # 与 _on_no_runnable 同款双条件(CHS 口径): 本地超窗且全局无进展才判死 + stall = self._bp.stall_window_s + if self._now() - entered_at > stall and await self._quota.progress_age_s() > stall: + raise AllSourcesExhausted( + scope=self._scope, + reason="stalled", + retry_after_s=self._retry.backoff_base_s, + per_source_reasons=reasons, + ) picked, gate_rejections = await self._pick_runnable(reasons, attempt_fails) if picked is None: await self._on_no_runnable(gate_rejections, reasons, entered_at) @@ -217,7 +227,10 @@ class RetryMW: outcome = await self._attempt(request, *picked, reasons, attempt_fails) if isinstance(outcome, LLMResponse): return outcome - fails += 1 + # 429 = 服务端调度指令(gRPC pushback 语义,迭代 5): 按 Retry-After + # 退避但不消耗重试预算——饱和窗口里等待而非死亡;其余失败照常计数 + if _failure_reason(outcome.exc) != "rate_limited": + fails += 1 if fails >= self._retry.max_attempts: raise AllSourcesExhausted( scope=self._scope, @@ -226,7 +239,7 @@ class RetryMW: per_source_reasons=reasons, ) from outcome.exc if not outcome.immediate: - await self._sleep(self._backoff_delay(fails, outcome.exc)) + await self._sleep(self._backoff_delay(max(fails, 1), outcome.exc)) # —— 选源与准入(CHS _pick_runnable 120-167)—— diff --git a/tests/unit/test_retry.py b/tests/unit/test_retry.py index 9b3b69f..1bfc0bd 100644 --- a/tests/unit/test_retry.py +++ b/tests/unit/test_retry.py @@ -557,3 +557,47 @@ class TestDemotionInsertPosition: srcs = [_src("a"), _src("b"), _src("c")] out = _demote_call_failures(srcs, {"a": 2}, None) assert [s.name for s in out] == ["b", "c", "a"] + + +class TestRateLimitPushback: + """迭代 5(设计 §3.38): 429 是服务端调度指令,不耗重试预算;时间上限兜底。""" + + async def test_429_does_not_consume_retry_budget(self): + # 3 连 429 后成功——若 429 计预算,max_attempts=3 时第 4 次不会发生 + mw, _, _, transport, sleep, _ = _harness( + [_src("a")], + [ + TransientError("t1", status_code=429, retry_after_s=1.0), + TransientError("t2", status_code=429, retry_after_s=1.0), + TransientError("t3", status_code=429, retry_after_s=1.0), + _ok(), + ], + ) + resp = await mw(_REQ) + assert resp.content == "ok" + assert len(transport.calls) == 4 + assert len(sleep.delays) == 3 # 每次 429 仍按 Retry-After 退避 + + async def test_429_storm_bounded_by_stall_window(self): + # 持续 429 且时钟推进超 stall_window → stalled 兜底,不无限循环 + clock = FakeClock() + script = [TransientError(str(i), status_code=429, retry_after_s=30.0) for i in range(99)] + mw, _, _, _, _, _ = _harness([_src("a")], script, clock=clock) + + async def advancing_sleep(seconds): + clock.advance(seconds) + + mw._sleep = advancing_sleep + with pytest.raises(AllSourcesExhausted) as ei: + await mw(_REQ) + assert ei.value.reason == "stalled" + + async def test_non_429_transient_still_consumes_budget(self): + mw, _, _, transport, _, _ = _harness( + [_src("a")], + [TransientError("1"), TransientError("2"), TransientError("3")], + ) + with pytest.raises(AllSourcesExhausted) as ei: + await mw(_REQ) + assert ei.value.reason == "retry_exhausted" + assert len(transport.calls) == 3