fix: address M2 verifier findings in soak harness
This commit is contained in:
@@ -71,6 +71,9 @@ PGW_TELEMETRY_BACKEND=none # sqlite | postgres | none(必填)
|
||||
# EMBED__EXPECTED_DIM=768 # 可选: 维度校验,不符抛 ResultInvalid
|
||||
|
||||
# ══ SOAK scope(压测 harness 专用;tools/soak/run_soak.py --scope SOAK)══
|
||||
# 网关保护(设计 §8.1 签字值,run_soak 强制: 必配且 ≤ 100/600,否则拒跑)
|
||||
# SOAK__GLOBAL__MAX_CONCURRENCY=100
|
||||
# SOAK__GLOBAL__RPM=600
|
||||
# 健康源 + 故障源混编(findings §3): 坏 key 源 / 黑洞源 / 紧看门狗源 / 紧闸源
|
||||
# SOAK__MINIMAX__1__BASE_URL= # 健康源(真实网关)
|
||||
# SOAK__MINIMAX__1__API_KEY=
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
---
|
||||
type: finding
|
||||
node_id: finding:m2-verifier-fixes
|
||||
title: "M2 verifier 三项 Important 补齐(不变量接线/网关保护/P3 验收)"
|
||||
date: 2026-07-21
|
||||
---
|
||||
|
||||
# M2 verifier 三项 Important 补齐
|
||||
|
||||
独立 verifier(全新上下文)验收 M2:0 Critical,核心库逐条达标;3 Important + 4 Minor 全部当日修复。
|
||||
|
||||
| # | 问题 | 修复 | 证据 |
|
||||
|---|---|---|---|
|
||||
| I1 | 不变量 1(记账归零/gate 可再准入)未接线进 harness | `run_soak._live_checks` 跑后活检 + 新增 `inv_gate_reenterable`(探针悬挂检测,拿到即归还) | P3 验收报告两项 PASS |
|
||||
| I2 | 签字网关保护(100/600)无落点 | `_guard` 硬门:压测 scope 必配 `__GLOBAL__` 键且 ≤ 签字值;.env.example SOAK 段写明 | 缺键拒跑 |
|
||||
| I3 | P3 中等规模验收未跑 | 300 calls × 2 workers 真实跑:297 ok,7 不变量 PASS,结构化基线 0.942,p50/p95=16.9s/36.4s;真实网关自然出现 22 次空补全全被 Transient 重试吸收 | tests/outputs/soak/soak_20260721_015028.md |
|
||||
| M4-M7 | settle 置位偏离未声明/对标未打钩/场景帽/故障存在性断言 | docstring 声明;chsanalyzer.md 十项打钩;`capped_budget`;`inv_fault_errors_present` | 单测 16 例 |
|
||||
@@ -30,6 +30,11 @@
|
||||
"id": "plan:m2-distributed",
|
||||
"label": "M2 分布式实现计划",
|
||||
"type": "plan"
|
||||
},
|
||||
{
|
||||
"id": "finding:m2-verifier-fixes",
|
||||
"label": "M2 verifier 三项 Important 补齐(不变量接线/网关保护/P3 验收)",
|
||||
"type": "finding"
|
||||
}
|
||||
],
|
||||
"links": [
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# Research Wiki 索引
|
||||
|
||||
> 自动生成,更新时间:2026-07-21 03:46 UTC
|
||||
> 自动生成,更新时间:2026-07-21 06:00 UTC
|
||||
|
||||
## design (4)
|
||||
- [2026-07-20-m1-core-design](designs/2026-07-20-m1-core-design.md) `design:2026-07-20-m1-core-design`
|
||||
@@ -8,8 +8,9 @@
|
||||
- [M1 核心里程碑设计:公共签名冻结与治理栈落地](designs/m1-core-design.md) `design:m1-core-design`
|
||||
- [M2 分布式:Redis 治理后端+背压+Postgres 遥测+pricing+Embedding+压测 harness](designs/m2-distributed.md) `design:m2-distributed`
|
||||
|
||||
## finding (2)
|
||||
## finding (3)
|
||||
- [2026-07-20-m2-soak-workload](findings/2026-07-20-m2-soak-workload.md) `finding:2026-07-20-m2-soak-workload`
|
||||
- [M2 verifier 三项 Important 补齐(不变量接线/网关保护/P3 验收)](findings/m2-verifier-fixes.md) `finding:m2-verifier-fixes`
|
||||
- [M2 真实数据压测: 场景矩阵与数据清单](findings/m2-soak-workload.md) `finding:m2-soak-workload`
|
||||
|
||||
## plan (4)
|
||||
|
||||
@@ -17,3 +17,5 @@
|
||||
- [2026-07-21 03:46 UTC] 新增 plan: M2 分布式实现计划 (plan:m2-distributed)
|
||||
- [2026-07-21 03:46 UTC] 新增边: plan:m2-distributed --implements--> design:m2-distributed
|
||||
- [2026-07-21 03:46 UTC] 重建索引: 12 篇页面
|
||||
- [2026-07-21 06:00 UTC] 新增 finding: M2 verifier 三项 Important 补齐(不变量接线/网关保护/P3 验收) (finding:m2-verifier-fixes)
|
||||
- [2026-07-21 06:00 UTC] 重建索引: 13 篇页面
|
||||
|
||||
@@ -11,20 +11,22 @@
|
||||
| M2 后 | VLM scope 治理核心(governance/limiter/scripts/provider_gate/selector/streaming + VLM invoker) | Redis 六道闸+契约测试、多源选源、跨进程熔断、背压 stall |
|
||||
| M3 后 | OCR scope(MonkeyOcrParseInvoker → `OcrLayoutPort`)、`core/eval/judge.py` 收编,**全量迁移** | OCR 端口族 + MonkeyOCR transport |
|
||||
|
||||
**能力对标清单**(库必须逐项对等,任一缺失即库的边界缺口):
|
||||
**能力对标清单**(库必须逐项对等,任一缺失即库的边界缺口;✅ = M2 验收打钩 2026-07-21,证据为库测试):
|
||||
|
||||
| # | 能力 | 项目侧证据 |
|
||||
|---|---|---|
|
||||
| 1 | 六道闸原子限流(全局/单源 × 并发/RPM/TPM),拒绝零副作用 | `app/coordination/scripts.py:6-34` |
|
||||
| 2 | TPM 预扣入场、settle 按实际 usage 落回 acquire 窗口多退少补 | `limiter.py:55-62`、`scripts.py:44-48` |
|
||||
| 3 | 并发 lease 带 TTL(防进程死亡泄漏)+ release/settle 幂等 | `limiter.py:46-63`、`tests/contracts_limiter.py:66-73` |
|
||||
| 4 | 窗口 id 用 Redis 服务器时钟(多进程口径统一) | `limiter.py:95-98` |
|
||||
| 5 | 跨进程熔断:单探针半开、探针租约 TTL、epoch fencing、force_open | `scripts.py:82-225`、`provider_gate.py` |
|
||||
| 6 | 背压 stall 双条件判卡死(本地等够 + 全局无进展)+ mark_progress 全局活性 | `governance.py:270-281`、`limiter.py:193-209` |
|
||||
| 7 | 换源重试:失败跨源累计、退避取 Retry-After 较大值、源冷却备忘 | `governance.py:105-118, 244-261` |
|
||||
| 8 | 错误分类:429 body 细分 insufficient_quota、工件级失败不熔断 | `invokers.py:144-166`、`governance.py:237-239` |
|
||||
| 9 | scope 级不可用结构化错误(reason 7 种 + retry_after_s),供 worker 延期重投 | `errors.py:140-180`、`workers/tracking.py:406-428` |
|
||||
| 10 | 流式三层活性看门狗;thinking token 刷活性不计结果 | `streaming.py`、`invokers.py:55-79` |
|
||||
| # | 能力 | 项目侧证据 | M2 打钩 |
|
||||
|---|---|---|---|
|
||||
| 1 | 六道闸原子限流(全局/单源 × 并发/RPM/TPM),拒绝零副作用 | `app/coordination/scripts.py:6-34` | ✅ `backends/redis/limiter.py` Lua 逐字移植;契约 redis 参数全绿 |
|
||||
| 2 | TPM 预扣入场、settle 按实际 usage 落回 acquire 窗口多退少补 | `limiter.py:55-62`、`scripts.py:44-48` | ✅ 契约 `test_prededuct_and_settle_refund[redis]` |
|
||||
| 3 | 并发 lease 带 TTL(防进程死亡泄漏)+ release/settle 幂等 | `limiter.py:46-63`、`tests/contracts_limiter.py:66-73` | ✅ 真实等待变体 `test_variant_lease_expiry_reclaims_slot` + 幂等契约 |
|
||||
| 4 | 窗口 id 用 Redis 服务器时钟(多进程口径统一) | `limiter.py:95-98` | ✅ `test_rpm_window_rollover_resets_quota` + 跨连接 RPM 用例 |
|
||||
| 5 | 跨进程熔断:单探针半开、探针租约 TTL、epoch fencing、force_open | `scripts.py:82-225`、`provider_gate.py` | ✅ `backends/redis/breaker.py`;熔断契约 + 7 个真实等待变体 |
|
||||
| 6 | 背压 stall 双条件判卡死(本地等够 + 全局无进展)+ mark_progress 全局活性 | `governance.py:270-281`、`limiter.py:193-209` | ✅ `middleware/retry.py` + `test_backpressure` 四象限 + 跨连接进度可见 |
|
||||
| 7 | 换源重试:失败跨源累计、退避取 Retry-After 较大值、源冷却备忘 | `governance.py:105-118, 244-261` | ✅ M1 已交付(`test_retry`);M2 联合验证下重验 |
|
||||
| 8 | 错误分类:429 body 细分 insufficient_quota、工件级失败不熔断 | `invokers.py:144-166`、`governance.py:237-239` | ✅ M1 已交付(`test_openai_compat`/ResultInvalid 不熔断) |
|
||||
| 9 | scope 级不可用结构化错误(reason + retry_after_s),供 worker 延期重投 | `errors.py:140-180`、`workers/tracking.py:406-428` | ✅ M1 错误模型(5 值 reason 勘误后)+ M2 增 stalled 真实触发路径 |
|
||||
| 10 | 流式三层活性看门狗;thinking token 刷活性不计结果 | `streaming.py`、`invokers.py:55-79` | ✅ M1 已交付(`test_streaming`) |
|
||||
|
||||
> 注: OCR invoker 能力(MonkeyOCR 两端点)不在本清单——归 M3。
|
||||
|
||||
## 2. 现状盘点(行数 wc -l 实测)
|
||||
|
||||
|
||||
@@ -134,7 +134,12 @@ class _RedisPermit:
|
||||
logger.warning("permit release 降级(租约将由 TTL 回收): {}", exc)
|
||||
|
||||
async def settle(self, actual_tokens: int) -> None:
|
||||
"""按实际 usage 结算 TPM 差额;释放侧失败降级 warning,不冒泡。"""
|
||||
"""按实际 usage 结算 TPM 差额;释放侧失败降级 warning,不冒泡。
|
||||
|
||||
有意偏离 CHS(limiter.py:62 成功后才置位): 本库释放侧失败被降级
|
||||
吞掉、调用方不重试,故尝试前置位——与内存版 flag 语义对齐,幂等
|
||||
性不受影响(2026-07-21 verifier M4 声明)。
|
||||
"""
|
||||
if self._settled:
|
||||
return
|
||||
self._settled = True
|
||||
|
||||
@@ -195,3 +195,50 @@ class TestInvariants:
|
||||
_row("c", session="r-p1"), # 非 P3 剔除
|
||||
]
|
||||
assert structured_success_rate(rows, session_suffix="-p3") == pytest.approx(0.5)
|
||||
|
||||
|
||||
class TestLiveInvariantsAndCaps:
|
||||
def test_capped_budget_clamps_to_signed_values(self):
|
||||
from tools.soak.scoreboard import capped_budget
|
||||
|
||||
assert capped_budget("P1", 9999) == 500
|
||||
assert capped_budget("P6", 100) == 100
|
||||
|
||||
async def test_gate_reenterable_detects_hung_probe(self):
|
||||
from tools.soak.scoreboard import inv_gate_reenterable
|
||||
|
||||
class _Decision:
|
||||
def __init__(self, allowed, state, is_probe=False):
|
||||
self.allowed = allowed
|
||||
self.state = state
|
||||
self.is_probe = is_probe
|
||||
self.retry_after_s = 30.0
|
||||
|
||||
class _Gate:
|
||||
def __init__(self, decision):
|
||||
self._d = decision
|
||||
self.released = 0
|
||||
|
||||
async def try_enter(self, name, owner):
|
||||
return self._d
|
||||
|
||||
async def release_probe(self, entry):
|
||||
self.released += 1
|
||||
|
||||
healthy = _Gate(_Decision(True, "closed"))
|
||||
await inv_gate_reenterable(healthy, ["s1"])
|
||||
probe_gate = _Gate(_Decision(True, "half_open", is_probe=True))
|
||||
await inv_gate_reenterable(probe_gate, ["s1"])
|
||||
assert probe_gate.released == 1 # 探针当场归还,不留新悬挂
|
||||
await inv_gate_reenterable(_Gate(_Decision(False, "open")), ["s1"]) # 冷却合法
|
||||
with pytest.raises(AssertionError, match="悬挂"):
|
||||
await inv_gate_reenterable(_Gate(_Decision(False, "half_open")), ["s1"])
|
||||
|
||||
def test_fault_errors_presence(self):
|
||||
from tools.soak.scoreboard import inv_fault_errors_present
|
||||
|
||||
rows = [_row("a", source="bad_1", error="SourceDeadError: 401"), _row("b")]
|
||||
inv_fault_errors_present(rows, fault_source_names=["bad_1"])
|
||||
inv_fault_errors_present(rows, fault_source_names=[]) # 未配故障源直接通过
|
||||
with pytest.raises(AssertionError):
|
||||
inv_fault_errors_present([_row("c")], fault_source_names=["bad_1"])
|
||||
|
||||
+68
-4
@@ -36,17 +36,33 @@ def _merged_env() -> dict[str, str]:
|
||||
return {k: v for k, v in {**dotenv_values(_ROOT / ".env"), **os.environ}.items() if v}
|
||||
|
||||
|
||||
def _guard(env: dict[str, str], workers: int) -> None:
|
||||
_SIGNED_MAX_CONCURRENCY = 100 # 网关保护签字值(设计 §8.1,2026-07-20 人类)
|
||||
_SIGNED_MAX_RPM = 600
|
||||
|
||||
|
||||
def _guard(env: dict[str, str], workers: int, scope: str) -> None:
|
||||
redis_url = env.get("REDIS_URL", "")
|
||||
if not redis_url.rstrip("/").endswith("/3"):
|
||||
raise SystemExit(f"拒跑: REDIS_URL 必须指向专用 db3,当前 {redis_url!r}")
|
||||
pg = env.get("PGW_TELEMETRY_PG_DSN", "")
|
||||
if pg and not pg.rstrip("/").endswith("/polygateway"):
|
||||
raise SystemExit(f"拒跑: PG DSN 必须指向 polygateway 专用库,当前库名不符")
|
||||
raise SystemExit("拒跑: PG DSN 必须指向 polygateway 专用库,当前库名不符")
|
||||
if workers > 1 and (
|
||||
env.get("PGW_LIMITER_BACKEND") != "redis" or env.get("PGW_BREAKER_BACKEND") != "redis"
|
||||
):
|
||||
raise SystemExit("拒跑: --workers>1 需要 PGW_LIMITER_BACKEND/PGW_BREAKER_BACKEND=redis")
|
||||
# 网关保护(签字值 100/600): 压测 scope 必须配全局闸且不超签字上限
|
||||
conc = env.get(f"{scope}__GLOBAL__MAX_CONCURRENCY")
|
||||
rpm = env.get(f"{scope}__GLOBAL__RPM")
|
||||
if not conc or not rpm:
|
||||
raise SystemExit(
|
||||
f"拒跑: 压测须配网关保护 {scope}__GLOBAL__MAX_CONCURRENCY(≤{_SIGNED_MAX_CONCURRENCY})"
|
||||
f" 与 {scope}__GLOBAL__RPM(≤{_SIGNED_MAX_RPM})——设计 §8.1 签字值"
|
||||
)
|
||||
if int(conc) > _SIGNED_MAX_CONCURRENCY or int(rpm) > _SIGNED_MAX_RPM:
|
||||
raise SystemExit(
|
||||
f"拒跑: 全局闸 {conc}/{rpm} 超签字上限 {_SIGNED_MAX_CONCURRENCY}/{_SIGNED_MAX_RPM}"
|
||||
)
|
||||
|
||||
|
||||
async def _flush_db3(redis_url: str) -> None:
|
||||
@@ -86,8 +102,10 @@ async def _worker_async(args: argparse.Namespace, worker_idx: int) -> None:
|
||||
sem = asyncio.Semaphore(args.concurrency)
|
||||
stats = {"calls": 0, "ok": 0, "failed": 0, "cancelled": 0, "tokens": 0}
|
||||
rss_samples = [_rss_mb()]
|
||||
from tools.soak.scoreboard import capped_budget
|
||||
|
||||
deadline = time.monotonic() + args.max_hours * 3600
|
||||
budget_calls = args.budget_calls // args.workers
|
||||
budget_calls = capped_budget(args.scenario, args.budget_calls) // args.workers
|
||||
budget_tokens = min(args.budget_tokens, TOKEN_HARD_CAP) // args.workers
|
||||
inflight: set[asyncio.Task] = set()
|
||||
|
||||
@@ -175,6 +193,17 @@ def _scoreboard(args: argparse.Namespace, env: dict[str, str]) -> None:
|
||||
_check("RPM 从未击穿(分钟桶)", sb.inv_rpm_never_exceeded, rows, rpm_conf)
|
||||
for r in results:
|
||||
_check(f"RSS 平稳(w)", sb.inv_rss_stable, r["rss_mb"], max_growth_mb=args.max_rss_growth_mb)
|
||||
# 不变量 1(记账归零 + gate 可再准入): 活后端检查,仅 redis 后端可跨进程复查
|
||||
if env.get("PGW_LIMITER_BACKEND") == "redis":
|
||||
asyncio.run(_live_checks(args, env, verdicts))
|
||||
else:
|
||||
verdicts.append(("记账归零/gate 可再准入", "SKIP — memory 后端跨进程不可查"))
|
||||
# 不变量 2c: P5/P6 且配置了故障源 → 故障必须真实发生
|
||||
if args.scenario in ("P5", "P6"):
|
||||
fault_names = [
|
||||
name for name in rpm_conf if rpm_conf[name] <= 10
|
||||
] # 紧闸源;坏 key/黑洞源由错误分布人工核对(findings §4 条 2 校准留待 P5 实跑)
|
||||
_check("故障混编生效(P5/P6)", sb.inv_fault_errors_present, rows, fault_source_names=fault_names)
|
||||
rate = sb.structured_success_rate(rows, session_suffix="-p3")
|
||||
verdicts.append(("P3 结构化成功率", f"{rate:.3f}(基线首跑建立)"))
|
||||
report = sb.render_report(args.run_id, rows, verdicts)
|
||||
@@ -185,6 +214,41 @@ def _scoreboard(args: argparse.Namespace, env: dict[str, str]) -> None:
|
||||
raise SystemExit("硬不变量被击穿,见报告")
|
||||
|
||||
|
||||
async def _live_checks(
|
||||
args: argparse.Namespace, env: dict[str, str], verdicts: list[tuple[str, str]]
|
||||
) -> None:
|
||||
"""不变量 1 活后端复查(跑后): 限流 inflight 归零 + 熔断门可再准入。"""
|
||||
from polygateway.backends.redis.breaker import RedisGate
|
||||
from polygateway.backends.redis.limiter import RedisLimiter
|
||||
from polygateway.config import GatewaySettings
|
||||
|
||||
from tools.soak import scoreboard as sb
|
||||
|
||||
settings = GatewaySettings.from_env(args.scope, env=env)
|
||||
names = [s.name for s in settings.sources]
|
||||
limiter = RedisLimiter.from_url(
|
||||
env["REDIS_URL"],
|
||||
scope=settings.scope,
|
||||
sources={s.name: s for s in settings.sources},
|
||||
global_limits=settings.global_limits,
|
||||
lease_ttl_s=settings.lease_ttl_s,
|
||||
)
|
||||
gate = RedisGate.from_url(env["REDIS_URL"], config=settings.breaker, scope=settings.scope)
|
||||
try:
|
||||
for name, invariant in (
|
||||
("记账归零(inflight)", sb.inv_accounting_zeroed(limiter, names)),
|
||||
("gate 可再准入(探针不悬挂)", sb.inv_gate_reenterable(gate, names)),
|
||||
):
|
||||
try:
|
||||
await invariant
|
||||
verdicts.append((name, "PASS"))
|
||||
except AssertionError as exc:
|
||||
verdicts.append((name, f"FAIL — {exc}"))
|
||||
finally:
|
||||
await limiter.aclose()
|
||||
await gate.aclose()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description="PolyGateway 真实数据压测")
|
||||
parser.add_argument("--scenario", required=True, choices=["P1", "P2", "P3", "P4", "P5", "P6"])
|
||||
@@ -200,7 +264,7 @@ def main() -> None:
|
||||
if args.run_id is None:
|
||||
args.run_id = time.strftime("soak_%Y%m%d_%H%M%S")
|
||||
env = _merged_env()
|
||||
_guard(env, args.workers)
|
||||
_guard(env, args.workers, args.scope)
|
||||
print(f"run_id={args.run_id}: FLUSHDB db3 + namespace 隔离")
|
||||
asyncio.run(_flush_db3(env["REDIS_URL"]))
|
||||
if args.workers == 1:
|
||||
|
||||
@@ -131,3 +131,51 @@ def write_report(run_id: str, content: str, out_dir: Path | str = "tests/outputs
|
||||
path = out / f"{run_id}.md"
|
||||
path.write_text(content, encoding="utf-8")
|
||||
return path
|
||||
|
||||
|
||||
# —— 活后端检查与预算帽(2026-07-21 verifier I1/M6/M7 补齐)——
|
||||
|
||||
SCENARIO_CALL_CAPS = { # 设计 §8.1 签字的单场景请求数上限
|
||||
"P1": 500,
|
||||
"P2": 450,
|
||||
"P3": 2500,
|
||||
"P4": 2500,
|
||||
"P5": 1000,
|
||||
"P6": 8000,
|
||||
}
|
||||
|
||||
|
||||
def capped_budget(scenario: str, requested_calls: int) -> int:
|
||||
"""预算帽: 请求数不越过签字上限(超出取上限并由调用方打印告知)。"""
|
||||
cap = SCENARIO_CALL_CAPS[scenario]
|
||||
return min(requested_calls, cap)
|
||||
|
||||
|
||||
async def inv_gate_reenterable(gate, sources: list[str]) -> None:
|
||||
"""不变量 1b: 跑后熔断门可再准入,探针不悬挂。
|
||||
|
||||
run 结束后已无 in-flight 调用,若某源仍处 HALF_OPEN 拒入 = 死探针
|
||||
悬挂(只能等 TTL);OPEN 冷却中属故障源的合法状态,不算击穿。
|
||||
拿到的探针当场归还(release_probe),不留新悬挂。
|
||||
"""
|
||||
for name in sources:
|
||||
decision = await gate.try_enter(name, "scoreboard-probe")
|
||||
if decision.allowed:
|
||||
if decision.is_probe:
|
||||
await gate.release_probe(decision)
|
||||
continue
|
||||
assert str(decision.state) != "half_open", (
|
||||
f"源 {name} 跑后仍 HALF_OPEN 拒入(探针悬挂,retry_after={decision.retry_after_s:.1f}s)"
|
||||
)
|
||||
|
||||
|
||||
def inv_fault_errors_present(rows: list[Row], *, fault_source_names: list[str]) -> None:
|
||||
"""不变量 2c(P5/P6): 配置了故障源则错误必然出现且落在故障源上。
|
||||
|
||||
注入"比例"的精确吻合依赖具体混编配置,自动断言留待 P5 实跑校准
|
||||
(findings §4 条 2);此处先钉存在性: 故障源零错误 = 故障根本没被打到。
|
||||
"""
|
||||
if not fault_source_names:
|
||||
return
|
||||
fault_errors = [r for r in rows if r.get("error") and r["source_name"] in fault_source_names]
|
||||
assert fault_errors, f"故障源 {fault_source_names} 零错误行——故障混编未生效"
|
||||
|
||||
Reference in New Issue
Block a user