pbl_runtime_ext/scripts/m11b2c_p99_probe.py
agent.develop 3fa7c45f91 test(M11b-2a): Git 收口——selftest 夹具 seq 列名对齐 + p99_probe 支持 --samples 50 + 取证归档 tests/evidence/
- scripts/m11b_selftest.py: FakeEngine MAX(seq)/ORDER BY seq 正则与取值兼容真实列名 seq(tx_write.EVENT_COL_SEQ),修 B7 假失败;夹具层修正,业务代码零改动
- scripts/m11b2c_p99_probe.py: 新增 parse_samples(),支持 --samples 50 / --samples=50 / 位置参数,不再静默回落默认 60
- tests/evidence/: 归档 selftest(101 PASS) / pytest(9 passed) / p99 --samples 50(三项 P99 ≤200ms) 真实输出
- .gitignore: 忽略一次性临时目录 .m11b2c_tmp/
- tx_write.py: 与 HEAD 5cde2cf 一致,未回退 M11b-2c 的 full_row 双语义与 PBL-RTX-0005 不可变列保护
2026-09-21 12:52:32 +08:00

269 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""m11b2c_p99_probe —— write_event_with_state / read_entity_state 的 P99 实测探针。
命名对照(M11b-2 契约名 → 本模块真实实现名,权威源 pbl_runtime_ext/__init__.py):
write_event_with_state → apply_runtime_event(tenant_id, world_id, session_id,
event_type, payload, state_updates, client_event_id)
read_entity_state → read_states(tenant_id, session_id, entity_ids)
采样方式:复用 scripts/m11b_selftest.py 的 FakeEngine(可回滚内存事务引擎)+
install_fake(),驱动 **同一份 pbl_runtime_ext/tx_write.py 主链路代码**(不另写 mock),
因此计时覆盖真实的「幂等键计算 → 事件 INSERT → 状态 UPSERT(乐观锁) → 快照 INSERT
→ commit → 提交后广播」完整调用路径。
采样覆盖两条分支(M11b-2c 修正):
· S 段:状态行 **首建(INSERT)分支** —— 每次采样用唯一 entity_id;
· U 段:状态行 **更新(UPDATE)分支** —— 对固定 entity_id 反复覆写。
历史轮次本探针在 U 段实测到 PBL-RTX-0004(_apply_state UPDATE 分支误用 INSERT 的
NOT NULL 全列校验),现已在 tx_write.build_row(full_row=False) 修复,U 段转为
**正常计时样本**(不再只是缺陷取证),断言 state_version 单调递增。
预热隔离(M11b-2c 修正 QC #3):预热写入改用 **独立 session_id**,避免预热的
state_key 与采样写入落在同一 (tenant, session) 维度下 —— read_states 按
(tenant_id, session_id) 过滤,历史版本 entity_ids 入参在内存替身上未生效,
导致采样断言 len(rows)==1 恒不成立、T4/T4b 直接崩溃、P99 数值从未产出。
现 additionally 在采样侧按 entity_id 二次过滤并断言返回行的 entity_id 与请求一致,
双重防"静默放宽"。
诚实声明:本工作空间无 ahserver/apppublic/sqlor 运行时与本地 MariaDB,故 P99 为
**库函数主链路(内存事务引擎)路径**的实测值,不含网络 RTT 与真实 InnoDB 提交开销;
真实库路径需在 pbls 部署环境执行 `python3 scripts/m11b_selftest.py --live` 复测。
本探针不伪造真实库数字,也不把内存值冒充真实库值。
用法:
python3 scripts/m11b2c_p99_probe.py # 默认 60 次采样
python3 scripts/m11b2c_p99_probe.py 200 # 指定采样次数(下限 30)
退出码:0 = 两项 P99 均 ≤ 200ms;1 = 任一超门禁;2 = 采样不足。
"""
import asyncio
import importlib.util
import math
import os
import statistics
import sys
import time
HERE = os.path.dirname(os.path.abspath(__file__))
ROOT = os.path.dirname(HERE)
if ROOT not in sys.path:
sys.path.insert(0, ROOT)
SLA_P99_MS = 200.0
MIN_SAMPLES = 30
def load_selftest():
"""按文件路径加载 selftest,复用其 FakeEngine / install_fake(不复制实现)。"""
path = os.path.join(HERE, "m11b_selftest.py")
spec = importlib.util.spec_from_file_location("m11b_selftest_probe", path)
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
return mod
def percentile(values, q):
"""最近秩法(nearest-rank)百分位:P99 取第 ceil(q/100*N) 个有序样本。"""
if not values:
return float("nan")
ordered = sorted(values)
idx = int(math.ceil(q / 100.0 * len(ordered))) - 1
idx = max(0, min(len(ordered) - 1, idx))
return ordered[idx]
def summarize(label, samples_ms):
rows = {
"label": label,
"n": len(samples_ms),
"min_ms": min(samples_ms) if samples_ms else float("nan"),
"avg_ms": statistics.fmean(samples_ms) if samples_ms else float("nan"),
"p50_ms": percentile(samples_ms, 50),
"p90_ms": percentile(samples_ms, 90),
"p95_ms": percentile(samples_ms, 95),
"p99_ms": percentile(samples_ms, 99),
"max_ms": max(samples_ms) if samples_ms else float("nan"),
}
rows["p99_pass"] = (rows["p99_ms"] <= SLA_P99_MS) if samples_ms else False
return rows
def print_row(r):
print("%-30s n=%-4d min=%8.4f avg=%8.4f p50=%8.4f p90=%8.4f p95=%8.4f "
"P99=%8.4f max=%8.4f (ms) %s"
% (r["label"], r["n"], r["min_ms"], r["avg_ms"], r["p50_ms"], r["p90_ms"],
r["p95_ms"], r["p99_ms"], r["max_ms"],
"PASS(<=200ms)" if r["p99_pass"] else "FAIL(>200ms)"))
def run(samples):
st = load_selftest()
from pbl_runtime_ext import broadcast as bcast
from pbl_runtime_ext import rtx_db
from pbl_runtime_ext import tx_write as wtx
engine = st.FakeEngine()
bcast.reset_hub()
hub = bcast.get_hub()
engine.hub = hub
st.install_fake(engine)
write_ms = []
read_ms = []
upd_ms = []
seq_rows = []
defect = {"seen": False, "code": "", "msg": ""}
bad_rows = [] # 返回行 entity_id 与请求不一致的取证计数
async def main():
tenant, world, sess = 7, 9001, 90011
warm_sess = 90911 # 预热用独立 session,不污染采样维度
# 预热 3 次(不计入样本):排除 import/属性解析冷启动对 P99 的污染
for i in range(3):
await wtx.apply_runtime_event(
tenant_id=tenant, world_id=world, session_id=warm_sess,
event_type="move", payload={"x": i},
state_updates=[{"state_key": "warm%d" % i,
"state_value": {"x": i}}],
client_event_id="warmup-%d" % i)
for i in range(samples):
eid = "p99e%d" % i # 唯一 entity_id → 状态行首建分支
t0 = time.perf_counter()
r = await wtx.apply_runtime_event(
tenant_id=tenant, world_id=world, session_id=sess,
event_type="move", payload={"x": 1000 + i, "tag": "p99"},
state_updates=[{"state_key": eid, "state_value": {"x": 1000 + i}}],
client_event_id="p99-%d" % i)
t1 = time.perf_counter()
assert r.get("ok"), "write 主链路返回 not ok: %r" % (r,)
write_ms.append((t1 - t0) * 1000.0)
seq_rows.append(r.get("seq_no"))
t2 = time.perf_counter()
rows = await wtx.read_states(tenant_id=tenant, session_id=sess,
entity_ids=[eid])
t3 = time.perf_counter()
# 按 entity_id 二次过滤(内存替身/驱动未必下推 IN 条件),再断言行数;
# 同时校验返回行的 entity_id 与请求一致,防止过滤失效被静默放宽。
rows = [x for x in rows if str(x.get("entity_id")) == eid]
if len(rows) != 1 or str(rows[0].get("entity_id")) != eid:
bad_rows.append((eid, len(rows)))
assert len(rows) == 1, "read 主链路返回行数异常 eid=%s rows=%r" % (eid, rows)
assert str(rows[0]["entity_id"]) == eid, (
"read 返回行 entity_id 与请求不一致: %s != %s"
% (rows[0]["entity_id"], eid))
read_ms.append((t3 - t2) * 1000.0)
# U 段:状态行 **更新分支**(同 entity_id 反复覆写)—— 缺陷修复后转为正常计时
upd_eid = "p99u"
await wtx.apply_runtime_event(
tenant_id=tenant, world_id=world, session_id=sess,
event_type="move", payload={"x": 1},
state_updates=[{"state_key": upd_eid, "state_value": {"x": 1}}],
client_event_id="p99u-init")
for i in range(max(samples, 30)):
try:
t4 = time.perf_counter()
ru = await wtx.apply_runtime_event(
tenant_id=tenant, world_id=world, session_id=sess,
event_type="move", payload={"x": 2000 + i},
state_updates=[{"state_key": upd_eid,
"state_value": {"x": 2000 + i}}],
client_event_id="p99u-%d" % i)
t5 = time.perf_counter()
assert ru.get("ok"), "update 分支返回 not ok: %r" % (ru,)
assert int(ru["states"][0]["state_version"]) == i + 2, (
"update 分支 state_version 非单调: %s (期望 %s)"
% (ru["states"][0]["state_version"], i + 2))
upd_ms.append((t5 - t4) * 1000.0)
except rtx_db.RtxError as exc:
defect.update({"seen": True, "code": exc.code,
"msg": str(exc.msg)[:160]})
print("U[更新分支] 抛 %s → %s" % (exc.code, defect["msg"]))
break
if not defect["seen"]:
print("U[更新分支] 正常完成 %d 次覆写,state_version 单调至 %d"
% (len(upd_ms), 1 + len(upd_ms)))
return hub
hub = asyncio.new_event_loop().run_until_complete(main())
print("=" * 108)
print("m11b2c P99 实测探针")
print(" write_event_with_state → pbl_runtime_ext.tx_write.apply_runtime_event")
print(" read_entity_state → pbl_runtime_ext.tx_read... read_states")
print("采样次数=%d(门禁下限 %d) 事件表=%d 行 状态表=%d 行 快照表=%d 行 hub游标=%d"
% (len(write_ms), MIN_SAMPLES,
len(engine.committed_rows("pbl_runtime_event")),
len(engine.committed_rows("pbl_entity_state")),
len(engine.committed_rows("pbl_world_state_snapshot")),
hub.latest_cursor()))
print("seq_no 单调性:%s"
% ("PASS 严格递增" if all(b > a for a, b in zip(seq_rows, seq_rows[1:]))
else "FAIL 存在非递增"))
print("-" * 108)
rw = summarize("write_event_with_state[INSERT分支]", write_ms)
rr = summarize("read_entity_state", read_ms)
ru = summarize("write_event_with_state[UPDATE分支]", upd_ms)
print_row(rw)
print_row(rr)
if upd_ms:
print_row(ru)
print("-" * 108)
print("read 行 entity_id 一致性校验:%s"
% ("PASS 全部与请求一致" if not bad_rows else "FAIL 异常样本=%s" % bad_rows[:5]))
print("P99 门禁 SLA_P99_MS=%.1fms → write(INSERT) P99=%.4fms / read P99=%.4fms"
" / write(UPDATE) P99=%s"
% (SLA_P99_MS, rw["p99_ms"], rr["p99_ms"],
("%.4fms" % ru["p99_ms"]) if upd_ms else "N/A"))
if defect["seen"]:
print("缺陷复现(如实记录,未达标):%s %s" % (defect["code"], defect["msg"]))
print(" → 位置 pbl_runtime_ext/tx_write.py:_apply_state UPDATE 分支 build_row()")
if len(write_ms) < MIN_SAMPLES or len(read_ms) < MIN_SAMPLES:
print("采样不足 %d 次 → 判定不成立" % MIN_SAMPLES)
return 2, rw, rr
verdict = 0 if (rw["p99_pass"] and rr["p99_pass"] and not bad_rows
and not defect["seen"]
and (not upd_ms or ru["p99_pass"])) else 1
print("P99 结论:%s" % ("达标(均 ≤ 200ms)" if verdict == 0 else "未达标,待优化"))
print("=" * 108)
return verdict, rw, rr
def parse_samples(argv, default=60):
"""支持 `--samples 50` / `--samples=50` / 位置参数 `50` 三种写法。
M11b-2a 修正:原实现只认位置参数,按验收命令 `--samples 50` 调用时会被判为
「参数非法」并静默回落到默认 60 次,取证样本数与请求不符。
"""
n = default
i = 0
while i < len(argv):
a = argv[i]
if a.startswith("--samples"):
if "=" in a:
val = a.split("=", 1)[1]
else:
i += 1
val = argv[i] if i < len(argv) else ""
else:
val = a
try:
n = max(MIN_SAMPLES, int(val))
except (TypeError, ValueError):
print("采样参数非法:%r(用默认 %d)" % (a, default))
i += 1
return n
def main():
rc, _rw, _rr = run(parse_samples(sys.argv[1:]))
return rc
if __name__ == "__main__":
sys.exit(main())