#!/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())