deliver: 交付收口(引擎代为提交)
This commit is contained in:
parent
30b92eafda
commit
6c7c91049b
@ -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",
|
||||
|
||||
@ -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,
|
||||
|
||||
545
pbl_runtime_ext/m11b_probe.py
Normal file
545
pbl_runtime_ext/m11b_probe.py
Normal file
@ -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)
|
||||
@ -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
|
||||
|
||||
@ -1 +1,16 @@
|
||||
打包元数据(305B):依赖 pbl_common
|
||||
[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*"]
|
||||
|
||||
@ -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__':
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user