fix: address M2.5 verifier findings before merge

AIMD ceiling now respects per-source max_concurrency and the pacer is
assembled explicitly in the client; MIN_CALLS parses as strict int;
acceptance doc corrects source-5 attempt count to 549; design and
migration notes aligned with implemented 429/stall/suppression
semantics and AIMD constants documented.
This commit is contained in:
2026-07-21 21:10:02 -04:00
parent 9d28e03d49
commit a06761917e
7 changed files with 61 additions and 5 deletions
+2
View File
@@ -42,6 +42,8 @@ LLM_CIRCUIT_BREAKER_COOLDOWN=60 # 或 LLM__BREAKER__COOLDOWN_S
# LLM__BREAKER__FAIL_RATE=0.6 # 窗口失败率阈值(429 不计入) # LLM__BREAKER__FAIL_RATE=0.6 # 窗口失败率阈值(429 不计入)
# LLM__BREAKER__WINDOW_S=60 # 失败率窗口(双 30s 桶) # LLM__BREAKER__WINDOW_S=60 # 失败率窗口(双 30s 桶)
# LLM__BREAKER__MAX_COOLDOWN_S=300 # 开路指数退避封顶(缺省 max(300, cooldown)) # LLM__BREAKER__MAX_COOLDOWN_S=300 # 开路指数退避封顶(缺省 max(300, cooldown))
# ── AIMD 自适应并发(M2.5,库常量非 env 键): 每源初始 8,429 ×0.5,成功 +1/limit,
# ── ceiling = max(64, 源级 MAX_CONCURRENCY);禁用需构造函数注入自定义 pacer ──
# LLM__QUOTA_FULL=wait # wait(默认) | fail_fast # LLM__QUOTA_FULL=wait # wait(默认) | fail_fast
# ══ 装配选择(PGW_*)══ # ══ 装配选择(PGW_*)══
@@ -108,7 +108,7 @@
### 3.38 迭代 5 补遗: 429 免重试预算(gRPC pushback 语义;2026-07-21 第六轮数据驱动) ### 3.38 迭代 5 补遗: 429 免重试预算(gRPC pushback 语义;2026-07-21 第六轮数据驱动)
**动因**: 第六轮实测网关限速是**按账号速率**而非并发——外部负载高峰时我方仅 26 req/min 仍吃 16% 429,AIMD(并发型)压不住速率型限流;429 每次消耗 1/3 重试预算,饱和窗口里调用被"429 链"干净杀死(16 retry_exhausted/336)。 **动因**: 第六轮实测网关限速是**按账号速率**而非并发——外部负载高峰时我方仅 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 而非快速失败(生产语义,迁移文档标注)。 **修正**: 一切 429 均视为**服务端调度指令而非失败**(gRPC A6 pushback / SRE 语义;无 Retry-After 头时用本地指数退避代行): 照常退避(与 Retry-After 取大)、照常喂 AIMD/健康分,但**不消耗重试预算**;新增**调用级时间兜底**在重试循环顶部——与 quota-wait 同款**双条件**(本地超 `stall_window_s` **且**全局无进展)按既有 `stalled` 语义抛出,防无限循环(全局持续有进展时单调用可等待超过 stall_window,非严格上限)。对下游的承诺变化: 饱和期调用最长等待 stall_window 而非快速失败(生产语义,迁移文档标注)。
### 3.39 迭代 6 补遗: 健康证据抑制连续通道(2026-07-21,第八轮取证驱动) ### 3.39 迭代 6 补遗: 健康证据抑制连续通道(2026-07-21,第八轮取证驱动)
@@ -141,7 +141,7 @@
| 旧行为 | 处置 | | 旧行为 | 处置 |
|---|---| |---|---|
| 连续失败 ≥ 阈值开路 | **保留**(通道 1) | | 连续失败 ≥ 阈值开路 | **保留为通道 1,但受 §3.39 健康证据抑制**(窗口样本充足且失败率低时连败不开路) |
| 成功清零连续失败计数 | **保留**(仅作用于通道 1 计数;窗口计数不清) | | 成功清零连续失败计数 | **保留**(仅作用于通道 1 计数;窗口计数不清) |
| 阈值自动抬升 max(配置, 并发×2) | **有意废除**(病灶 2;M2 设计错误,迁移文档更新) | | 阈值自动抬升 max(配置, 并发×2) | **有意废除**(病灶 2;M2 设计错误,迁移文档更新) |
| SourceDead force_open / 探针租约 / epoch fencing / release_probe / 半开单探针 | **保留**(逐字不动) | | SourceDead force_open / 探针租约 / epoch fencing / release_probe / 半开单探针 | **保留**(逐字不动) |
@@ -9,7 +9,7 @@
|---|---|---| |---|---|---|
| 成功率 | 58.1%(4648/8000) | **98.96%(7917/8000)** | | 成功率 | 58.1%(4648/8000) | **98.96%(7917/8000)** |
| 终态失败构成 | 3318 retry_exhausted + 34 其他 | **37 retry_exhausted(0.46%)+ 45 语料毒负载 400(0.56%)+ 1 结构化** | | 终态失败构成 | 3318 retry_exhausted + 34 其他 | **37 retry_exhausted(0.46%)+ 45 语料毒负载 400(0.56%)+ 1 结构化** |
| 坏源吸流(源5 尝试) | 11070(76%) | ~300(3%,熔断+健康路由压制) | | 坏源吸流(源5 尝试) | 11070(76%) | 549(真实尝试的 9.3%,压制 20 倍;独立核验重算更正) |
| 用时 / 成本 | 1h47m / 23.4 元 | 1h17m / 38.7 元(token 2321 万) | | 用时 / 成本 | 1h47m / 23.4 元 | 1h17m / 38.7 元(token 2321 万) |
| 缓存命中 | 2373 | 4153 | | 缓存命中 | 2373 | 4153 |
| P3 结构化成功率(故障池内) | 0.218 | 0.853 | | P3 结构化成功率(故障池内) | 0.218 | 0.853 |
+3 -1
View File
@@ -188,5 +188,7 @@ stack = ExtractionProviderStack(
| 开路时长指数递增,`retry_after_s` 上限由 cooldown_s(60s)变为 max_cooldown_s(缺省 300s) | **G1 消费点注意**: tracking.py arq defer 延期时长最多放大 5 倍;属期望行为(反复坏的源就该等更久) | | 开路时长指数递增,`retry_after_s` 上限由 cooldown_s(60s)变为 max_cooldown_s(缺省 300s) | **G1 消费点注意**: tracking.py arq defer 延期时长最多放大 5 倍;属期望行为(反复坏的源就该等更久) |
| 阈值自动抬升由全局并发改为源级并发 | CHS .env 注释约定的本意(防并发误熔)保留;全局并发抬升是 M2 错误移植,已废除 | | 阈值自动抬升由全局并发改为源级并发 | CHS .env 注释约定的本意(防并发误熔)保留;全局并发抬升是 M2 错误移植,已废除 |
| 缺省选源 round_robin → health_aware(EWMA×在途 P2C) | 显式配置 `LLM__SELECTOR=round_robin` 者行为不变;迁移时建议直接吃新缺省 | | 缺省选源 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` | | 429 不消耗重试预算(pushback 语义),调用级时间兜底 = 双条件 stall(本地超窗且全局无进展) | 饱和期调用延迟大幅拉长(全局有进展时可超 stall_window);CHS 快失败偏好者可调小 `LLM__BACKPRESSURE__STALL_WINDOW_S` |
| **AIMD 自适应并发**(库常量: 每源初始 8、429 削减 ×0.5、成功 +1/limit、下限 1、ceiling=max(64, 源级 max_concurrency)) | 冷启动每源并发从"无上限"变为 8 起步爬升——高吞吐下游首分钟吞吐低于旧行为;无 env 开关(韧性常量,同 jitter 系数先例),需要禁用者构造函数注入自定义 pacer |
| 连败通道受健康证据抑制(窗口样本 ≥ min_calls 且失败率 < fail_rate 时 5 连败不开路) | CHS "连败必开"在高流量健康源上不再成立(防随机噪声误熔唯一好源);冷启动/低流量语义不变 |
| 结构化重问缺省 1→2 | 解析失败时最多多一跳成本;`PGW_STRUCTURED_MAX_RETRIES=1` 可显式还原 | | 结构化重问缺省 1→2 | 解析失败时最多多一跳成本;`PGW_STRUCTURED_MAX_RETRIES=1` 可显式还原 |
+6
View File
@@ -25,6 +25,7 @@ from polygateway.middleware.telemetry import TelemetryEmitter, TelemetryMW
from polygateway.pricing import PricingTable from polygateway.pricing import PricingTable
from polygateway.providers import get_provider from polygateway.providers import get_provider
from polygateway.sources import ( from polygateway.sources import (
AdaptivePacer,
HealthAwareSelector, HealthAwareSelector,
LeastInflightSelector, LeastInflightSelector,
RoundRobinSelector, RoundRobinSelector,
@@ -97,6 +98,10 @@ class GatewayClient:
backpressure=backpressure, backpressure=backpressure,
quota_full=quota_full, quota_full=quota_full,
cooldown_memo=SourceCooldownMemo(now=now), cooldown_memo=SourceCooldownMemo(now=now),
# AIMD ceiling 尊重源级静态并发上限(独立核验 I1: 不得静默钳制大于 64 的配置)
pacer=AdaptivePacer(
ceiling=float(max([64, *(s.max_concurrency for s in sources if s.max_concurrency)]))
),
emitter=emitter, emitter=emitter,
now=now, now=now,
sleep=sleep, sleep=sleep,
@@ -128,6 +133,7 @@ class GatewayClient:
) )
) )
self._structured_available = structured_strategy is not None self._structured_available = structured_strategy is not None
self._terminal = terminal # 内部引用: 装配自省/测试用
self._handler = compose(middlewares, terminal) self._handler = compose(middlewares, terminal)
self._transport = transport self._transport = transport
self._telemetry = telemetry self._telemetry = telemetry
+2 -1
View File
@@ -231,7 +231,8 @@ def _load_breaker(
# 派生规则: 探针租约须撑过一次最慢调用,且不短于冷却期(第三项保证守卫恒成立) # 派生规则: 探针租约须撑过一次最慢调用,且不短于冷却期(第三项保证守卫恒成立)
probe_ttl_s = max(2 * slowest, cooldown_s, probe_floor) probe_ttl_s = max(2 * slowest, cooldown_s, probe_floor)
# M2.5 失败率通道参数(可选键,库缺省——韧性参数缺省先例同 backpressure) # M2.5 失败率通道参数(可选键,库缺省——韧性参数缺省先例同 backpressure)
min_calls = int(_opt_float(env, f"{scope}__BREAKER__MIN_CALLS", 10)) raw_min = env.get(f"{scope}__BREAKER__MIN_CALLS")
min_calls = int(_cast(raw_min, "int", f"{scope}__BREAKER__MIN_CALLS")) if raw_min else 10
fail_rate = _opt_float(env, f"{scope}__BREAKER__FAIL_RATE", 0.6) fail_rate = _opt_float(env, f"{scope}__BREAKER__FAIL_RATE", 0.6)
window_s = _opt_float(env, f"{scope}__BREAKER__WINDOW_S", 60.0) window_s = _opt_float(env, f"{scope}__BREAKER__WINDOW_S", 60.0)
max_cooldown_s = _opt_float(env, f"{scope}__BREAKER__MAX_COOLDOWN_S", max(300.0, cooldown_s)) max_cooldown_s = _opt_float(env, f"{scope}__BREAKER__MAX_COOLDOWN_S", max(300.0, cooldown_s))
+45
View File
@@ -119,3 +119,48 @@ class TestAdaptivePacer:
assert pacer.admit("s1") is True assert pacer.admit("s1") is True
pacer.leave("s1") # 多余 leave 不下穿 0 pacer.leave("s1") # 多余 leave 不下穿 0
assert pacer._inflight.get("s1", 0) == 0 assert pacer._inflight.get("s1", 0) == 0
class TestPacerAssembly:
"""独立核验 I1/M5: ceiling 尊重源级并发;取消路径在途归零。"""
def test_ceiling_respects_large_source_concurrency(self):
from polygateway.client import GatewayClient
from polygateway.config import GatewaySettings
env = {
"LLM__QWEN__1__BASE_URL": "https://gw.example/v1",
"LLM__QWEN__1__API_KEY": "sk-a",
"LLM__QWEN__1__MODEL": "m",
"LLM__QWEN__1__TIMEOUT_S": "60",
"LLM__QWEN__1__MAX_CONCURRENCY": "128",
"LLM_MAX_RETRIES": "3",
"LLM_RETRY_BASE_DELAY": "2.0",
"LLM_RETRY_MAX_DELAY": "30.0",
"LLM_CIRCUIT_BREAKER_THRESHOLD": "5",
"LLM_CIRCUIT_BREAKER_COOLDOWN": "60",
"PGW_CACHE_BACKEND": "none",
"PGW_TELEMETRY_BACKEND": "none",
}
client = GatewayClient.from_settings(GatewaySettings.from_env("LLM", env=env))
assert client._terminal._pacer._ceiling == pytest.approx(128.0)
def test_min_calls_rejects_float_value(self):
from polygateway.config import GatewaySettings
env = {
"LLM__QWEN__1__BASE_URL": "https://gw.example/v1",
"LLM__QWEN__1__API_KEY": "sk-a",
"LLM__QWEN__1__MODEL": "m",
"LLM__QWEN__1__TIMEOUT_S": "60",
"LLM_MAX_RETRIES": "3",
"LLM_RETRY_BASE_DELAY": "2.0",
"LLM_RETRY_MAX_DELAY": "30.0",
"LLM_CIRCUIT_BREAKER_THRESHOLD": "5",
"LLM_CIRCUIT_BREAKER_COOLDOWN": "60",
"LLM__BREAKER__MIN_CALLS": "10.5",
"PGW_CACHE_BACKEND": "none",
"PGW_TELEMETRY_BACKEND": "none",
}
with pytest.raises(ValueError, match="MIN_CALLS"):
GatewaySettings.from_env("LLM", env=env)