"""压测记分板: 六条硬不变量断言 + 报告产出(findings §4;吃自己遥测的狗粮)。 输入口径三类(M2 计划 T11): 不变量 2/3/5/6 = 纯函数(遥测行迭代器); 不变量 1(记账归零)= 活 Redis/后端检查(async);不变量 4(RSS)= 纯函数 (run_soak 周期采样 json)。断言函数抛 AssertionError 即不变量被击穿。 """ from __future__ import annotations import sqlite3 from collections import Counter from datetime import UTC, datetime from pathlib import Path from typing import Any Row = dict[str, Any] def load_rows(*db_paths: Path | str) -> list[Row]: """合并全部 worker 遥测库的 llm_calls 行(dict 形态)。""" rows: list[Row] = [] for path in db_paths: conn = sqlite3.connect(path) conn.row_factory = sqlite3.Row try: rows.extend(dict(r) for r in conn.execute("SELECT * FROM llm_calls")) finally: conn.close() return rows # —— 硬不变量(纯函数) —— def inv_rows_match_calls(rows: list[Row], *, expected_calls: int, tolerance: int = 0) -> None: """不变量 2a: 遥测行数 == 发出的请求数(±取消双记的已知语义容差)。""" assert abs(len(rows) - expected_calls) <= tolerance, ( f"遥测行数 {len(rows)} ≠ 请求数 {expected_calls}(容差 {tolerance})" ) def inv_call_ids_unique(rows: list[Row]) -> None: """不变量 2b: call_id 无重复。""" dupes = [cid for cid, n in Counter(r["call_id"] for r in rows).items() if n > 1] assert not dupes, f"call_id 重复: {dupes[:5]}" def _admit_second(row: Row, clock_offset_s: float) -> float: """还原准入时刻(服务器钟): created_at 是完成落库时刻,减去调用延迟。""" dt = datetime.fromisoformat(str(row["created_at"])) if dt.tzinfo is None: # SQLite datetime('now') 落库为 UTC naive;按本机时区解析会错位整时 dt = dt.replace(tzinfo=UTC) latency_ms = row.get("latency_ms") or 0 return dt.timestamp() - float(latency_ms) / 1000.0 + clock_offset_s def _edge_counts(admits: list[float], minute: int, slack_s: float) -> tuple[int, int]: """窗口内贴左界/贴右界(± slack_s)的行数。""" left = right = 0 for a in admits: if int(a // 60) != minute: continue if a - minute * 60 <= slack_s: left += 1 elif (minute + 1) * 60 - a <= slack_s: right += 1 return left, right def _jitter_breaches(admits: list[float], limit: int, slack_s: float) -> dict[int, int]: """贴边豁免的全局结算: 分侧借邻窗余量,余量一次性消耗不可重复认领。 verifier 对抗样例(2026-07-21)证明"合计余量 + 逐窗独立豁免"会漏报: 左贴边行只可能属于左邻窗(反之亦然),且同一邻窗余量只能被认领一次。 """ buckets = Counter(int(a // 60) for a in admits) consumed: Counter[int] = Counter() breaches: dict[int, int] = {} for minute in sorted(buckets): n = buckets[minute] if n <= limit: continue left_edge, right_edge = _edge_counts(admits, minute, slack_s) need = n - limit for edge, neighbor in ((left_edge, minute - 1), (right_edge, minute + 1)): room = max(0, limit - buckets.get(neighbor, 0)) - consumed[neighbor] take = min(edge, max(0, room), need) consumed[neighbor] += take need -= take if need > 0: breaches[minute] = n return breaches def inv_rpm_never_exceeded( rows: list[Row], per_source_rpm: dict[str, int], *, clock_offset_s: float = 0.0, boundary_slack_s: float = 2.0, ) -> None: """不变量 3: 任一限流器分钟窗口内单源准入数 ≤ RPM 配置。 口径与限流器同源(`backends/redis/limiter.py` `_window_id`): 窗口 = **Redis 服务器钟**的固定分钟;准入时刻 = created_at(完成落库)− latency。 滑动 60s 窗会对"窗尾+窗头"的合法背靠背流量误报,不采用(findings §4 条 3)。 created_at 秒级截断给准入时刻 ±秒级噪声,距窗口边界 ≤ boundary_slack_s 的 行允许归入有余量的邻窗。缓存命中行不计: 缓存在限流闸之前返回,未耗名额 也未打网关(遥测必录使其带原源名落库)。2026-07-21 P6 教训: 本机钟聚桶 + 完成时刻口径 + 计入缓存行,三重口径偏差曾把合规流量误判为击穿。 """ per_source: dict[str, list[float]] = {} for r in rows: if per_source_rpm.get(r["source_name"], 0) > 0 and not r.get("cache_hit"): per_source.setdefault(r["source_name"], []).append(_admit_second(r, clock_offset_s)) breaches: dict[tuple[str, int], int] = {} for source, admits in per_source.items(): limit = per_source_rpm[source] for minute, n in _jitter_breaches(admits, limit, boundary_slack_s).items(): breaches[(source, minute)] = n assert not breaches, f"RPM 击穿(准入时刻+服务器钟口径): {dict(list(breaches.items())[:5])}" def inv_rss_stable(samples_mb: list[float], *, max_growth_mb: float) -> None: """不变量 4: 进程 RSS 首末差 < 阈值(窗口计数器/租约表无无界增长)。""" assert len(samples_mb) >= 2, "RSS 采样不足" growth = samples_mb[-1] - samples_mb[0] assert growth < max_growth_mb, f"RSS 增长 {growth:.1f}MB ≥ 阈值 {max_growth_mb}MB" def structured_success_rate(rows: list[Row], *, session_suffix: str) -> float: """不变量 5 的量: 指定场景行的最终成功率(error 为空即成功)。""" scoped = [r for r in rows if str(r.get("session_id") or "").endswith(session_suffix)] if not scoped: return 0.0 ok = sum(1 for r in scoped if not r.get("error")) return ok / len(scoped) async def inv_accounting_zeroed(limiter, sources: list[str]) -> None: """不变量 1: 结束后全部源 inflight == 0(无泄漏租约)。活后端检查。""" for name in sources: stats = await limiter.source_stats(name) assert stats.inflight == 0, f"源 {name} 结束后 inflight={stats.inflight}(租约泄漏)" # —— 报告 —— def _percentile(values: list[float], q: float) -> float: if not values: return 0.0 ordered = sorted(values) return ordered[min(int(len(ordered) * q), len(ordered) - 1)] def render_report(run_id: str, rows: list[Row], verdicts: list[tuple[str, str]]) -> str: """结构化 Markdown 报告(不变量 6: 成本/延迟基线产出)。""" errors = Counter(str(r["error"]).split(":")[0] for r in rows if r.get("error")) latencies = [float(r["latency_ms"]) for r in rows if not r.get("error")] ttfts = [float(r["ttft_ms"]) for r in rows if r.get("ttft_ms") is not None] costs = [float(r["cost"]) for r in rows if r.get("cost") is not None] tokens = sum(int(r["prompt_tokens"]) + int(r["completion_tokens"]) for r in rows) cache_hits = sum(1 for r in rows if r.get("cache_hit")) lines = [ f"# Soak 报告: {run_id}", "", "## 不变量裁决", *(f"- {name}: {verdict}" for name, verdict in verdicts), "", "## 规模", f"- 遥测行数: {len(rows)};缓存命中: {cache_hits};总 token: {tokens}", f"- 错误分布: {dict(errors) or '无'}", "", "## 延迟基线(成功行)", f"- 总时长 p50/p95: {_percentile(latencies, 0.5):.0f} / {_percentile(latencies, 0.95):.0f} ms", f"- TTFT p50/p95: {_percentile(ttfts, 0.5):.0f} / {_percentile(ttfts, 0.95):.0f} ms", "", "## 成本", f"- 累计 cost: {sum(costs):.4f}(有单价行 {len(costs)}/{len(rows)})", ] return "\n".join(lines) + "\n" def write_report(run_id: str, content: str, out_dir: Path | str = "tests/outputs/soak") -> Path: out = Path(out_dir) out.mkdir(parents=True, exist_ok=True) 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 签字的单场景请求数上限(P7: M3 计划 T8) "P1": 500, "P2": 450, "P3": 2500, "P4": 2500, "P5": 1000, "P6": 8000, "P7": 1500, } 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} 零错误行——故障混编未生效" def inv_success_rate(ok: int, total: int, *, min_rate: float) -> None: """P7 不变量①: 调用级成功率(worker stats 口径,与 M2.5 验收同源)。""" assert total > 0, "零调用无法评估成功率" rate = ok / total assert rate >= min_rate, f"成功率 {rate:.4f}({ok}/{total})< 阈值 {min_rate}" def inv_fault_share(rows: list[Row], fault_source_names: list[str], *, max_share: float) -> None: """P7 不变量②: 坏源尝试占比受健康选源压制(排除缓存行)。""" attempts = [r for r in rows if not r.get("cache_hit")] assert attempts, "无尝试行" fault = sum(1 for r in attempts if r["source_name"] in fault_source_names) share = fault / len(attempts) assert share <= max_share, ( f"坏源尝试占比 {share:.3f}({fault}/{len(attempts)})> 阈值 {max_share}" ) def inv_errors_classified(rows: list[Row], known_prefixes: tuple[str, ...]) -> None: """P7 不变量⑦/⑧: 一切错误行可归入已知异常类名前缀(零未分类异常)。""" unknown = Counter( str(r["error"]).split(":")[0] for r in rows if r.get("error") and not str(r["error"]).startswith(known_prefixes) ) assert not unknown, f"未分类错误前缀: {dict(unknown)}" def inv_rss_absolute(samples_mb: list[float], *, max_mb: float) -> None: """P7 不变量⑤(计划 T8): 全程 RSS 绝对值有界(与增长口径互补)。""" assert samples_mb, "RSS 采样为空" peak = max(samples_mb) assert peak < max_mb, f"RSS 峰值 {peak:.1f}MB ≥ 绝对上限 {max_mb}MB" async def inv_healthy_gates_closed(gate, healthy_names: list[str]) -> None: """P7 不变量③(计划 T8): 跑后真源熔断门 CLOSED 可准入 = 真源零误熔的 机械化终态判据(误熔必由该源失败计数驱动,配合遥测"真源零错误行"覆盖全程)。 拿到探针当场归还,不留悬挂。""" for name in healthy_names: decision = await gate.try_enter(name, "scoreboard-healthy-probe") if decision.allowed and decision.is_probe: await gate.release_probe(decision) assert decision.allowed and str(decision.state) == "closed", ( f"真源 {name} 跑后门态 {decision.state}(allowed={decision.allowed})——疑似误熔" ) def inv_any_errors(rows: list[Row]) -> None: """不变量 2c 的兜底形态(P5/P6): 故障源混编池下错误行必然存在。 紧闸源被限流闸跳过不产错误行,故障名单无法从配置泛化推断——按源归因 交报告"错误分布"人工核对;此处只钉全局存在性。 """ assert any(r.get("error") for r in rows), "P5/P6 故障混编池零错误行——故障未生效"