diff --git a/pbl_runtime_ext/__init__.py b/pbl_runtime_ext/__init__.py index 633b298..16ac1c3 100644 --- a/pbl_runtime_ext/__init__.py +++ b/pbl_runtime_ext/__init__.py @@ -1 +1,12 @@ -包声明(98B) \ No newline at end of file +# -*- coding: utf-8 -*- +"""pbl_runtime_ext —— scense_runtime 薄扩展:单事务事件+状态写入与广播(M11a/M11b)。 + +铁律:**不改基表**(scense_runtime / world / scene / entity 基表只读), +扩展一律落在 pbl_* 自有表(pbl_runtime_event / pbl_entity_state / pbl_world_state_snapshot)。 + +包出口: + from pbl_runtime_ext.init import load_pbl_runtime_ext +""" + +__all__ = ["init", "tables", "tx_event", "api"] +__version__ = "1.0.0" diff --git a/pbl_runtime_ext/init.py b/pbl_runtime_ext/init.py index 6555d95..dd9da64 100644 --- a/pbl_runtime_ext/init.py +++ b/pbl_runtime_ext/init.py @@ -1 +1,245 @@ -挂载入口(1,758B):ensure_tables 建 3 表;load_pbl_runtime_ext(app,sor,ensure,pusher) 可注入广播通道,注册 api(apply_event/list_events/get_state/replay),返回体标 write_protected_base:['scense_runtime'] 供 QC 核验 \ No newline at end of file +# -*- coding: utf-8 -*- +"""pbl_runtime_ext 挂载入口(M11a/M11b)。 + +load_pbl_runtime_ext(env=None): + 1. ensure_tables 幂等建 3 表(pbl_runtime_event / pbl_entity_state / pbl_world_state_snapshot) + 2. **挂载即执行 self_check()**:写保护域(world/scene/entity/scense/scense_runtime/ + script_engine 基表)出现 U/D 调用、或单事务/幂等契约缺失 → 抛 RuntimeError 拒绝启动 + 3. 注册 api(append_event / list_events / find_by_idem / next_seq / self_check) + 4. 返回 api 字典 + +铁律:本模块是 scense_runtime 的**薄扩展**,只写 pbl_* 自有表,不改任何基表。 +""" + +from . import tables as _tables +from . import tx_event as _tx + +MODULE = "pbl_runtime_ext" +EXPECTED_TABLES = 3 + + +def ensure_tables(env=None, sor=None): + """幂等建 3 表。返回 (ok_list, bad_list)。""" + return _tables.ensure_tables(env=env, sor=sor) + + +def api(): + """对外 API 契约。""" + return { + "append_event": _tx.append_event, + "list_events": _tx.list_events, + "find_by_idem": _tx.find_by_idem, + "next_seq": _tx.next_seq, + "TxError": _tx.TxError, + "self_check": self_check, + "ensure_tables": ensure_tables, + "EVENT_TABLE": _tx.EVENT_TABLE, + "STATE_TABLE": _tx.STATE_TABLE, + "SNAPSHOT_TABLE": _tx.SNAPSHOT_TABLE, + } + + +def self_check(env=None): + """模块自检:表契约 + 写保护 + 单事务/幂等契约。返回 (all_ok, msgs)。""" + msgs = [] + all_ok = True + + tbl_ok, tbl_msgs = _tables.self_check() + if not tbl_ok: + all_ok = False + msgs.extend(tbl_msgs) + if len(_tables.TABLES) != EXPECTED_TABLES: + all_ok = False + msgs.append("表数=%d 应为 %d" % (len(_tables.TABLES), EXPECTED_TABLES)) + + tx_ok, tx_msgs = _tx.self_check() + if not tx_ok: + all_ok = False + msgs.extend(tx_msgs) + + # 幂等探针:同 idem_key 二次投递必须走 dedup 分支(离线用假 sor 验证) + probe_ok, probe_msgs = _probe_idempotent() + if not probe_ok: + all_ok = False + msgs.extend(probe_msgs) + + if all_ok: + msgs.append("SELF_CHECK %s: PASS %d/%d" % (MODULE, len(_tables.TABLES), len(_tables.TABLES))) + return all_ok, msgs + + +class _FakeSor(object): + """离线自检用内存 sor(只实现 C/R/U),验证单事务+幂等真实行为。""" + + def __init__(self): + self.rows = {} + self.seq = 0 + self.updated = [] + + def C(self, tbl, row): + self.seq += 1 + row = dict(row) + row["id"] = self.seq + self.rows.setdefault(tbl, []).append(row) + return self.seq + + def R(self, tbl, where="", fields="*", order="", limit=100): + out = [] + for r in self.rows.get(tbl, []): + hit = True + for clause in str(where).split(" AND "): + if "=" not in clause: + continue + k, v = clause.split("=", 1) + v = v.strip().strip("'") + if str(r.get(k.strip())) != v: + hit = False + break + if hit: + out.append(r) + return out[:limit] + + def U(self, tbl, values, where=""): + self.updated.append((tbl, dict(values), where)) + return len(self.rows.get(tbl, [])) + + def sqlExe(self, sql): + return 0 + + +class _FakeEnv(object): + def __init__(self): + self.sor = _FakeSor() + self.broadcasts = [] + + def broadcast(self, channel, payload): + self.broadcasts.append((channel, payload)) + return 1 + + +def _probe_idempotent(): + """幂等 + 单事务 + 广播探针。返回 (all_ok, msgs)。""" + msgs = [] + all_ok = True + env = _FakeEnv() + try: + r1 = _tx.append_event(env, tenant_id=7, world_id=1, session_id=11, + event_type="move", payload={"x": 1}, + idem_key="probe-key-1", + state_updates=[{"state_key": "pos", "state_value": {"x": 1}}]) + r2 = _tx.append_event(env, tenant_id=7, world_id=1, session_id=11, + event_type="move", payload={"x": 1}, + idem_key="probe-key-1", + state_updates=[{"state_key": "pos", "state_value": {"x": 1}}]) + except Exception as exc: # noqa: BLE001 + return False, ["探针[单事务写入] FAIL:%r" % exc] + + if r1.get("dedup") is not False: + all_ok = False + msgs.append("探针[首次写入] FAIL:dedup=%r 应为 False" % r1.get("dedup")) + else: + msgs.append("探针[首次写入] PASS seq_no=%s state=%s" % (r1.get("seq_no"), r1.get("state"))) + if r2.get("dedup") is not True: + all_ok = False + msgs.append("探针[幂等去重] FAIL:dedup=%r 应为 True" % r2.get("dedup")) + else: + msgs.append("探针[幂等去重] PASS 同 idem_key 未重复写") + n_ev = len(env.sor.rows.get(_tx.EVENT_TABLE, [])) + if n_ev != 1: + all_ok = False + msgs.append("探针[事件行数] FAIL:%d 应为 1" % n_ev) + else: + msgs.append("探针[事件行数] PASS 1 行") + if len(env.sor.rows.get(_tx.STATE_TABLE, [])) != 1: + all_ok = False + msgs.append("探针[状态行数] FAIL:应为 1") + else: + msgs.append("探针[状态行数] PASS 1 行") + if len(env.sor.rows.get(_tx.SNAPSHOT_TABLE, [])) != 1: + all_ok = False + msgs.append("探针[快照行数] FAIL:应为 1") + else: + msgs.append("探针[快照行数] PASS 1 行") + if not env.broadcasts: + all_ok = False + msgs.append("探针[提交后广播] FAIL:未广播") + else: + msgs.append("探针[提交后广播] PASS 通道=%d" % len(env.broadcasts)) + + # 缺租户必须拒 + try: + _tx.append_event(env, tenant_id=None, world_id=1, session_id=11, event_type="move") + all_ok = False + msgs.append("探针[缺租户必须拒] FAIL:未抛 TxError") + except _tx.TxError as exc: + if exc.code != "PBL-TENANT-0001": + all_ok = False + msgs.append("探针[缺租户必须拒] FAIL:code=%s" % exc.code) + else: + msgs.append("探针[缺租户必须拒] PASS PBL-TENANT-0001") + + # 状态写失败必须整体回滚(事件标 rolled_back) + env2 = _FakeEnv() + try: + _tx.append_event(env2, tenant_id=7, world_id=1, session_id=12, + event_type="move", payload={}, + state_updates=[{"bad": "no state_key"}]) + all_ok = False + msgs.append("探针[状态失败回滚] FAIL:未抛 TxError") + except _tx.TxError as exc: + rolled = [u for u in env2.sor.updated if u[1].get("state") == "rolled_back"] + if not rolled: + all_ok = False + msgs.append("探针[状态失败回滚] FAIL:事件未标 rolled_back(%r)" % exc.code) + else: + msgs.append("探针[状态失败回滚] PASS 事件标 rolled_back,无脏数据") + return all_ok, msgs + + +def load_pbl_runtime_ext(env=None): + """挂载入口:建表 → 自检(不过即抛)→ 注册契约 → 返回 api 字典。""" + srv = env + if srv is None: + try: + from ahserver.serverenv import ServerEnv + srv = ServerEnv() + except Exception: # noqa: BLE001 + srv = None + + if srv is not None: + ensure_tables(srv) + + ok, msgs = self_check(srv) + for m in msgs: + _log(srv, m) + if not ok: + raise RuntimeError( + "%s self_check FAILED(fail-closed,拒绝启动):%s" + % (MODULE, "; ".join([m for m in msgs if "FAIL" in m or "应为" in m][:6])) + ) + + if srv is not None: + setattr(srv, "pbl_runtime_ext_api", api()) + modules = getattr(srv, "modules", None) + if isinstance(modules, list) and MODULE not in modules: + modules.append(MODULE) + return api() + + +def _log(srv, msg): + logger = getattr(srv, "logger", None) if srv is not None else None + if logger is not None and hasattr(logger, "info"): + try: + logger.info("[%s] %s" % (MODULE, msg)) + return + except Exception: # noqa: BLE001 + pass + print("[%s] %s" % (MODULE, msg)) + + +if __name__ == "__main__": + _ok, _msgs = self_check() + for _m in _msgs: + print(_m) + print("RESULT: %s" % ("PASS" if _ok else "FAIL")) + raise SystemExit(0 if _ok else 1) diff --git a/pbl_runtime_ext/tables.py b/pbl_runtime_ext/tables.py index 289822f..09d8bdf 100644 --- a/pbl_runtime_ext/tables.py +++ b/pbl_runtime_ext/tables.py @@ -1 +1,142 @@ -3 表 DDL(4,739B):pbl_rt_event(13字段,append-only,uk_rte_tenant_sess_seq,idx tx/type/time) / pbl_rt_state(11字段,uk_rts_tenant_sess_entity,state longtext 全量覆盖,state_version 乐观锁,last_event_id/last_tx_id) / pbl_rt_broadcast(14字段,tx_id/event_id/channel/audience/message/state pending-sent-failed/retry_count/last_error/sent_time);ddl()/all_ddl() \ No newline at end of file +# -*- coding: utf-8 -*- +"""pbl_runtime_ext 表定义与幂等建表(M11a)。 + +3 表(mariadb 方言,BIGINT AUTO_INCREMENT,tenant_id 强制打头,无 FK/ENUM/TIMESTAMP): + * pbl_runtime_event —— 运行时事件流(append-only,幂等键去重) + * pbl_entity_state —— 实体状态(薄扩展,不改 entity 基表) + * pbl_world_state_snapshot —— 世界状态快照(单事务写入后广播基线) + +与 modules/pbl_runtime_ext/models/*.json 及 sql/pbl_runtime_ext.sql 同构。 +""" + +TABLES = ["pbl_runtime_event", "pbl_entity_state", "pbl_world_state_snapshot"] + +DDL = [ + """ +CREATE TABLE IF NOT EXISTS `pbl_runtime_event` ( + `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键', + `tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户ID(多租户强制打头)', + `event_code` VARCHAR(64) NOT NULL COMMENT '事件编码', + `idem_key` VARCHAR(128) NOT NULL DEFAULT '' COMMENT '幂等键(同键重复投递只落一条)', + `world_id` BIGINT NOT NULL DEFAULT 0 COMMENT '世界ID(引用 world 基表,只读)', + `session_id` BIGINT NOT NULL DEFAULT 0 COMMENT '游戏会话ID(引用 scense 基表,只读)', + `entity_id` BIGINT NOT NULL DEFAULT 0 COMMENT '实体ID(引用 entity 基表,只读)', + `event_type` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '事件类型', + `payload` TEXT COMMENT '事件负载 JSON', + `seq_no` BIGINT NOT NULL DEFAULT 0 COMMENT '会话内单调序号', + `source` VARCHAR(32) NOT NULL DEFAULT 'runtime' COMMENT '来源:runtime/agent/script', + `state` VARCHAR(16) NOT NULL DEFAULT 'applied' COMMENT 'applied/rejected/rolled_back', + `tx_group` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '事务组(同组同事务)', + `broadcast` TINYINT NOT NULL DEFAULT 1 COMMENT '是否已广播 0/1', + `created_by` BIGINT NOT NULL DEFAULT 0 COMMENT '创建人', + `created_at` DATETIME COMMENT '创建时间', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_tenant_idem` (`tenant_id`, `idem_key`), + KEY `ix_tenant_session_seq` (`tenant_id`, `session_id`, `seq_no`), + KEY `ix_tenant_world` (`tenant_id`, `world_id`), + KEY `ix_tenant_type` (`tenant_id`, `event_type`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='运行时事件流(append-only)' +""", + """ +CREATE TABLE IF NOT EXISTS `pbl_entity_state` ( + `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键', + `tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户ID(多租户强制打头)', + `state_code` VARCHAR(64) NOT NULL COMMENT '状态编码', + `world_id` BIGINT NOT NULL DEFAULT 0 COMMENT '世界ID', + `session_id` BIGINT NOT NULL DEFAULT 0 COMMENT '会话ID', + `entity_id` BIGINT NOT NULL DEFAULT 0 COMMENT '实体ID(引用 entity 基表,只读)', + `state_key` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '状态键', + `state_value` TEXT COMMENT '状态值 JSON', + `version` INT NOT NULL DEFAULT 1 COMMENT '乐观锁版本', + `last_event_id` BIGINT NOT NULL DEFAULT 0 COMMENT '最后触发事件ID', + `created_at` DATETIME COMMENT '创建时间', + `updated_at` DATETIME COMMENT '更新时间', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_tenant_entity_key` (`tenant_id`, `session_id`, `entity_id`, `state_key`), + KEY `ix_tenant_world` (`tenant_id`, `world_id`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='实体状态薄扩展' +""", + """ +CREATE TABLE IF NOT EXISTS `pbl_world_state_snapshot` ( + `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键', + `tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户ID(多租户强制打头)', + `snapshot_code` VARCHAR(64) NOT NULL COMMENT '快照编码', + `world_id` BIGINT NOT NULL DEFAULT 0 COMMENT '世界ID', + `session_id` BIGINT NOT NULL DEFAULT 0 COMMENT '会话ID', + `seq_no` BIGINT NOT NULL DEFAULT 0 COMMENT '快照对应事件序号', + `state` TEXT COMMENT '世界状态 JSON', + `entity_count` INT NOT NULL DEFAULT 0 COMMENT '实体数', + `checksum` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '状态校验和', + `created_by` BIGINT NOT NULL DEFAULT 0 COMMENT '创建人', + `created_at` DATETIME COMMENT '创建时间', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_tenant_snapshot_code` (`tenant_id`, `snapshot_code`), + KEY `ix_tenant_session_seq` (`tenant_id`, `session_id`, `seq_no`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='世界状态快照' +""", +] + + +def get_ddl(): + """返回本模块全部建表语句(list[str])。""" + return [x.strip() for x in DDL] + + +def ensure_tables(env=None, sor=None): + """幂等建表(CREATE TABLE IF NOT EXISTS)。返回 (ok_list, bad_list)。""" + ok, bad = [], [] + runner = sor + if runner is None and env is not None: + runner = getattr(env, "sor", None) or getattr(env, "db", None) + for stmt in get_ddl(): + name = stmt.split("`")[1] if "`" in stmt else "?" + try: + if runner is not None and hasattr(runner, "sqlExe"): + runner.sqlExe(stmt) + elif env is not None and hasattr(env, "sqlExe"): + env.sqlExe(stmt) + else: + if "AUTO_INCREMENT" not in stmt or "tenant_id" not in stmt: + raise ValueError("DDL 形态不合规:%s" % name) + ok.append(name) + except Exception as exc: # noqa: BLE001 + bad.append((name, str(exc)[:160])) + return ok, bad + + +def self_check(): + """离线自检:3 表齐全 / tenant_id 打头 / 无禁用方言 / 幂等键唯一约束在位。""" + msgs = [] + all_ok = True + forbidden = ("BIGSERIAL", "SERIAL", "nextval", "FOREIGN KEY", + "REFERENCES", "ENUM(", "TIMESTAMP") + for stmt in get_ddl(): + name = stmt.split("`")[1] + body = stmt.upper() + if name not in TABLES: + all_ok = False + msgs.append("未知表 %s" % name) + for kw in forbidden: + if kw in body: + all_ok = False + msgs.append("%s 命中禁用方言 %s" % (name, kw)) + cols = [ln.strip().split("`")[1] for ln in stmt.splitlines() if ln.strip().startswith("`")] + biz = [c for c in cols if c != "id"] + if not biz or biz[0] != "tenant_id": + all_ok = False + msgs.append("%s 首个业务列=%s(应为 tenant_id)" % (name, biz[:1])) + if "AUTO_INCREMENT" not in body: + all_ok = False + msgs.append("%s 缺 AUTO_INCREMENT 主键" % name) + if "IF NOT EXISTS" not in body: + all_ok = False + msgs.append("%s 非幂等建表(缺 IF NOT EXISTS)" % name) + if "UNIQUE KEY `uk_tenant_idem`" not in DDL[0]: + all_ok = False + msgs.append("pbl_runtime_event 缺幂等唯一键 uk_tenant_idem") + if len(TABLES) != 3: + all_ok = False + msgs.append("表数=%d 应为 3" % len(TABLES)) + if all_ok: + msgs.append("SELF_CHECK pbl_runtime_ext.tables: PASS %d/%d" % (len(TABLES), len(TABLES))) + return all_ok, msgs diff --git a/pbl_runtime_ext/tx_event.py b/pbl_runtime_ext/tx_event.py index 0167fca..a69e7b9 100644 --- a/pbl_runtime_ext/tx_event.py +++ b/pbl_runtime_ext/tx_event.py @@ -1 +1,301 @@ -单事务核心(14,404B):EVENT_TYPES 18 项(与 pbl_runtime_event_type 对齐)/AUDIENCE 4 项;_next_seq 会话内单调序号;apply_event(单事务组装 stmts: START TRANSACTION→INSERT event→逐条 state_updates(SELECT 判存在→merge 或全量覆盖→UPDATE 带 state_version+1 / INSERT)→逐条 broadcast 登记(state='pending')→COMMIT;dry_run 返回 SQL 计划不落库;执行失败 ROLLBACK+fail 审计+PBL-RUNTIME-0001;成功后 write_audit + 事务外 _deliver_broadcasts);_deliver_broadcasts(pusher 缺失或抛错只标 failed+last_error,不向上抛);register_pusher/_get_pusher 通道无关注入;replay_failed_broadcasts(取 state=failed 且 retry_count