feat: add AIMD adaptive concurrency pacing per source

P6 round 2 showed routing convergence turns account-level 429s into
the binding constraint (62% throttle rate at full concurrency). Each
source now carries a local AIMD limit: multiplicative cut on 429,
additive growth on success. Over-limit picks queue via the existing
quota-wait poll instead of burning retry budget or tripping the
circuit-open verdict.
This commit is contained in:
2026-07-21 10:38:48 -04:00
parent 910857fd13
commit 58061c7536
6 changed files with 157 additions and 2 deletions
@@ -86,6 +86,14 @@
实现(选源器无关,不触碰端口签名): `attempt_fails: dict[str, int]``__call__` **局部变量**逐层传参给 `_pick_runnable`(RetryMW 实例被并发调用共享,严禁实例属性,独立审查 M2);`_pick_runnable` 拿到 `order()` 结果后,把**本调用内已失败 ≥2 次**的源移到序列尾部(存在其他候选时)。效果: 健康源失败 1 次仍是首选(429/空补全属瞬态,原地退避重试最优);失败 2 次让位给次优源试一次;全部候选都失败过则按原序回退——自适应而非教条,调用结束状态即弃。病灶 5(预算与池结构无关)由此**间接闭环**: 预算数值不动,但尝试链质量由健康排序+降权保证(推演: 稳态下 3 次尝试大概率全落健康源,1−0.12³ ≈ 99.8%)。 实现(选源器无关,不触碰端口签名): `attempt_fails: dict[str, int]``__call__` **局部变量**逐层传参给 `_pick_runnable`(RetryMW 实例被并发调用共享,严禁实例属性,独立审查 M2);`_pick_runnable` 拿到 `order()` 结果后,把**本调用内已失败 ≥2 次**的源移到序列尾部(存在其他候选时)。效果: 健康源失败 1 次仍是首选(429/空补全属瞬态,原地退避重试最优);失败 2 次让位给次优源试一次;全部候选都失败过则按原序回退——自适应而非教条,调用结束状态即弃。病灶 5(预算与池结构无关)由此**间接闭环**: 预算数值不动,但尝试链质量由健康排序+降权保证(推演: 稳态下 3 次尝试大概率全落健康源,1−0.12³ ≈ 99.8%)。
### 3.35 迭代 1 补遗: AIMD 自适应并发(2026-07-21,P6 第二轮数据驱动)
**动因**: 三件套修复后 P6 第二轮实测——路由完全收敛(源5 尝试 11070→23),但 64 并发全压健康源后**网关账号级 429 率从 34% 恶化到 62%**,18 分钟即积累 180+ 终态失败(预算上限 160),纯退避重试救不动持续性背压。这正是 §3.4 预留的 AIMD 场景,数据证明它是达成 98% 的必要件,提前启用(用户已授权多轮迭代)。
**机制**(Netflix concurrency-limits 损失型客户端建议;进程本地,同健康记忆): 每源浮点并发上限 `limit`,初始 8、下限 1、上限 = worker 并发;**429 即乘性削减 ×0.7,真实成功即加性增长 +1/limit**。RetryMW 本地记每源在途数(选中 +1,尝试收尾 -1);`_pick_runnable``在途 ≥ limit` 的源跳过(reason=`adaptive_paced`,**不计 gate_rejections**——全被 paced 时走既有 quota-wait 轮询等待,绝不误判 CircuitOpen,也不消耗重试预算)。效果: 多余调用排队而非失败,并发自动收敛到网关可持续水位;429 同时仍喂健康 EWMA(排序信号)但不入熔断(§3.1 不变)。
**测试**: 单测削减/增长/上下限;RetryMW 集成——429 后准入收紧、被 paced 源跳过不耗预算、全 paced 走轮询非 CircuitOpen、在途数在异常/取消路径归零。
### 3.4 不做与预留(方案 C 组件的接入点) ### 3.4 不做与预留(方案 C 组件的接入点)
- 账号级 429 共享退避: 不做;预留 = 冷却备忘 key 从 source_name 换 account_key 即可接入(LiteLLM 先例,治理粒度=配额粒度原则记入 ARCHITECTURE)。 - 账号级 429 共享退避: 不做;预留 = 冷却备忘 key 从 source_name 换 account_key 即可接入(LiteLLM 先例,治理粒度=配额粒度原则记入 ARCHITECTURE)。
+9 -1
View File
@@ -8,7 +8,15 @@ SCOPE_REASONS = frozenset(
{"circuit_open", "retry_exhausted", "stalled", "quota_exhausted", "no_sources"} {"circuit_open", "retry_exhausted", "stalled", "quota_exhausted", "no_sources"}
) )
SOURCE_REASONS = frozenset( SOURCE_REASONS = frozenset(
{"network_error", "timeout", "rate_limited", "source_dead", "circuit_open", "cooldown"} {
"network_error",
"timeout",
"rate_limited",
"source_dead",
"circuit_open",
"cooldown",
"adaptive_paced", # M2.5 §3.35: AIMD 超限排队(quota-wait 通道)
}
) )
+13 -1
View File
@@ -33,7 +33,7 @@ from polygateway.errors import (
from polygateway.middleware.breaker import BreakerGate from polygateway.middleware.breaker import BreakerGate
from polygateway.middleware.ratelimit import QuotaGate from polygateway.middleware.ratelimit import QuotaGate
from polygateway.ports import OutcomeAwareSelector from polygateway.ports import OutcomeAwareSelector
from polygateway.sources import SourceCooldownMemo from polygateway.sources import AdaptivePacer, SourceCooldownMemo
from polygateway.streaming import StreamLivenessTimeout from polygateway.streaming import StreamLivenessTimeout
from polygateway.types import LLMResponse from polygateway.types import LLMResponse
@@ -118,6 +118,7 @@ class RetryMW:
backpressure: BackpressurePolicy, backpressure: BackpressurePolicy,
quota_full: str = "wait", quota_full: str = "wait",
cooldown_memo: SourceCooldownMemo | None = None, cooldown_memo: SourceCooldownMemo | None = None,
pacer: AdaptivePacer | None = None,
emitter: object | None = None, emitter: object | None = None,
now: Callable[[], float] = time.monotonic, now: Callable[[], float] = time.monotonic,
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep, sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
@@ -137,6 +138,8 @@ class RetryMW:
self._memo = cooldown_memo or SourceCooldownMemo(now=now) self._memo = cooldown_memo or SourceCooldownMemo(now=now)
# M2.5: 选源器可选健康喂数端口,构造期 isinstance 判定一次(设计 §3.2) # M2.5: 选源器可选健康喂数端口,构造期 isinstance 判定一次(设计 §3.2)
self._outcome_sink = selector if isinstance(selector, OutcomeAwareSelector) else None self._outcome_sink = selector if isinstance(selector, OutcomeAwareSelector) else None
# M2.5 §3.35: AIMD 自适应并发——429 收紧、成功回涨,超限调用排队不烧预算
self._pacer = pacer or AdaptivePacer(ceiling=64.0)
self._emitter = emitter self._emitter = emitter
self._now = now self._now = now
self._sleep = sleep self._sleep = sleep
@@ -184,6 +187,10 @@ class RetryMW:
gate_rejections += 1 gate_rejections += 1
reasons[cand.name] = "cooldown" reasons[cand.name] = "cooldown"
continue continue
if not self._pacer.admit(cand.name):
# AIMD 超限: 不计 gate_rejections → 走 quota-wait 排队,不误判熔断
reasons.setdefault(cand.name, "adaptive_paced")
continue
permit = await self._quota.try_acquire(cand) permit = await self._quota.try_acquire(cand)
if permit is None: if permit is None:
reasons.setdefault(cand.name, "rate_limited") reasons.setdefault(cand.name, "rate_limited")
@@ -196,6 +203,7 @@ class RetryMW:
if entry is None: if entry is None:
await self._settle_and_release(permit, 0) await self._settle_and_release(permit, 0)
if entry.allowed: if entry.allowed:
self._pacer.enter(cand.name)
return (cand, permit, entry), gate_rejections return (cand, permit, entry), gate_rejections
gate_rejections += 1 gate_rejections += 1
reasons[cand.name] = "circuit_open" reasons[cand.name] = "circuit_open"
@@ -261,6 +269,7 @@ class RetryMW:
await self._record_quietly(self._breaker.record_success(entry)) await self._record_quietly(self._breaker.record_success(entry))
await self._record_quietly(self._quota.mark_progress()) await self._record_quietly(self._quota.mark_progress())
self._feed_outcome(source.name, ok=True) self._feed_outcome(source.name, ok=True)
self._pacer.on_success(source.name)
response = self._build_response(source, result, call_id, started) response = self._build_response(source, result, call_id, started)
await self._emit(request, source, call_id, started, response=response) await self._emit(request, source, call_id, started, response=response)
return response return response
@@ -284,12 +293,15 @@ class RetryMW:
reasons[source.name] = reason reasons[source.name] = reason
attempt_fails[source.name] = attempt_fails.get(source.name, 0) + 1 attempt_fails[source.name] = attempt_fails.get(source.name, 0) + 1
self._feed_outcome(source.name, ok=False) self._feed_outcome(source.name, ok=False)
if reason == "rate_limited":
self._pacer.on_backpressure(source.name)
await self._record_quietly(self._breaker.record_failure(entry, reason, dead)) await self._record_quietly(self._breaker.record_failure(entry, reason, dead))
if not dead: if not dead:
actual = source.est_tokens # 保守: 失败请求可能已被网关计费(CHS 同款) actual = source.est_tokens # 保守: 失败请求可能已被网关计费(CHS 同款)
await self._emit(request, source, call_id, started, error=exc) await self._emit(request, source, call_id, started, error=exc)
return _Failed(exc, immediate=dead) return _Failed(exc, immediate=dead)
finally: finally:
self._pacer.leave(source.name)
await self._settle_and_release(permit, actual) await self._settle_and_release(permit, actual)
async def _on_rejected( async def _on_rejected(
+39
View File
@@ -82,6 +82,45 @@ class HealthAwareSelector:
return [head] + [s for s in ranked if s.name != head.name] return [head] + [s for s in ranked if s.name != head.name]
class AdaptivePacer:
"""AIMD 自适应并发(M2.5 设计 §3.35;Netflix concurrency-limits 损失型)。
429 是网关的"降速"信号: 乘性削减该源并发上限(×0.7),真实成功加性
增长(+1/limit),上限收敛到网关可持续水位;超限调用在 RetryMW 的
quota-wait 轮询里排队而非烧重试预算。进程本地,属 client 实例。
"""
_INITIAL = 8.0
_CUT = 0.7
_FLOOR = 1.0
def __init__(self, *, ceiling: float) -> None:
if ceiling < self._FLOOR:
raise ValueError("ceiling 不得小于下限 1")
self._ceiling = ceiling
self._limit: dict[str, float] = {}
self._inflight: dict[str, int] = {}
def limit(self, source_name: str) -> float:
return self._limit.get(source_name, min(self._INITIAL, self._ceiling))
def on_backpressure(self, source_name: str) -> None:
self._limit[source_name] = max(self._FLOOR, self.limit(source_name) * self._CUT)
def on_success(self, source_name: str) -> None:
cur = self.limit(source_name)
self._limit[source_name] = min(self._ceiling, cur + 1.0 / cur)
def admit(self, source_name: str) -> bool:
return self._inflight.get(source_name, 0) < self.limit(source_name)
def enter(self, source_name: str) -> None:
self._inflight[source_name] = self._inflight.get(source_name, 0) + 1
def leave(self, source_name: str) -> None:
self._inflight[source_name] = max(0, self._inflight.get(source_name, 0) - 1)
class SourceCooldownMemo: class SourceCooldownMemo:
"""进程本地的源冷却备忘(CHS governance.py:107 同款)。 """进程本地的源冷却备忘(CHS governance.py:107 同款)。
+40
View File
@@ -72,3 +72,43 @@ class TestHealthAwareSelector:
def test_missing_stats_defaults_to_zero_inflight(self): def test_missing_stats_defaults_to_zero_inflight(self):
sel = HealthAwareSelector(rng=lambda: 0.0) sel = HealthAwareSelector(rng=lambda: 0.0)
assert [s.name for s in sel.order([_src("s1")], {})] == ["s1"] assert [s.name for s in sel.order([_src("s1")], {})] == ["s1"]
class TestAdaptivePacer:
"""AIMD 自适应并发(M2.5 设计 §3.35): 429 乘性削减,成功加性增长。"""
def test_initial_and_bounds(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
assert pacer.limit("s1") == pytest.approx(8.0)
for _ in range(200):
pacer.on_backpressure("s1")
assert pacer.limit("s1") == pytest.approx(1.0) # 下限 1
for _ in range(2000):
pacer.on_success("s1")
assert pacer.limit("s1") == pytest.approx(32.0) # 上限 = ceiling
def test_cut_and_growth_math(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
pacer.on_backpressure("s1")
assert pacer.limit("s1") == pytest.approx(8.0 * 0.7)
before = pacer.limit("s1")
pacer.on_success("s1")
assert pacer.limit("s1") == pytest.approx(before + 1.0 / before)
def test_inflight_gate(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
for _ in range(200):
pacer.on_backpressure("s1") # limit → 1
assert pacer.admit("s1") is True
pacer.enter("s1")
assert pacer.admit("s1") is False # 在途 1 ≥ limit 1
pacer.leave("s1")
assert pacer.admit("s1") is True
pacer.leave("s1") # 多余 leave 不下穿 0
assert pacer._inflight.get("s1", 0) == 0
+48
View File
@@ -98,6 +98,7 @@ def _harness(
global_limits=_NO_GLOBAL, global_limits=_NO_GLOBAL,
rng=lambda: 0.0, rng=lambda: 0.0,
selector=None, selector=None,
pacer=None,
): ):
clock = clock or FakeClock() clock = clock or FakeClock()
limiter = InMemoryLimiter( limiter = InMemoryLimiter(
@@ -125,6 +126,7 @@ def _harness(
now=clock, now=clock,
sleep=sleep, sleep=sleep,
rng=rng, rng=rng,
pacer=pacer,
) )
return mw, limiter, gate, transport, sleep, clock return mw, limiter, gate, transport, sleep, clock
@@ -443,3 +445,49 @@ class TestM25Orchestration:
await mw(_REQ) await mw(_REQ)
g = gate._gates["a"] g = gate._gates["a"]
assert g.a0 + g.a1 == 0 assert g.a0 + g.a1 == 0
class TestAdaptivePacing:
"""AIMD 接线(设计 §3.35): 429 收紧准入,paced 源等待而非烧预算。"""
async def test_paced_source_waits_without_consuming_budget(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
for _ in range(200):
pacer.on_backpressure("a") # limit → 1
pacer.enter("a") # 模拟一个在途占满名额
mw, _, _, transport, _, _ = _harness(
[_src("a")], [_ok()], quota_full="fail_fast", pacer=pacer
)
with pytest.raises(AllSourcesExhausted) as ei:
await mw(_REQ)
assert ei.value.reason == "quota_exhausted" # 走配额等待通道,非 CircuitOpen
assert ei.value.per_source_reasons.get("a") == "adaptive_paced"
assert transport.calls == [] # 未发起尝试 → 不烧重试预算
async def test_429_cuts_limit_success_grows_it(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
mw, _, _, _, _, _ = _harness(
[_src("a")],
[TransientError("throttled", status_code=429), _ok()],
pacer=pacer,
)
await mw(_REQ)
# 429 削减一次(8→5.6),随后成功加性增长(5.6 + 1/5.6)
assert pacer.limit("a") == pytest.approx(8.0 * 0.7 + 1.0 / (8.0 * 0.7))
async def test_inflight_returns_to_zero_after_call(self):
from polygateway.sources import AdaptivePacer
pacer = AdaptivePacer(ceiling=32.0)
mw, _, _, _, _, _ = _harness(
[_src("a"), _src("b")],
[TransientError("x"), _ok()],
pacer=pacer,
)
await mw(_REQ)
assert pacer._inflight.get("a", 0) == 0
assert pacer._inflight.get("b", 0) == 0