fix: make RPM jitter exemption side-aware with one-shot room

Verifier adversarial cases showed pooled-room exemption could hide real
breaches: edge rows may only borrow from the neighbor on their own side,
and neighbor room is consumed globally so two windows cannot claim the
same slot. Also parse SQLite created_at as UTC, guard the dispatch
semaphore on progress-callback failure, and print a caveat that rescore
live checks reflect current db3 state.
This commit is contained in:
2026-07-21 07:42:54 -04:00
parent 8c178a249c
commit 5dfd88c4d6
3 changed files with 87 additions and 24 deletions
+31
View File
@@ -237,6 +237,37 @@ class TestInvariants:
] ]
inv_rpm_never_exceeded(rows, {"s1": 5}) inv_rpm_never_exceeded(rows, {"s1": 5})
def test_rpm_jitter_is_side_aware(self):
# verifier 对抗样例①: 左贴边行只能借左邻余量——左邻已满、右邻
# 全空时,7 行(2 行贴左界)仍是真实击穿,不得借右邻豁免
rows = [
_row(f"L{i}", created_at=f"2026-07-20T10:00:{s:02d}")
for i, s in enumerate([10, 20, 30, 40, 50]) # 左邻(10:00)满额 5
]
rows += [
_row(f"M{i}", created_at=f"2026-07-20T10:01:{s:02d}")
for i, s in enumerate([0, 1, 10, 20, 30, 40, 50]) # 本窗 7 行,2 行贴左界
]
with pytest.raises(AssertionError):
inv_rpm_never_exceeded(rows, {"s1": 5})
def test_rpm_jitter_room_not_double_claimed(self):
# verifier 对抗样例②: 两个超限窗不得重复认领中间窗的同 1 个余量
rows = [
_row(f"A{i}", created_at=f"2026-07-20T10:00:{s:02d}")
for i, s in enumerate([10, 20, 30, 40, 50, 59]) # 6 行,1 行贴右界
]
rows += [
_row(f"B{i}", created_at=f"2026-07-20T10:01:{s:02d}")
for i, s in enumerate([10, 20, 30, 40]) # 中间窗 4 行,仅 1 余量
]
rows += [
_row(f"C{i}", created_at=f"2026-07-20T10:02:{s:02d}")
for i, s in enumerate([1, 10, 20, 30, 40, 50]) # 6 行,1 行贴左界
]
with pytest.raises(AssertionError):
inv_rpm_never_exceeded(rows, {"s1": 5})
def test_rss_stable(self): def test_rss_stable(self):
inv_rss_stable([100.0, 105.0, 110.0], max_growth_mb=50.0) inv_rss_stable([100.0, 105.0, 110.0], max_growth_mb=50.0)
with pytest.raises(AssertionError): with pytest.raises(AssertionError):
+12 -3
View File
@@ -87,7 +87,9 @@ def _rss_mb() -> float:
return int(out.stdout.strip()) / 1024.0 return int(out.stdout.strip()) / 1024.0
async def _paced_dispatch(generator, *, sem, spawn, should_stop, on_dispatched) -> set[asyncio.Task]: async def _paced_dispatch(
generator, *, sem, spawn, should_stop, on_dispatched
) -> set[asyncio.Task]:
"""有界分发: 先占并发名额再建任务,名额由任务收尾释放。 """有界分发: 先占并发名额再建任务,名额由任务收尾释放。
2026-07-21 P6 教训: 无界 create_task 曾在 15s 内入队全部预算, 2026-07-21 P6 教训: 无界 create_task 曾在 15s 内入队全部预算,
@@ -105,8 +107,12 @@ async def _paced_dispatch(generator, *, sem, spawn, should_stop, on_dispatched)
if should_stop(): if should_stop():
break break
await sem.acquire() await sem.acquire()
on_dispatched() try:
task = asyncio.create_task(_run(spawn(kind, kwargs))) on_dispatched()
task = asyncio.create_task(_run(spawn(kind, kwargs)))
except BaseException:
sem.release() # 名额已占而任务未建(如进度采样失败),不泄漏
raise
inflight.add(task) inflight.add(task)
task.add_done_callback(inflight.discard) task.add_done_callback(inflight.discard)
return inflight return inflight
@@ -336,6 +342,9 @@ def main() -> None:
raise SystemExit(f"拒跑: 找不到 data/soak/result_{args.run_id}_*.json") raise SystemExit(f"拒跑: 找不到 data/soak/result_{args.run_id}_*.json")
args.workers = len(found) args.workers = len(found)
_guard(env, args.workers, args.scope) _guard(env, args.workers, args.scope)
print(
"rescore: 记账归零/gate 两项活检查反映**当前** db3 状态;若 db3 已被后续 run 复用请忽略"
)
_scoreboard(args, env) _scoreboard(args, env)
return return
if args.budget_calls is None or args.budget_tokens is None: if args.budget_calls is None or args.budget_tokens is None:
+44 -21
View File
@@ -9,7 +9,7 @@ from __future__ import annotations
import sqlite3 import sqlite3
from collections import Counter from collections import Counter
from datetime import datetime from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -47,22 +47,50 @@ def inv_call_ids_unique(rows: list[Row]) -> None:
def _admit_second(row: Row, clock_offset_s: float) -> float: def _admit_second(row: Row, clock_offset_s: float) -> float:
"""还原准入时刻(服务器钟): created_at 是完成落库时刻,减去调用延迟。""" """还原准入时刻(服务器钟): created_at 是完成落库时刻,减去调用延迟。"""
done = datetime.fromisoformat(str(row["created_at"])).timestamp() dt = datetime.fromisoformat(str(row["created_at"]))
if dt.tzinfo is None:
# SQLite datetime('now') 落库为 UTC naive;按本机时区解析会错位整时
dt = dt.replace(tzinfo=timezone.utc)
latency_ms = row.get("latency_ms") or 0 latency_ms = row.get("latency_ms") or 0
return done - float(latency_ms) / 1000.0 + clock_offset_s return dt.timestamp() - float(latency_ms) / 1000.0 + clock_offset_s
def _movable_edge_rows( def _edge_counts(admits: list[float], minute: int, slack_s: float) -> tuple[int, int]:
admits: list[float], buckets: Counter[int], minute: int, limit: int, slack_s: float """窗口内贴左界/贴右界(± slack_s)的行数。"""
) -> int: left = right = 0
"""超限窗口内可归邻窗的贴边行数(受邻窗余量约束)。""" for a in admits:
near_edge = sum( if int(a // 60) != minute:
1 continue
for a in admits if a - minute * 60 <= slack_s:
if int(a // 60) == minute and min(a - minute * 60, (minute + 1) * 60 - a) <= slack_s left += 1
) elif (minute + 1) * 60 - a <= slack_s:
room = max(0, limit - buckets.get(minute - 1, 0)) + max(0, limit - buckets.get(minute + 1, 0)) right += 1
return min(near_edge, room) 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( def inv_rpm_never_exceeded(
@@ -89,13 +117,8 @@ def inv_rpm_never_exceeded(
breaches: dict[tuple[str, int], int] = {} breaches: dict[tuple[str, int], int] = {}
for source, admits in per_source.items(): for source, admits in per_source.items():
limit = per_source_rpm[source] limit = per_source_rpm[source]
buckets = Counter(int(a // 60) for a in admits) for minute, n in _jitter_breaches(admits, limit, boundary_slack_s).items():
for minute, n in buckets.items(): breaches[(source, minute)] = n
if n <= limit:
continue
movable = _movable_edge_rows(admits, buckets, minute, limit, boundary_slack_s)
if n - movable > limit:
breaches[(source, minute)] = n
assert not breaches, f"RPM 击穿(准入时刻+服务器钟口径): {dict(list(breaches.items())[:5])}" assert not breaches, f"RPM 击穿(准入时刻+服务器钟口径): {dict(list(breaches.items())[:5])}"