diff --git a/pbl_runtime_ext/__init__.py b/pbl_runtime_ext/__init__.py index 63c806a..f299e28 100644 --- a/pbl_runtime_ext/__init__.py +++ b/pbl_runtime_ext/__init__.py @@ -31,6 +31,8 @@ from .tx_write import (apply_runtime_event, assert_append_only, # noqa: F401 idem_key_of, poll_events, read_states) from .m11b_api import (pbl_runtime_broadcast_pull, # noqa: F401 pbl_runtime_broadcast_stats) +from .m11b_probe import FakeEngine, run_probes # noqa: F401 +from .m11b_probe import self_check as m11b_probe_self_check # noqa: F401 __all__ = [ # 挂载 @@ -40,6 +42,8 @@ __all__ = [ # M11b 契约 "apply_runtime_event", "poll_events", "read_states", "assert_append_only", "pbl_runtime_broadcast_pull", "pbl_runtime_broadcast_stats", + # M11b 写入路径离线探针(state_updates 非空端到端覆盖) + "run_probes", "m11b_probe_self_check", "FakeEngine", "idem_key_of", "contracts", # M11b 基础设施 "RtxError", "dbname", "q_all", "q_one", "transaction", diff --git a/pbl_runtime_ext/init.py b/pbl_runtime_ext/init.py index 045f3e8..7940517 100644 --- a/pbl_runtime_ext/init.py +++ b/pbl_runtime_ext/init.py @@ -23,8 +23,9 @@ try: from . import partitions as _part from . import latency as _lat from . import rtx_db as _rtx + from . import m11b_probe as _probe except Exception as _exc: # noqa: BLE001 缺依赖时不阻断 M11a 能力 - _m11b = _wtx = _bcast = _part = _lat = _rtx = None + _m11b = _wtx = _bcast = _part = _lat = _rtx = _probe = None _M11B_IMPORT_ERROR = str(_exc) else: _M11B_IMPORT_ERROR = None @@ -116,7 +117,8 @@ def self_check(env=None): else: for mod, label in ((_wtx, "tx_write"), (_m11b, "m11b_api"), (_bcast, "broadcast"), (_part, "partitions"), - (_lat, "latency"), (_rtx, "rtx_db")): + (_lat, "latency"), (_rtx, "rtx_db"), + (_probe, "m11b_probe")): ok_i, msgs_i = mod.self_check() all_ok = all_ok and ok_i msgs.extend(["[%s] %s" % (label, m) for m in msgs_i]) @@ -254,6 +256,8 @@ def _registrations(): _api = None for name in list(reg.keys()): reg[name] = getattr(_api, name, None) if _api else None + if _probe is not None: + reg["m11b_probe_self_check"] = _probe.self_check if _m11b is not None: reg.update({ "apply_runtime_event": _m11b.apply_runtime_event, diff --git a/pbl_runtime_ext/m11b_probe.py b/pbl_runtime_ext/m11b_probe.py new file mode 100644 index 0000000..bbecd46 --- /dev/null +++ b/pbl_runtime_ext/m11b_probe.py @@ -0,0 +1,545 @@ +# -*- coding: utf-8 -*- +"""m11b_probe —— M11b 写入主链路的**离线可回滚探针**(响应 QC #1/#2 漏检根因)。 + +为什么需要这个文件 +------------------ +上一版 tx_write.apply_runtime_event 在 `state_updates` 非空分支残留占位死代码 +``await tx.execute("SELECT 1", (None, None))``,而离线 self_check 从不调用 +apply_runtime_event,导致该运行时缺陷(参数数与占位符数不匹配 → 驱动抛错 → +整个带状态更新的事务回滚)被漏检,self_check PASS 属**假阴性**。 + +本模块提供一台**内存事务引擎**(staged/committed 双区 + begin/commit/rollback + +故障注入),用真实代码路径跑完 apply_runtime_event,覆盖 6 类承诺: + + P1 正常写入(事件 + 2 条实体状态 + 快照):单事务、seq 单调、state_version 递增、 + **广播严格在 commit 之后**(无幽灵更新)、事件行零 UPDATE + P2 幂等重放:同 client_event_id 命中 → dedup=True,不新增事件行、**不广播** + P3 乐观锁冲突:base_version 不匹配 → 抛 PBL-STATE-CONFLICT,rollback、无 commit、 + 无广播、事件行不存在(真回滚) + P4 中途故障注入:状态 INSERT 抛错 → 事件行随事务消失(原子性 + 无脏数据) + P5 append-only 守卫:对 pbl_runtime_event 的 UPDATE/DELETE/TRUNCATE 一律拒绝 + P6 机械门禁自身有效性:**SQL 占位符数 == 参数个数**(正是漏检上一版缺陷的检查) + +另附 P7 契约三处同步核对(实现 / __init__.py 导出 / init.py env 注册)。 + +铁律:本文件只做**离线**探针,不连库、不建表;真实库端到端由 +scripts/m11b_selftest.py --live 负责,连不上时显式声明「环境受限未验证」。 +""" + +import asyncio +import json +import re + +from . import broadcast as _hub +from . import rtx_db +from .m11b_config import (ERR_APPEND_ONLY, ERR_STATE_CONFLICT, EVENT_TABLE, + SNAPSHOT_TABLE, STATE_TABLE) + +__all__ = ["FakeEngine", "run_probes", "self_check"] + +# 契约名(三处同步核对用) +CONTRACT_NAMES = ("apply_runtime_event", "poll_events", "read_states", + "assert_append_only", "pbl_runtime_broadcast_pull", + "pbl_runtime_broadcast_stats") + +_PH_RE = re.compile(r"%s") + + +def _arity_violation(sql, params): + """机械门禁:SQL 中 %s 占位符数必须等于参数个数。 + + 上一版 ``tx.execute("SELECT 1", (None, None))``(0 占位符 / 2 参数)就是靠这条 + 检查在自检阶段暴露,而不是等运行时驱动抛错。返回违规描述或 None。 + """ + text = str(sql) + want = len(_PH_RE.findall(text)) + got = 0 if params is None else len(params) + if want != got: + return ("占位符/参数数不匹配:占位符=%d 参数=%d SQL=%s" + % (want, got, text[:120])) + return None + + +def _is_write_on(sql, table): + """判断 SQL 是否是对某表的 UPDATE/DELETE/TRUNCATE/REPLACE。""" + text = " ".join(str(sql).split()).lower() + if table.lower() not in text: + return False + return bool(re.match(r"^\s*(update|delete|truncate|replace)", text)) + + +class FakeEngine(object): + """内存事务引擎:staged 区写入,commit 才落到 committed 区。 + + 参数: + fail_on : None | ("insert", table) | ("state_update", None) | ("execute", kw) + seed : {table: [row, ...]} 预置已提交行(用于幂等/乐观锁场景) + """ + + def __init__(self, fail_on=None, seed=None): + self.fail_on = fail_on + self.journal = [] # 事件序列:begin/stmt:*/commit/rollback/publish + self.errors = [] # 机械门禁违规(占位符数、append-only 写) + self._id = 0 + self.committed = {} + self.staged = {} + self.tx_obj = None # 由 _Session 注入,供探针核对事务生命周期 + for tbl, rows in (seed or {}).items(): + self.committed.setdefault(tbl, []).extend([dict(r) for r in rows]) + + # ---------------- 行可见性:staged + committed(同事务内可读自己写的) + def _all(self, table): + out = list(self.committed.get(table, [])) + out.extend(self.staged.get(table, [])) + return out + + def _next_id(self): + self._id += 1 + return self._id + + # ---------------- 事务原语(与 rtx_db.Transaction 同签名) + def begin(self): + self.journal.append("begin") + + def _should_fail(self, kind, table=None, sql=None): + if not self.fail_on: + return False + f_kind, f_target = self.fail_on + if f_kind != kind: + return False + if f_kind == "insert": + return table == f_target + if f_kind == "state_update": + return table == STATE_TABLE + if f_kind == "execute": + return str(f_target or "") in str(sql or "") + return False + + def execute(self, sql, params=None): + bad = _arity_violation(sql, params) + if bad: + self.errors.append(bad) + raise rtx_db.RtxError("PBL-RTX-0005", bad) + for tbl in (EVENT_TABLE, STATE_TABLE, SNAPSHOT_TABLE): + if _is_write_on(sql, tbl): + if tbl == EVENT_TABLE: + self.errors.append( + "append-only 违规:对 %s 执行了写语句 %s" % (tbl, str(sql)[:80])) + raise rtx_db.RtxError(ERR_APPEND_ONLY, self.errors[-1]) + if self._should_fail("state_update", tbl, sql): + raise RuntimeError("故障注入:状态 UPDATE 失败") + self.journal.append("execute") + return self._respond(sql, params) + + def insert(self, table, row): + if self._should_fail("insert", table, table): + raise RuntimeError("故障注入:%s INSERT 失败" % table) + self.journal.append("insert") + rid = self._next_id() + rec = dict(row) + rec["id"] = rid + self.staged.setdefault(table, []).append(rec) + return rid + + def commit(self): + self.journal.append("commit") + for tbl, rows in self.staged.items(): + self.committed.setdefault(tbl, []).extend(rows) + self.staged = {} + + def rollback(self): + self.journal.append("rollback") + self.staged = {} + + # ---------------- 只读响应(按 SQL 关键字路由,语义与真实表一致) + def _respond(self, sql, params): + text = " ".join(str(sql).split()) + low = text.lower() + params = list(params or []) + + if "max(" in low and EVENT_TABLE in text: + sid = params[1] if len(params) > 1 else None + tid = params[0] if params else None + seqs = [int(r.get("seq") or 0) for r in self._all(EVENT_TABLE) + if str(r.get("tenant_id")) == str(tid) + and str(r.get("session_id")) == str(sid)] + return [{"max_seq": max(seqs) if seqs else 0}] + + if "from " + EVENT_TABLE in low and "event_uid" in low: + tid, uid = params[0], params[1] + for r in self._all(EVENT_TABLE): + if str(r.get("tenant_id")) == str(tid) and str(r.get("event_uid")) == str(uid): + return [{"id": r.get("id"), "seq": r.get("seq"), + "event_type": r.get("event_type"), + "state_version": r.get("state_version")}] + return [] + + if "from " + STATE_TABLE in low: + tid, sid = params[0], params[1] + eid = params[2] if len(params) > 2 else None + out = [] + for r in self._all(STATE_TABLE): + if str(r.get("tenant_id")) != str(tid) or str(r.get("session_id")) != str(sid): + continue + if eid is not None and str(r.get("entity_id")) != str(eid): + continue + out.append({"id": r.get("id"), + "state_version": r.get("state_version"), + "entity_id": r.get("entity_id")}) + return out + + return [] + + # ---------------- 断言辅助 + def rows_of(self, table): + return self.committed.get(table, []) + + def index_of(self, marker): + try: + return self.journal.index(marker) + except ValueError: + return -1 + + +class _FakeTx(object): + """与 rtx_db.Transaction 对外的最小契约一致(execute/insert/commit/rollback/stmt_count)。""" + + def __init__(self, engine): + self.engine = engine + self.stmt_count = 0 + self.committed = False + self.rolled_back = False + + async def execute(self, sql, params=None): + self.stmt_count += 1 + return self.engine.execute(sql, params) + + async def insert(self, table, row): + self.stmt_count += 1 + return self.engine.insert(table, row) + + async def commit(self): + self.committed = True + self.engine.commit() + + async def rollback(self): + self.rolled_back = True + self.engine.rollback() + + async def __aenter__(self): + self.engine.begin() + return self + + async def __aexit__(self, exc_type, exc, tb): + # 与真实 rtx_db 一致:异常时若未显式 rollback 则兜底回滚 + if exc_type is not None and not self.rolled_back: + self.engine.rollback() + return False + + +class _Session(object): + """临时接管 rtx_db.transaction / hub.publish,跑完即还原(不污染真实运行时)。""" + + def __init__(self, engine, world_id, tenant_id): + self.engine = engine + self.world_id = world_id + self.tenant_id = tenant_id + self.published = [] + self._orig_tx = None + self._hub = None + self._orig_publish = None + self.tx_obj = None + + def __enter__(self): + from . import tx_write as wtx + engine = self.engine + published = self.published + + def _transaction(db=None): + tx = _FakeTx(engine) + self.tx_obj = tx + engine.tx_obj = tx # 探针可读:事务对象生命周期状态 + return tx + + self._orig_tx = wtx.rtx_db.transaction + wtx.rtx_db.transaction = _transaction + + hub = _hub.get_hub() + self._hub = hub + self._orig_publish = hub.publish + + def _publish(channel, evt): + engine.journal.append("publish") + published.append((channel, evt)) + return self._orig_publish(channel, evt) + + hub.publish = _publish + return self + + def __exit__(self, exc_type, exc, tb): + from . import tx_write as wtx + try: + wtx.rtx_db.transaction = self._orig_tx + if self._hub is not None and self._orig_publish is not None: + self._hub.publish = self._orig_publish + except Exception: # noqa: BLE001 + pass + return False + + +async def _apply(**kw): + from . import tx_write as wtx + return await wtx.apply_runtime_event(**kw) + + +def _base_kwargs(world_id, tenant_id, session_id, **over): + kw = dict(tenant_id=tenant_id, world_id=world_id, session_id=session_id, + event_type="entity.move", + payload={"x": 1, "y": 2}, + state_updates=[{"entity_id": "e1", "state": {"hp": 10}}, + {"entity_id": "e2", "state": {"hp": 7}}], + actor_id="tester", client_event_id=None) + kw.update(over) + return kw + + +# ---------------------------------------------------------------- P1 +async def p1_normal_write(msgs): + """事件 + 2 状态 + 快照:单事务、提交后广播、事件行零 UPDATE。""" + eng = FakeEngine() + tid, wid, sid = "t1", "w-p1", "s1" + with _Session(eng, wid, tid): + res = await _apply(**_base_kwargs( + wid, tid, sid, client_event_id="cli-p1", + snapshot={"reason": "probe", "entities": {"e1": {"hp": 10}}})) + ok = True + ev_rows = eng.rows_of(EVENT_TABLE) + st_rows = eng.rows_of(STATE_TABLE) + sn_rows = eng.rows_of(SNAPSHOT_TABLE) + if not res.get("ok") or res.get("dedup"): + ok = False + if res.get("state_rows") != 2 or len(st_rows) != 2: + ok = False + if len(ev_rows) != 1 or ev_rows[0].get("seq") != 1: + ok = False + if len(sn_rows) != 1: + ok = False + if sorted(int(s["state_version"]) for s in res.get("states") or []) != [1, 1]: + ok = False + ci, pi = eng.index_of("commit"), eng.index_of("publish") + if ci < 0 or pi < 0 or pi < ci: + ok = False + if eng.errors: + ok = False + if not eng.tx_obj or not eng.tx_obj.committed or eng.tx_obj.rolled_back: + ok = False + if res.get("broadcast_ms") is None: + ok = False + msgs.append("%s P1 正常写入(事件+2状态+快照)单事务提交:stmt=%d 事件行=%d 状态行=%d 快照行=%d " + "广播在commit后=%s 占位符违规=%s" + % ("PASS" if ok else "FAIL", + (eng.tx_obj.stmt_count if eng.tx_obj else -1), len(ev_rows), + len(st_rows), len(sn_rows), (pi > ci >= 0), eng.errors or 0)) + return ok + + +# ---------------------------------------------------------------- P2 +async def p2_idempotent_replay(msgs): + """同 client_event_id 重放:命中已有事件 → dedup,不新增行、不广播。""" + tid, wid, sid = "t1", "w-p2", "s1" + eng = FakeEngine(seed={EVENT_TABLE: [ + {"id": 900, "tenant_id": tid, "event_uid": "cli-p2", "session_id": sid, + "seq": 7, "event_type": "entity.move", "state_version": 3}]}) + with _Session(eng, wid, tid): + res = await _apply(**_base_kwargs(wid, tid, sid, client_event_id="cli-p2")) + ok = bool(res.get("dedup")) and res.get("seq") == 7 and res.get("event_id") == 900 + if eng.rows_of(EVENT_TABLE).__len__() != 1: + ok = False + if eng.index_of("publish") >= 0: + ok = False + if eng.rows_of(STATE_TABLE): + ok = False + if eng.errors: + ok = False + msgs.append("%s P2 幂等重放:dedup=%s 事件行数=%d 广播次数=%d(重放不重复写、不重复广播)" + % ("PASS" if ok else "FAIL", res.get("dedup"), + len(eng.rows_of(EVENT_TABLE)), eng.journal.count("publish"))) + return ok + + +# ---------------------------------------------------------------- P3 +async def p3_optimistic_conflict(msgs): + """base_version 不匹配 → 冲突回滚,事件行不存在、不广播、无脏数据。""" + tid, wid, sid = "t1", "w-p3", "s1" + eng = FakeEngine(seed={STATE_TABLE: [ + {"id": 500, "tenant_id": tid, "session_id": sid, "entity_id": "e1", + "state_json": '{"hp":1}', "state_version": 5}]}) + code = None + with _Session(eng, wid, tid): + try: + await _apply(**_base_kwargs(wid, tid, sid, client_event_id="cli-p3", + base_version=3)) + except rtx_db.RtxError as exc: + code = exc.code + ok = (code == ERR_STATE_CONFLICT) + if eng.index_of("commit") >= 0: + ok = False + if eng.index_of("rollback") < 0: + ok = False + if eng.index_of("publish") >= 0: + ok = False + if eng.rows_of(EVENT_TABLE) or eng.rows_of(SNAPSHOT_TABLE): + ok = False + if eng.errors: + ok = False + if not eng.tx_obj or eng.tx_obj.committed or not eng.tx_obj.rolled_back: + ok = False + msgs.append("%s P3 乐观锁冲突回滚:错误码=%s commit=%d rollback=%d 事件行残留=%d 广播=%d" + % ("PASS" if ok else "FAIL", code, eng.index_of("commit"), + eng.index_of("rollback"), len(eng.rows_of(EVENT_TABLE)), + eng.journal.count("publish"))) + return ok + + +# ---------------------------------------------------------------- P4 +async def p4_midway_failure_atomic(msgs): + """状态 INSERT 中途抛错 → 事件行随事务消失(原子性,非「标 rolled_back」)。""" + tid, wid, sid = "t1", "w-p4", "s1" + eng = FakeEngine(fail_on=("insert", STATE_TABLE)) + raised = None + with _Session(eng, wid, tid): + try: + await _apply(**_base_kwargs(wid, tid, sid, client_event_id="cli-p4", + snapshot={"reason": "probe"})) + except Exception as exc: # noqa: BLE001 + raised = type(exc).__name__ + ok = raised is not None + if eng.index_of("rollback") < 0: + ok = False + if eng.index_of("commit") >= 0 or eng.index_of("publish") >= 0: + ok = False + if eng.rows_of(EVENT_TABLE) or eng.rows_of(STATE_TABLE) or eng.rows_of(SNAPSHOT_TABLE): + ok = False + if not eng.tx_obj or eng.tx_obj.committed or not eng.tx_obj.rolled_back: + ok = False + msgs.append("%s P4 中途故障原子性:抛出=%s rollback=%d 已提交事件行=%d 已提交状态行=%d 广播=%d" + % ("PASS" if ok else "FAIL", raised, eng.index_of("rollback"), + len(eng.rows_of(EVENT_TABLE)), len(eng.rows_of(STATE_TABLE)), + eng.journal.count("publish"))) + return ok + + +# ---------------------------------------------------------------- P5 +def p5_append_only(msgs): + from . import tx_write as wtx + blocked = [] + for op in ("UPDATE", "DELETE", "TRUNCATE"): + try: + wtx.assert_append_only(op, EVENT_TABLE) + except rtx_db.RtxError as exc: + if exc.code == ERR_APPEND_ONLY: + blocked.append(op) + ok = len(blocked) == 3 + # 引擎侧二次拦截:即使绕过 assert,FakeEngine 也拒绝任何对事件表的写 + eng = FakeEngine() + try: + eng.execute("UPDATE %s SET state_version=%%s WHERE id=%%s" % EVENT_TABLE, + (1, 2)) + engine_guard = False + except rtx_db.RtxError as exc: + engine_guard = exc.code == ERR_APPEND_ONLY + msgs.append("%s P5 append-only:%s 被代码守卫拦截=%s 引擎兜底拦截=%s" + % ("PASS" if (ok and engine_guard) else "FAIL", "/".join(blocked) or "无", + ok, engine_guard)) + return ok and engine_guard + + +# ---------------------------------------------------------------- P6 +def p6_gate_effective(msgs): + """门禁自身有效性:0 占位符 / 2 参数(上一版缺陷形态)必须被拒绝。""" + eng = FakeEngine() + caught = None + try: + eng.execute("SELECT 1", (None, None)) + except rtx_db.RtxError as exc: + caught = exc.code + ok = caught == "PBL-RTX-0005" and bool(eng.errors) + msgs.append("%s P6 占位符/参数数门禁生效:'SELECT 1'+2参数 被拒(错误码=%s)" + % ("PASS" if ok else "FAIL", caught)) + return ok + + +# ---------------------------------------------------------------- P7 +def p7_triple_sync(msgs): + """契约三处同步:实现定义 / __init__.py 导出 / init.py env 注册。""" + import pbl_runtime_ext as pkg + from . import init as pkg_init + from . import tx_write as wtx + try: + from . import m11b_api as mapi + except Exception: # noqa: BLE001 + mapi = None + reg = pkg_init._registrations() + bad = [] + for name in CONTRACT_NAMES: + impl = getattr(wtx, name, None) or (getattr(mapi, name, None) if mapi else None) + exported = hasattr(pkg, name) + registered = callable(reg.get(name)) + if not (callable(impl) and exported and registered): + bad.append("%s(impl=%s export=%s env=%s)" + % (name, callable(impl), exported, registered)) + ok = not bad + msgs.append("%s P7 契约三处同步:%d/%d 对齐 %s" + % ("PASS" if ok else "FAIL", + len(CONTRACT_NAMES) - len(bad), len(CONTRACT_NAMES), + ("缺失/失配: " + "; ".join(bad)) if bad else "")) + return ok + + +def run_probes(): + """跑全部探针,返回 (ok, msgs)。任何探针异常都记 FAIL(不静默跳过)。""" + msgs = [] + results = [] + + async def _run_all(): + for coro in (p1_normal_write(msgs), p2_idempotent_replay(msgs), + p3_optimistic_conflict(msgs), p4_midway_failure_atomic(msgs)): + try: + results.append(await coro) + except Exception as exc: # noqa: BLE001 + msgs.append("FAIL 探针执行异常:%r" % exc) + results.append(False) + + try: + loop = asyncio.new_event_loop() + try: + loop.run_until_complete(_run_all()) + finally: + loop.close() + except Exception as exc: # noqa: BLE001 + msgs.append("FAIL 事件循环不可用:%r" % exc) + results.append(False) + + for sync_probe in (p5_append_only, p6_gate_effective, p7_triple_sync): + try: + results.append(bool(sync_probe(msgs))) + except Exception as exc: # noqa: BLE001 + msgs.append("FAIL %s 异常:%r" % (sync_probe.__name__, exc)) + results.append(False) + + ok = bool(results) and all(results) + msgs.append("m11b_probe %s:7 项探针(state_updates 非空端到端 P1~P4 + 守卫 P5 + " + "门禁有效性 P6 + 三处同步 P7)" % ("PASS" if ok else "FAIL")) + return ok, msgs + + +def self_check(): + """与包内其它模块一致的自检出口。""" + return run_probes() + + +if __name__ == "__main__": + _ok, _msgs = run_probes() + for _m in _msgs: + print(_m) + raise SystemExit(0 if _ok else 1) diff --git a/pbl_runtime_ext/tx_write.py b/pbl_runtime_ext/tx_write.py index dc9ec4f..9b6ae39 100644 --- a/pbl_runtime_ext/tx_write.py +++ b/pbl_runtime_ext/tx_write.py @@ -338,19 +338,11 @@ async def apply_runtime_event(tenant_id=None, world_id=None, session_id=None, st = await _apply_state(tx, tenant_id, session_id, upd, seq_no, base_version, actor, now) state_events.append(st) - if state_events: - ver = max(int(s["state_version"]) for s in state_events) - await tx.execute( - "UPDATE %s SET %s=%%s WHERE id=%%s" - % (EVENT_TABLE, EVENT_COL_STATE_VERSION) - if False else - "UPDATE %s SET %s=%%s WHERE id=%%s" - % ("pbl_runtime_event_shadow_noop", EVENT_COL_STATE_VERSION) - if False else - "SELECT 1", (None, None)) # 占位:见下方真实回填 - # 事件行的 state_version 回填(append-only 允许同事务内 - # 在 INSERT 之前无法得知版本,故改为写入前不可知 → 不回填, - # 由 state_events 返回值携带,避免对 append-only 表 UPDATE) + # 4b) 事件行 state_version:pbl_runtime_event 为 append-only, + # 同事务内 INSERT 之后**不再 UPDATE 回填**(对 append-only 表 + # 的任何 UPDATE 都被 assert_append_only 拦截)。本轮 UPSERT 后的 + # 最新 state_version 只通过返回值 states[] 下发给调用方, + # 事件行保持写入时的 state_version(=本次事件 seq 对应的水位)。 # 5) 快照(掉线补齐/离线兜底) if snapshot is not None: await tx.insert(SNAPSHOT_TABLE, build_row(SNAPSHOT_TABLE, { @@ -554,6 +546,20 @@ def self_check(): msgs.append("tx_write FAIL:%s 非协程" % name) if callable(assert_append_only): msgs.append("tx_write PASS:3 协程契约 + 1 同步守卫签名正确") + + # ---- 端到端写入探针(QC #2:把 state_updates 非空路径纳入自检覆盖)---- + # 上一版占位死代码 tx.execute("SELECT 1", (None, None)) 只在 state_updates 非空 + # 时执行,离线自检不调用 apply_runtime_event 因而漏检(假阴性)。此处强制跑 + # m11b_probe 的内存事务引擎探针,任何占位/参数数不匹配/幽灵广播都会 FAIL。 + try: + from . import m11b_probe as _probe # 延迟导入:避免模块级循环依赖 + p_ok, p_msgs = _probe.run_probes() + ok = ok and p_ok + msgs.extend(["probe| " + m for m in p_msgs]) + except Exception as exc: # noqa: BLE001 + ok = False + msgs.append("tx_write FAIL:m11b_probe 不可用(写入路径未被自检覆盖)%r" % exc) + if ok: msgs.append("tx_write.self_check PASS") return ok, msgs diff --git a/pyproject.toml b/pyproject.toml index a823f3f..57767b3 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1 +1,16 @@ -打包元数据(305B):依赖 pbl_common \ No newline at end of file +[build-system] +requires = ["setuptools>=45", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "pbl_runtime_ext" +version = "1.0.0" +description = "scense_runtime 薄扩展:单事务事件+状态写入与广播(M11a/M11b)" +requires-python = ">=3.8" +dependencies = [ + "sqlor", +] + +[tool.setuptools.packages.find] +where = ["."] +include = ["pbl_runtime_ext*"] diff --git a/scripts/load_path.py b/scripts/load_path.py index b167547..937ba60 100644 --- a/scripts/load_path.py +++ b/scripts/load_path.py @@ -1,19 +1,26 @@ #!/usr/bin/env python3 # -*- coding: utf-8 -*- -"""pbl_runtime_ext RBAC 路径注册(硬门禁 6.6 / QC #11)。 +"""pbl_runtime_ext RBAC 路径注册(硬门禁 6.6 / QC #11 / M11b QC #8)。 约定: - 路径 = 模块自动路由 `/pbl_runtime_ext/api/<契约>.dspy`,不带端口、不带 /wss 前缀; -- 角色 `logined` = 登录即可访问的读接口;写接口按角色分级(teacher/admin); -- 由 apps/pbls/build.sh 第 8 步调用 `register()`;rbac CLI 不在位时打印清单(不静默跳过)。 +- 角色 `logined` = 登录即可访问;本模块 8 个契约全部要求登录(含写接口), + **禁止通配符**(module-development-spec 硬规定:每条路径显式列出); +- 由 apps/pbls/build.sh 第 8 步调用 `register()`;rbac CLI 不在位时打印 PENDING 清单 + (返回 False,不静默跳过、不冒充成功)。 + +维护规则:wwwroot/api/ 下 .dspy 增删必须同步本清单 —— `check_paths()` 会在 +自检阶段比对磁盘文件与 PATHS,缺登记或多登记都判 FAIL。 """ import os import subprocess import sys MODULE = 'pbl_runtime_ext' +HERE = os.path.dirname(os.path.abspath(__file__)) +WWWROOT_API = os.path.join(os.path.dirname(HERE), 'wwwroot', 'api') -# (path, role) +# (path, role) —— M11a 6 个 + M11b 广播 2 个,共 8 个显式路径 PATHS = [ ('/pbl_runtime_ext/api/pbl_runtime_event_append.dspy', 'logined'), ('/pbl_runtime_ext/api/pbl_runtime_event_poll.dspy', 'logined'), @@ -23,24 +30,49 @@ PATHS = [ ('/pbl_runtime_ext/api/pbl_session_member_add.dspy', 'logined'), ('/pbl_runtime_ext/api/pbl_runtime_broadcast_pull.dspy', 'logined'), ('/pbl_runtime_ext/api/pbl_runtime_broadcast_stats.dspy', 'logined'), - ] +def on_disk(): + """磁盘上真实存在的 .dspy 文件名集合(小写,用于比对)。""" + try: + return set(f for f in os.listdir(WWWROOT_API) if f.endswith('.dspy')) + except OSError: + return set() + + +def check_paths(): + """清单 ↔ 磁盘一致性自检。返回 (ok, missing_in_list, stale_in_list)。""" + listed = set(os.path.basename(p) for p, _ in PATHS) + disk = on_disk() + missing = sorted(disk - listed) # 有文件没登记 → 上线必 403 + stale = sorted(listed - disk) # 登记了但文件不存在 → 死路径 + return (not missing and not stale), missing, stale + + def register(): tool = os.environ.get('RBAC_SET_PERM', 'set_role_perm.py') done, missing = 0, [] for path, role in PATHS: - if subprocess.call([sys.executable if os.environ.get('PY') else 'python3', - tool, role, path], - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) == 0: + cmd = [os.environ.get('PY', 'python3'), tool, role, path] + try: + rc = subprocess.call(cmd, stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL) + except Exception: # noqa: BLE001 工具不在位也不能崩,记 PENDING + rc = 1 + if rc == 0: done += 1 else: missing.append((path, role)) - print('[%s] rbac paths: total=%d ok=%d pending=%d' %(len(PATHS), done, len(missing))) + print('[%s] rbac paths: total=%d ok=%d pending=%d' + % (MODULE, len(PATHS), done, len(missing))) for path, role in missing: - print(' PENDING %%-12s %s' %(role, path)) - return len(missing) == 0 + print(' PENDING %-12s %s' % (role, path)) + ok_list, absent, stale = check_paths() + if not ok_list: + print(' PATHS-DISK-MISMATCH missing_in_list=%s stale_in_list=%s' + % (absent, stale)) + return len(missing) == 0 and ok_list if __name__ == '__main__':