feat: exempt 429 pushback from retry budget with stall ceiling

Round 6 hit account-level rate throttling the concurrency AIMD cannot
absorb: at 26 req/min the gateway still returned 16% 429s and each one
burned a third of the retry budget. Retry-After-guided 429s now back
off without consuming attempts (gRPC pushback semantics); the retry
loop gains a per-call ceiling using the same dual-condition stall
verdict as quota-wait (local window exceeded AND no global progress).
This commit is contained in:
2026-07-21 13:19:16 -04:00
parent 6b98a89bb3
commit 69968f2e8b
4 changed files with 66 additions and 2 deletions
@@ -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)。
+2
View File
@@ -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` 可显式还原 |
+15 -2
View File
@@ -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)——
+44
View File
@@ -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