203 lines
8.5 KiB
Python
203 lines
8.5 KiB
Python
#!/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 → 提交后广播」完整调用路径。
|
||
|
||
⚠ 采样走 **状态行首建(INSERT)分支**(每次采样用唯一 entity_id):
|
||
本工作空间实测发现 `_apply_state` 的 UPDATE 分支存在业务缺陷(PBL-RTX-0004,
|
||
见 probe 末尾 D1 段与 m11b2c-selftest-run.log),该缺陷属业务实现范畴,
|
||
M11b-2c-A 任务边界禁止修改 pbl_runtime_ext/*.py,故本探针不覆盖 UPDATE 分支,
|
||
并在输出中如实标注,交由后续任务修复后复测。
|
||
|
||
诚实声明:本工作空间无 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 = []
|
||
seq_rows = []
|
||
defect = {"seen": False, "code": "", "msg": ""}
|
||
|
||
async def main():
|
||
tenant, world, sess = 7, 9001, 90011
|
||
# 预热 3 次(不计入样本):排除 import/属性解析冷启动对 P99 的污染
|
||
for i in range(3):
|
||
await wtx.apply_runtime_event(
|
||
tenant_id=tenant, world_id=world, session_id=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()
|
||
assert len(rows) == 1, "read 主链路返回行数异常: %r" % (rows,)
|
||
read_ms.append((t3 - t2) * 1000.0)
|
||
|
||
# D1:状态行 **更新分支**(同 entity_id 二次写入)—— 本工作空间实测缺陷取证
|
||
try:
|
||
await wtx.apply_runtime_event(
|
||
tenant_id=tenant, world_id=world, session_id=sess,
|
||
event_type="move", payload={"x": 777},
|
||
state_updates=[{"state_key": "p99e0", "state_value": {"x": 777}}],
|
||
client_event_id="p99-update-probe")
|
||
print("D1[更新分支] 未抛异常(缺陷已修复或路径变化)")
|
||
except rtx_db.RtxError as exc:
|
||
defect = {"seen": True, "code": exc.code, "msg": str(exc.msg)[:160]}
|
||
print("D1[更新分支] 抛 %s → %s" % (exc.code, defect["msg"]))
|
||
|
||
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", write_ms)
|
||
rr = summarize("read_entity_state", read_ms)
|
||
print_row(rw)
|
||
print_row(rr)
|
||
print("-" * 108)
|
||
print("P99 门禁 SLA_P99_MS=%.1fms → write P99=%.4fms / read P99=%.4fms"
|
||
% (SLA_P99_MS, rw["p99_ms"], rr["p99_ms"]))
|
||
if defect["seen"]:
|
||
print("已知缺陷(非本任务范围,如实记录):%s %s" % (defect["code"], defect["msg"]))
|
||
print(" → 位置 pbl_runtime_ext/tx_write.py:_apply_state UPDATE 分支 build_row() "
|
||
"NOT NULL 全列校验;修复后需复测 UPDATE 分支 P99")
|
||
if len(write_ms) < MIN_SAMPLES:
|
||
print("采样不足 %d 次 → 判定不成立" % MIN_SAMPLES)
|
||
return 2, rw, rr
|
||
verdict = 0 if (rw["p99_pass"] and rr["p99_pass"]) else 1
|
||
print("P99 结论:%s" % ("达标(均 ≤ 200ms)" if verdict == 0 else "未达标,待优化"))
|
||
print("=" * 108)
|
||
return verdict, rw, rr
|
||
|
||
|
||
def main():
|
||
n = 60
|
||
if len(sys.argv) > 1:
|
||
try:
|
||
n = max(MIN_SAMPLES, int(sys.argv[1]))
|
||
except ValueError:
|
||
print("采样参数非法:%s(用默认 60)" % sys.argv[1])
|
||
rc, _rw, _rr = run(n)
|
||
return rc
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|