#!/usr/bin/env python3 # -*- coding: utf-8 -*- """pbl_runtime_event 表 DDL 单一权威生成器(M11b-1)。 范围(任务边界,勿越界):只负责 **schema 定义 + 按月 RANGE 分区 + append-only 约束** 的 DDL 生成与校验;不含事务写入逻辑、广播代码、轮询代码(那些属 M11b-2/-3, 在 pbl_runtime_ext/pbl_runtime_ext/{tx_write,partitions,broadcast}.py)。 归属模块(QC #22 决策 (b)):本表属 pbls 项目 generated_modules 中的 ``pbl_runtime_ext``(见 projects/pbls/pbls_spec.json 的 generated_modules 列表), **不写在共享模块 world_sync 里**(world_sync 只保留 W-05 同步职责,已回退无关改动)。 设计决策(详见 projects/pbls/docs/02-develop/dev-notes-m11b1-pbl_runtime_event.md): 1. 逻辑主键是 ``id``(AUTO_INCREMENT),但 MySQL/MariaDB 分区表要求 「每个唯一键必须包含全部分区列」,所以**物理主键是复合键 (id, created_at)**; ``id`` 仍然全局唯一(AUTO_INCREMENT 保证),只是不再单独作 PRIMARY KEY。 2. 分区键选 ``created_at``(DATETIME,不用 TIMESTAMP:分区表达式不允许 TIMESTAMP), 采用 ``PARTITION BY RANGE COLUMNS(created_at)`` + 命名月分区 p + 末尾 ``pmax VALUES LESS THAN (MAXVALUE)`` 兜底,预建 6 个月。 3. append-only 双保险: (a) 权限层 —— 迁移脚本 REVOKE app 账号的 UPDATE/DELETE(deploy 按需执行); (b) 触发器层 —— BEFORE UPDATE / BEFORE DELETE 触发器调用存储过程 ``pbl_assert_append_only()`` SIGNAL SQLSTATE '45000' 拒绝。 注:MySQL 系不允许在分区表上写视图/事件,但触发器是允许的; 护栏防的是「误用与 ORM 隐式写入」,不防拥有 SET 权限的 DBA (break-glass 变量 ``@pbl_guard_ctx`` 仅用于一次性历史数据回填)。 4. 与 sqlor ORM 的兼容性:sor.C() 只发 INSERT(不冲突);sor.U()/sor.D() 会被 触发器拒绝 —— 因此本模块所有写路径必须走 INSERT(tx_event/tx_write 已如此), 迁移脚本与生成器**不产生任何** UPDATE/DELETE 语句。 用法: python3 scripts/pbl_runtime_event_ddl.py --print # 打印建表 DDL python3 scripts/pbl_runtime_event_ddl.py --emit [目录] # 生成迁移 SQL(幂等) python3 scripts/pbl_runtime_event_ddl.py --verify # 离线静态自检 python3 scripts/pbl_runtime_event_ddl.py --apply # 连库执行(读 env/*.json) python3 scripts/pbl_runtime_event_ddl.py --verify-db # 连库验证分区裁剪/落分区/禁改禁删 退出码:0 成功 / 1 校验失败 / 2 环境不可用(DB 连不上,如实上报不谎报)。 """ import argparse import json import os import re import sys from datetime import date # 相对路径基准:本文件 = modules/pbl_runtime_ext/scripts/xxx.py → 上溯 3 级 = 机构工作空间 HERE = os.path.dirname(os.path.abspath(__file__)) MODULE_DIR = os.path.dirname(HERE) # modules/pbl_runtime_ext WORKSPACE = os.path.dirname(os.path.dirname(MODULE_DIR)) # 机构工作空间根 DEV_NOTES = "projects/pbls/docs/02-develop/dev-notes-m11b1-pbl_runtime_event.md" TABLE = "pbl_runtime_event" PARTITION_COLUMN = "created_at" PARTITION_PREFIX = "p" # p202609 ... PARTITION_MAXVALUE = "pmax" MONTHS_PRECREATE = 6 # 任务要求:预建 >= 6 个月分区 RETENTION_MONTHS = 14 # 月分区保留数(超出由维护过程 DROP PARTITION) # ---------------------------------------------------------------- 列定义(唯一真源) # (列名, 类型串, COMMENT);顺序 = 物理列顺序 COLUMNS = [ ("id", "BIGINT NOT NULL AUTO_INCREMENT", "自增ID(逻辑主键;物理主键=(id,created_at))"), ("tenant_id", "VARCHAR(32) NOT NULL DEFAULT ''", "租户ID(多租户强制打头)"), ("event_uid", "VARCHAR(64) NOT NULL", "事件UID(uuid,应用层生成)"), ("event_code", "VARCHAR(64) NOT NULL DEFAULT ''", "事件编码"), ("idem_key", "VARCHAR(128) NOT NULL DEFAULT ''", "幂等键(同键重复投递只落一条)"), ("world_id", "BIGINT NOT NULL DEFAULT 0", "世界ID(引用 world 基表,只读)"), ("scene_id", "BIGINT NOT NULL DEFAULT 0", "场景/会话ID(引用 scense 基表,只读)"), ("session_id", "BIGINT NOT NULL DEFAULT 0", "游戏会话ID(M11a 语义,与 scene_id 同域)"), ("entity_id", "VARCHAR(64) NOT NULL DEFAULT ''", "实体ID(引用 entity 基表,只读)"), ("actor_id", "VARCHAR(64) NOT NULL DEFAULT ''", "触发者(服务端会话解析,客户端禁填)"), ("event_type", "VARCHAR(64) NOT NULL", "事件类型"), ("payload", "TEXT NULL", "事件负载JSON(任务书字段)"), ("payload_json", "LONGTEXT NULL", "事件负载JSON(M11a 写路径列名,与 payload 同义留宽)"), ("causation_id", "VARCHAR(64) NOT NULL DEFAULT ''", "因果链:引发本事件的事件UID"), ("seq", "BIGINT NOT NULL DEFAULT 0", "会话内单调序号(轮询游标,应用层生成)"), ("seq_no", "BIGINT NOT NULL DEFAULT 0", "会话内单调序号(M11a 写路径列名)"), ("source", "VARCHAR(32) NOT NULL DEFAULT 'runtime'", "来源:runtime/agent/script/client_intent"), ("state", "VARCHAR(16) NOT NULL DEFAULT 'applied'", "applied/rejected/rolled_back"), ("state_version", "BIGINT NOT NULL DEFAULT 0", "本事件推进到的服务端权威版本"), ("tx_group", "VARCHAR(64) NOT NULL DEFAULT ''", "事务组(同组同事务)"), ("broadcast", "TINYINT NOT NULL DEFAULT 1", "是否参与广播 0/1"), ("created_by", "BIGINT NOT NULL DEFAULT 0", "创建人"), ("created_at", "DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP", "创建时间(分区键,禁TIMESTAMP)"), ("updated_at", "DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP", "记录时间戳(append-only:永不 UPDATE)"), ] # 索引:(索引名, 类型, [列]);分区表铁律——唯一键必须包含全部分区列 UNIQUE_KEYS = [ ("uk_event_uid", ["tenant_id", "event_uid", PARTITION_COLUMN]), ("uk_tenant_idem", ["tenant_id", "idem_key", PARTITION_COLUMN]), ] SECONDARY_KEYS = [ ("ix_scene_created", ["scene_id", PARTITION_COLUMN]), # 任务书要求的复合索引 ("ix_tenant_session_seq", ["tenant_id", "session_id", "seq"]), ("ix_tenant_session_seqno", ["tenant_id", "session_id", "seq_no"]), ("ix_tenant_world_type", ["tenant_id", "world_id", "event_type"]), ("ix_tx_group", ["tx_group"]), ("ix_created_at", [PARTITION_COLUMN]), ] PRIMARY_KEY = ["id", PARTITION_COLUMN] # M11a/M11b 现有写路径 INSERT 的列(自检据此断言 DDL 覆盖,防列名漂移) INSERT_COLUMNS = ["tenant_id", "event_uid", "session_id", "world_id", "actor_id", "event_type", "payload_json", "causation_id", "state_version", "seq", "occurred_at", "created_at"] # tx_event.py(M11a 单事务写路径)使用的列 TX_EVENT_COLUMNS = ["tenant_id", "event_code", "idem_key", "world_id", "session_id", "entity_id", "event_type", "payload", "seq_no", "source", "state", "tx_group", "broadcast", "created_by", "created_at", "updated_at"] # occurred_at 是 M11a 写路径与 partitions.py 曾用的分区列名 —— 迁移时必须一并落列, # 否则旧 INSERT 直接报错。保留为兼容列(分区键已改为 created_at)。 LEGACY_COLUMNS = ["occurred_at"] GUARD_ADMIN_CTX = "append_only_admin" def column_names(): return [c[0] for c in COLUMNS] + LEGACY_COLUMNS # ---------------------------------------------------------------- DDL 片段生成 def month_start(d): return date(d.year, d.month, 1) def add_months(d, n): total = d.year * 12 + (d.month - 1) + n return date(total // 12, total % 12 + 1, 1) def partition_plan(today=None, months=None): """返回 [(分区名, 上界 date), ...],含当月共 months 个月(>=6)。""" today = today or date.today() months = max(1, int(months or MONTHS_PRECREATE)) first = month_start(today) out = [] for i in range(months): ms = add_months(first, i) out.append(("%s%s" % (PARTITION_PREFIX, ms.strftime("%Y%m")), add_months(ms, 1))) return out def create_table_sql(today=None, months=None): """建表 DDL(含 RANGE COLUMNS 月分区 + pmax 兜底)。幂等:IF NOT EXISTS。""" lines = ["CREATE TABLE IF NOT EXISTS `%s` (" % TABLE] body = [] for name, typ, comment in COLUMNS: body.append(" `%s` %s COMMENT '%s'" % (name, typ, comment)) for legacy in LEGACY_COLUMNS: body.append(" `%s` DATETIME NULL COMMENT '兼容列(M11a 旧分区键,新写入可空)'" % legacy) body.append(" PRIMARY KEY (%s)" % ", ".join("`%s`" % c for c in PRIMARY_KEY)) for name, cols in UNIQUE_KEYS: body.append(" UNIQUE KEY `%s` (%s)" % (name, ", ".join("`%s`" % c for c in cols))) for name, cols in SECONDARY_KEYS: body.append(" KEY `%s` (%s)" % (name, ", ".join("`%s`" % c for c in cols))) lines.append(",\n".join(body)) lines.append(") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci") lines.append(" COMMENT='运行时事件流(append-only,按月RANGE分区,保留%d个月)' " % RETENTION_MONTHS) plan = partition_plan(today, months) parts = [" PARTITION `%s` VALUES LESS THAN (DATE '%s')" % (n, b.isoformat()) for n, b in plan] parts.append(" PARTITION `%s` VALUES LESS THAN (MAXVALUE)" % PARTITION_MAXVALUE) lines.append("PARTITION BY RANGE COLUMNS(`%s`) (\n%s\n);" % (PARTITION_COLUMN, ",\n".join(parts))) return "\n".join(lines) def append_only_sql(): """append-only 约束:存储过程 + BEFORE UPDATE/DELETE 触发器(幂等)。""" return """-- append-only 断言过程:非白名单操作一律 SIGNAL 拒绝 DROP PROCEDURE IF EXISTS `pbl_assert_append_only`; CREATE PROCEDURE `pbl_assert_append_only`( IN p_table VARCHAR(64), IN p_priv VARCHAR(16), IN p_detail VARCHAR(255)) READS SQL DATA BEGIN IF p_table <> '%(table)s' THEN RETURN; -- 只守护事件表,其他表不受影响 END IF; IF IFNULL(@pbl_guard_ctx, '') = '%(ctx)s' THEN RETURN; -- break-glass:一次性历史数据回填 END IF; SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'PBL_APPEND_ONLY_DENIED', MYSQL_ERRNO = 45001; END; -- BEFORE UPDATE 护栏:把违规语句原文带进错误消息,便于定位调用方 DROP TRIGGER IF EXISTS `trg_%(table)s_append_only_upd`; CREATE TRIGGER `trg_%(table)s_append_only_upd` BEFORE UPDATE ON `%(table)s` FOR EACH ROW BEGIN DECLARE v_sql TEXT; DECLARE v_msg VARCHAR(255); SET v_sql = ''; BEGIN DECLARE CONTINUE HANDLER FOR NOT FOUND SET v_sql = ''; SELECT INFO INTO v_sql FROM information_schema.processlist WHERE ID = CONNECTION_ID(); END; IF IFNULL(@pbl_guard_ctx, '') = '%(ctx)s' THEN SET v_msg = CONCAT('break-glass update on ', '%(table)s'); ELSE SET v_msg = CONCAT('UPDATE denied: ', LEFT(IFNULL(v_sql, 'n/a'), 120)); END IF; CALL pbl_assert_append_only('%(table)s', 'update', v_msg); IF IFNULL(@pbl_guard_ctx, '') <> '%(ctx)s' THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'PBL_APPEND_ONLY_DENIED_UPDATE', MYSQL_ERRNO = 45002; END IF; END; -- BEFORE DELETE 护栏 DROP TRIGGER IF EXISTS `trg_%(table)s_append_only_del`; CREATE TRIGGER `trg_%(table)s_append_only_del` BEFORE DELETE ON `%(table)s` FOR EACH ROW BEGIN DECLARE v_sql TEXT; DECLARE v_msg VARCHAR(255); SET v_sql = ''; BEGIN DECLARE CONTINUE HANDLER FOR NOT FOUND SET v_sql = ''; SELECT INFO INTO v_sql FROM information_schema.processlist WHERE ID = CONNECTION_ID(); END; IF IFNULL(@pbl_guard_ctx, '') = '%(ctx)s' THEN SET v_msg = CONCAT('break-glass delete on ', '%(table)s'); ELSE SET v_msg = CONCAT('DELETE denied: ', LEFT(IFNULL(v_sql, 'n/a'), 120)); END IF; CALL pbl_assert_append_only('%(table)s', 'delete', v_msg); IF IFNULL(@pbl_guard_ctx, '') <> '%(ctx)s' THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'PBL_APPEND_ONLY_DENIED_DELETE', MYSQL_ERRNO = 45003; END IF; END; -- 权限层加固(可选,需 DBA 账号执行;app 账号只保留增/查) -- GRANT INSERT, SELECT, UPDATE, DELETE ON `%(db)s`.`%(table)s` TO NULL; -- 占位,勿执行 -- REVOKE UPDATE, DELETE ON `%(db)s`.`%(table)s` FROM '%(app_user)s'@'%%'; """ % {"table": TABLE, "ctx": GUARD_ADMIN_CTX, "db": "%(dbname)s", "app_user": "%(app_user)s"} def maintenance_sql(today=None, months=None): """分区维护:过程(幂等补建月分区 / 到期 DROP 旧分区)+ 每月 EVENT。 QC #1 修复:正文含大量字面百分号(DATE_FORMAT(...,'%Y-%m-01')、'%Y%m'), 原先用「三引号 % 字典」渲染会抛 ValueError: unsupported format character。 现改用 __TABLE__/__PMAX__/__MONTHS__ token + str.replace, 字面百分号原样输出、无需再写双写转义。 """ months = int(months or MONTHS_PRECREATE) body = """-- 1) 幂等补建未来月分区:pmax 存在时用 REORGANIZE PARTITION 拆出新月份(不丢数据) DROP PROCEDURE IF EXISTS `pbl_runtime_event_ensure_partitions`; CREATE PROCEDURE `pbl_runtime_event_ensure_partitions`(IN p_months_ahead INT) MODIFIES SQL DATA BEGIN DECLARE v_done INT DEFAULT 0; DECLARE v_name VARCHAR(16); DECLARE v_upper DATE; DECLARE v_exists INT DEFAULT 0; DECLARE v_sql TEXT; DECLARE cur CURSOR FOR WITH RECURSIVE ms(month_start) AS ( SELECT DATE_FORMAT(CURDATE(), '%Y-%m-01') UNION ALL SELECT DATE_ADD(month_start, INTERVAL 1 MONTH) FROM ms WHERE month_start < DATE_ADD(DATE_FORMAT(CURDATE(), '%Y-%m-01'), INTERVAL p_months_ahead MONTH) ) SELECT CONCAT('p', DATE_FORMAT(month_start, '%Y%m')), DATE_ADD(month_start, INTERVAL 1 MONTH) FROM ms; DECLARE CONTINUE HANDLER FOR NOT FOUND SET v_done = 1; OPEN cur; read_loop: LOOP FETCH cur INTO v_name, v_upper; IF v_done = 1 THEN LEAVE read_loop; END IF; SELECT COUNT(*) INTO v_exists FROM information_schema.PARTITIONS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = '__TABLE__' AND PARTITION_NAME = v_name; IF v_exists = 0 THEN SET @pbl_part_sql = CONCAT( 'ALTER TABLE `__TABLE__` REORGANIZE PARTITION `__PMAX__` INTO (', 'PARTITION `', v_name, '` VALUES LESS THAN (DATE ''', v_upper, '''), ', 'PARTITION `__PMAX__` VALUES LESS THAN (MAXVALUE))'); PREPARE pbl_part_st FROM @pbl_part_sql; EXECUTE pbl_part_st; DEALLOCATE PREPARE pbl_part_st; END IF; END LOOP; CLOSE cur; END; -- 2) 到期归档:整月 DROP PARTITION(append-only 表唯一合规的清理手段) DROP PROCEDURE IF EXISTS `pbl_runtime_event_prune_partitions`; CREATE PROCEDURE `pbl_runtime_event_prune_partitions`(IN p_keep_months INT) MODIFIES SQL DATA BEGIN DECLARE v_done INT DEFAULT 0; DECLARE v_name VARCHAR(64); DECLARE v_cutoff DATE; DECLARE cur CURSOR FOR SELECT PARTITION_NAME FROM information_schema.PARTITIONS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = '__TABLE__' AND PARTITION_NAME IS NOT NULL AND PARTITION_NAME <> '__PMAX__' AND PARTITION_NAME < CONCAT('p', DATE_FORMAT(v_cutoff, '%Y%m')); DECLARE CONTINUE HANDLER FOR NOT FOUND SET v_done = 1; SET v_cutoff = DATE_ADD(DATE_FORMAT(CURDATE(), '%Y-%m-01'), INTERVAL -p_keep_months MONTH); OPEN cur; prune_loop: LOOP FETCH cur INTO v_name; IF v_done = 1 THEN LEAVE prune_loop; END IF; SET @pbl_drop_sql = CONCAT('ALTER TABLE `__TABLE__` DROP PARTITION `', v_name, '`'); PREPARE pbl_drop_st FROM @pbl_drop_sql; EXECUTE pbl_drop_st; DEALLOCATE PREPARE pbl_drop_st; END LOOP; CLOSE cur; END; -- 3) 每月自动补分区(需 event_scheduler=ON;否则用 cron,见开发说明) -- SET GLOBAL event_scheduler = ON; DROP EVENT IF EXISTS `ev___TABLE___partition_maintain`; CREATE EVENT IF NOT EXISTS `ev___TABLE___partition_maintain` ON SCHEDULE EVERY 1 MONTH STARTS TIMESTAMP(DATE_FORMAT(DATE_ADD(CURDATE(), INTERVAL 1 MONTH), '%Y-%m-01 03:00:00')) ON COMPLETION PRESERVE ENABLE COMMENT '每月预建未来 __MONTHS__ 个月分区' DO CALL pbl_runtime_event_ensure_partitions(__MONTHS__); """ return (body .replace("__TABLE__", TABLE) .replace("__PMAX__", PARTITION_MAXVALUE) .replace("__MONTHS__", str(months))) def migrate_from_legacy_sql(): """M11a 非分区表 → M11b 分区表的一次性改造(幂等:已分区则跳过)。""" return """-- 一次性改造(仅当 pbl_runtime_event 已存在且未分区时手工执行;先备份): -- 1) CREATE TABLE pbl_runtime_event_new ... (本文件 A 段 DDL,改名 _new) -- 2) SET @pbl_guard_ctx = '%(ctx)s'; -- break-glass,仅回填期间 -- INSERT INTO pbl_runtime_event_new (...) SELECT ... FROM pbl_runtime_event; -- SET @pbl_guard_ctx = NULL; -- 3) RENAME TABLE pbl_runtime_event TO pbl_runtime_event_m11a_bak, -- pbl_runtime_event_new TO pbl_runtime_event; -- 4) 重放 append-only 触发器段(B 段),核对 EXPLAIN 分区裁剪后删备份表。 -- 本脚本不自动执行改造:DROP/RENAME 生产表必须人工确认(见开发说明「回滚」段)。 SELECT 'm11b-1 migrate: 人工执行,见注释' AS note; """ % {"ctx": GUARD_ADMIN_CTX} def full_script(dbname="pbls", app_user="pbls", today=None, months=None): ddl = create_table_sql(today, months) guard = append_only_sql().replace("%(dbname)s", dbname).replace("%(app_user)s", app_user) maint = maintenance_sql(today, months) head = ("-- ===========================================================================\n" "-- %s 表 DDL + 按月 RANGE 分区 + append-only 约束(M11b-1)\n" "-- 由 scripts/pbl_runtime_event_ddl.py 生成,改表结构只改生成器再 --emit。\n" "-- 目标引擎 mariadb(projects/pbls/env/test.json: db.engine=mariadb)。\n" "-- 执行:mysql --protocol=tcp -u -p < 本文件\n" "-- 说明:触发器含多语句体,须用 mysql CLI(DELIMITER)执行;\n" "-- 应用侧 ensure_tables() 只执行 A 段单语句 CREATE TABLE。\n" "-- 开发说明:%s\n" "-- ===========================================================================\n\n" % (TABLE, DEV_NOTES)) return (head + "-- ---------------------------------------------------------------- A. 建表(分区)\n" + ddl + "\n\n" + "DELIMITER $$\n\n" + "-- ---------------------------------------------------------------- B. append-only\n" + _with_delimiter(guard) + "\n-- ---------------------------------------------------------------- C. 分区维护\n" + _with_delimiter(maint) + "\nDELIMITER ;\n\n" + "-- ---------------------------------------------------------------- D. 一次性改造\n" + migrate_from_legacy_sql()) def _with_delimiter(body): """把过程/触发器体内语句交给 mysql CLI:整段以 $$ 结束。""" out = [] for chunk in re.split(r"(?<=;)\n(?=(DROP|CREATE|GRANT|REVOKE|SELECT|--|SET))", body): c = chunk.strip() if not c: continue if c.startswith("--") or c.startswith("SET ") or c.startswith("GRANT") or c.startswith("REVOKE"): out.append(c if c.endswith(";") else c + ";") elif c.startswith("DROP PROCEDURE") or c.startswith("DROP TRIGGER") or c.startswith("DROP EVENT"): out.append(c if c.endswith(";") else c + ";") else: stmt = c.rstrip(";") out.append(stmt + "$$") return "\n\n".join(out) + "\n" # ---------------------------------------------------------------- 迁移文件产出 MIGRATION_FILES = [ ("20260919_m11b_01_create_pbl_runtime_event.sql", lambda db: create_table_sql() + "\n"), ("20260919_m11b_02_pbl_runtime_event_append_only.sql", lambda db: "DELIMITER $$\n\n" + _with_delimiter( append_only_sql().replace("%(dbname)s", db).replace("%(app_user)s", db)) + "\nDELIMITER ;\n"), ("20260919_m11b_03_pbl_runtime_event_partition_maintenance.sql", lambda db: "DELIMITER $$\n\n" + _with_delimiter(maintenance_sql()) + "\nDELIMITER ;\n"), ] def emit(out_dir=None, dbname="pbls"): out_dir = out_dir or os.path.join(MODULE_DIR, "scripts", "migrations") os.makedirs(out_dir, exist_ok=True) written = [] for name, fn in MIGRATION_FILES: path = os.path.join(out_dir, name) with open(path, "w", encoding="utf-8") as f: f.write("-- 生成器: scripts/pbl_runtime_event_ddl.py --emit (M11b-1)\n" "-- 开发说明: %s\n-- 幂等: CREATE ... IF NOT EXISTS / DROP ... IF EXISTS\n\n" % DEV_NOTES) f.write(fn(dbname)) written.append(path) canon = os.path.join(MODULE_DIR, "sql", "m11b_partitions.sql") os.makedirs(os.path.dirname(canon), exist_ok=True) with open(canon, "w", encoding="utf-8") as f: f.write(full_script(dbname=dbname)) written.append(canon) return written # ---------------------------------------------------------------- 离线静态自检 def verify(): """静态断言:返回 (ok, msgs)。不连库,可离线取证。""" msgs = [] ok = True ddl = create_table_sql() cols = column_names() # 1. 任务书字段齐备 for need in ["event_id", "scene_id", "entity_id", "event_type", "payload", PARTITION_COLUMN, "seq"]: hit = need in cols or (need == "event_id" and "event_uid" in cols) if not hit: ok = False msgs.append("FAIL 缺字段 %s(event_id 由 event_uid 承载)" % need) msgs.append("PASS 任务书字段齐备: %s" % ",".join(cols)) # 2. 物理主键含分区键;唯一键含分区键 if PARTITION_COLUMN not in PRIMARY_KEY: ok = False msgs.append("FAIL 物理主键未包含分区键 %s" % PARTITION_COLUMN) else: msgs.append("PASS 物理主键=%s(逻辑主键 id + 分区键)" % ",".join(PRIMARY_KEY)) for name, kcols in UNIQUE_KEYS: if PARTITION_COLUMN not in kcols: ok = False msgs.append("FAIL 唯一键 %s 未含分区键" % name) msgs.append("PASS 唯一键均含分区键: %s" % ",".join(n for n, _ in UNIQUE_KEYS)) # 3. 索引要求:scene_id+created_at 复合索引、event_uid 唯一索引 if not any(set(k) == {"scene_id", PARTITION_COLUMN} for _, k in SECONDARY_KEYS): ok = False msgs.append("FAIL 缺 scene_id+created_at 复合索引") else: msgs.append("PASS 存在 (scene_id, created_at) 复合索引 ix_scene_created") if not any(k[:2] == ["tenant_id", "event_uid"] for _, k in UNIQUE_KEYS): ok = False msgs.append("FAIL 缺 event_uid 唯一索引") else: msgs.append("PASS 存在 event_uid 唯一索引 uk_event_uid") # 4. 分区:>=6 个月 + MAXVALUE 兜底 + 分区键为 DATETIME plan = partition_plan() if len(plan) < MONTHS_PRECREATE: ok = False msgs.append("FAIL 预建分区数 %d < %d" % (len(plan), MONTHS_PRECREATE)) else: msgs.append("PASS 预建 %d 个月分区: %s + %s(MAXVALUE)" % (len(plan), ",".join(n for n, _ in plan), PARTITION_MAXVALUE)) if "PARTITION BY RANGE COLUMNS(`%s`)" % PARTITION_COLUMN not in ddl: ok = False msgs.append("FAIL 未声明 RANGE COLUMNS 分区") if "VALUES LESS THAN (MAXVALUE)" not in ddl: ok = False msgs.append("FAIL 缺 MAXVALUE 兜底分区") if re.search(r"`%s`\s+TIMESTAMP" % PARTITION_COLUMN, ddl): ok = False msgs.append("FAIL 分区键用了 TIMESTAMP(MySQL 不允许)") else: msgs.append("PASS 分区键类型 DATETIME(非 TIMESTAMP)") # 5. append-only:DDL/迁移全文不得出现对本表的 UPDATE/DELETE 语句 all_sql = ddl + append_only_sql() + maintenance_sql() for bad in [r"UPDATE\s+`%s`" % TABLE, r"DELETE\s+FROM\s+`%s`" % TABLE]: if re.search(bad, all_sql, re.I): ok = False msgs.append("FAIL 生成的 SQL 含违规写语句 %s" % bad) if "BEFORE UPDATE" not in append_only_sql() or "BEFORE DELETE" not in append_only_sql(): ok = False msgs.append("FAIL append-only 触发器缺失") else: msgs.append("PASS append-only: BEFORE UPDATE + BEFORE DELETE + pbl_assert_append_only()") # 6. 与既有写路径列名兼容(防 DDL 漂移把 M11a/M11b 代码打挂) for col in set(INSERT_COLUMNS) | set(TX_EVENT_COLUMNS): if col not in cols: ok = False msgs.append("FAIL 写路径引用列 %s 不在 DDL 中" % col) msgs.append("PASS 覆盖 M11a/M11b 写路径全部列名(%d 个)" % len(set(INSERT_COLUMNS) | set(TX_EVENT_COLUMNS))) # 7. 幂等性:所有 DDL 语句可重复执行 idem_ok = ("CREATE TABLE IF NOT EXISTS" in ddl and "DROP PROCEDURE IF EXISTS" in append_only_sql() and "DROP TRIGGER IF EXISTS" in append_only_sql() and "DROP EVENT IF EXISTS" in maintenance_sql()) if not idem_ok: ok = False msgs.append("FAIL DDL 非幂等(缺 IF EXISTS / DROP IF EXISTS)") else: msgs.append("PASS 幂等:IF NOT EXISTS + DROP ... IF EXISTS + 分区存在性判断") return ok, msgs # ---------------------------------------------------------------- 连库执行 / 验证 def _db_conf(): """读项目唯一事实源 env/test.json(禁止硬编码连接串)。""" for env_name in ("test", "prod"): path = os.path.join(WORKSPACE, "projects", "pbls", "env", "%s.json" % env_name) if os.path.exists(path): with open(path, encoding="utf-8") as f: conf = json.load(f) db = conf.get("db") or {} if db.get("host"): return env_name, db return None, None def _connect(): import pymysql env_name, db = _db_conf() if not db: raise RuntimeError("未找到 projects/pbls/env/*.json 的 db 段") return pymysql.connect(host=db["host"], port=int(db.get("port") or 3306), user=db.get("user"), password=db.get("password"), database=db.get("dbname"), charset="utf8mb4", autocommit=True), env_name, db def apply_ddl(): """连库执行 A 段建表 + B 段 append-only + C 段维护(逐语句,幂等)。""" conn, env_name, db = _connect() cur = conn.cursor() statements = [create_table_sql()] statements += _split_procedural(append_only_sql() .replace("%(dbname)s", db.get("dbname") or "pbls") .replace("%(app_user)s", db.get("user") or "pbls")) statements += _split_procedural(maintenance_sql()) done, failed = [], [] for stmt in statements: s = stmt.strip().rstrip(";").strip() if not s or s.startswith("--"): continue try: cur.execute(s) done.append(s.split("\n")[0][:70]) except Exception as exc: # noqa: BLE001 failed.append((s.split("\n")[0][:70], str(exc)[:200])) cur.close(); conn.close() print("[apply] env=%s db=%s 成功 %d 条,失败 %d 条" % (env_name, db.get("dbname"), len(done), len(failed))) for f in failed: print(" FAIL %s -> %s" % f) return 0 if not failed else 1 def _split_procedural(body): """按 DDL 边界切语句:CREATE PROCEDURE/TRIGGER/EVENT 取到结尾分号(含内部 ;)。""" out, buf, inside = [], "", False for line in body.splitlines(): low = line.strip().lower() if not inside and (low.startswith("create procedure") or low.startswith("create trigger") or low.startswith("create event")): inside = True buf += line + "\n" if inside and line.strip() in ("END;", "END;$$", ";"): out.append(buf); buf, inside = "", False elif not inside and low.endswith(";"): out.append(buf); buf = "" if buf.strip(): out.append(buf) return [x.replace("$$", "").rstrip() for x in out] def _explain_partitions(cur, sql, params=None): """EXPLAIN FORMAT=JSON 取执行计划实际扫描的分区清单(QC #6)。 不再用 PARTITIONS 扩展语法(EXPLAIN PARTITIONS):其分区列固定在结果 row[3],列序随 MariaDB/MySQL 版本变化,且该扩展语法在新版 MySQL 已废弃。JSON 计划里的 partitions 字段两版一致。 取不到(引擎不支持/解析失败)返回 [],由调用方判失败,不猜。 """ try: cur.execute("EXPLAIN FORMAT=JSON " + sql, tuple(params or ())) row = cur.fetchone() except Exception: # noqa: BLE001 return [] raw = row[0] if row else None if isinstance(raw, (bytes, bytearray)): raw = raw.decode("utf-8", "replace") if not raw: return [] try: plan = json.loads(raw) except (ValueError, TypeError): return [] found = [] def walk(node): if isinstance(node, dict): for key, val in node.items(): if key == "partitions" and isinstance(val, list): for p in val: if isinstance(p, str) and p and p not in found: found.append(p) else: walk(val) elif isinstance(node, list): for val in node: walk(val) walk(plan) return found def _row_in_partition(cur, part, uid): """SELECT ... FROM t PARTITION(p) 反查:探针行是否物理落在该分区。""" try: cur.execute("SELECT id FROM `%s` PARTITION (`%s`)" " WHERE tenant_id='t_probe' AND event_uid=%%s" % (TABLE, part), (uid,)) except Exception: # noqa: BLE001 return False return bool(cur.fetchone()) def _locate_row_partition(cur, uid, ts, part_names): """定位某探针行实际所在分区(QC #6)。 1) 首选按 information_schema 分区清单逐个 PARTITION(p) 反查物理位置——与 EXPLAIN 输出列序无关,最可靠; 2) 反查不中(引擎不支持 PARTITION 子句等)时退回 EXPLAIN FORMAT=JSON 解析 partitions(带 tenant_id+event_uid+created_at 等值条件,裁剪后应只剩 1 个分区)。 """ for part in part_names: if _row_in_partition(cur, part, uid): return part planned = _explain_partitions(cur, "SELECT id FROM `%s` WHERE tenant_id='t_probe'" " AND event_uid=%%s AND created_at=%%s" % TABLE, (uid, ts)) if len(planned) == 1: return planned[0] return ",".join(planned) or "?" def verify_db(): """连库验证:分区裁剪 / 数据落分区 / UPDATE·DELETE 被拒。""" conn, env_name, db = _connect() cur = conn.cursor() results = [] cur.execute("SELECT COUNT(*) FROM information_schema.PARTITIONS WHERE TABLE_SCHEMA=DATABASE()" " AND TABLE_NAME=%s AND PARTITION_NAME IS NOT NULL", (TABLE,)) n_parts = cur.fetchone()[0] results.append(("分区数", n_parts >= MONTHS_PRECREATE + 1, "实际 %d(含 pmax)" % n_parts)) # 分区清单:以 information_schema.PARTITIONS 为权威来源(按 ordinal 排序), # 后面「落哪个分区」「分区裁剪」两项判定都基于这份清单(QC #6) cur.execute("SELECT PARTITION_NAME, PARTITION_DESCRIPTION FROM information_schema.PARTITIONS " "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=%s " "AND PARTITION_NAME IS NOT NULL ORDER BY PARTITION_ORDINAL_POSITION", (TABLE,)) part_rows = [(r[0], r[1]) for r in cur.fetchall()] part_names = [n for n, _ in part_rows if n] results.append(("MAXVALUE 兜底", PARTITION_MAXVALUE in part_names, ",".join(part_names))) # 插入两条不同月份的探针行 → 应落在两个不同的物理分区 probe = "m11bprobe%s" % date.today().strftime("%H%M%S") probes = (("%s_a" % probe, "2026-01-05 00:00:00"), ("%s_b" % probe, "2027-03-05 00:00:00")) try: for i, (uid, ts) in enumerate(probes): cur.execute("INSERT INTO `%s` (tenant_id,event_uid,event_type,scene_id,seq,created_at)" " VALUES ('t_probe',%%s,'probe',0,%%s,%%s)" % TABLE, (uid, i + 1, ts)) except Exception as exc: # noqa: BLE001 results.append(("跨月插入", False, str(exc)[:160])) else: results.append(("跨月插入", True, "%d 条探针行(不同月份)" % len(probes))) located = [_locate_row_partition(cur, uid, ts, part_names) for uid, ts in probes] distinct = len(set(located)) == 2 in_list = all(p in part_names for p in located) results.append(("跨月落不同分区", distinct and in_list, " | ".join("%s -> %s" % (uid, p) for (uid, _), p in zip(probes, located)))) # 分区裁剪:created_at 区间只应命中该区间所属的少数分区,而非全部分区。 # 旧写法(PARTITIONS 扩展语法 + 取 row[3]):列序跨版本不稳定,新版 MySQL 已废弃该语法。 pruned = _explain_partitions(cur, "SELECT id FROM `%s` WHERE created_at>='2026-01-01'" " AND created_at<'2026-02-01'" % TABLE) results.append(("分区裁剪", bool(pruned) and all(p in part_names for p in pruned) and len(pruned) < max(1, len(part_names)), "命中 %d/%d 分区: %s" % (len(pruned), len(part_names), ",".join(pruned) or "N/A"))) # append-only:UPDATE / DELETE 必须被拒 for verb, sql in (("UPDATE", "UPDATE `%s` SET event_type='x' WHERE tenant_id='t_probe'" % TABLE), ("DELETE", "DELETE FROM `%s` WHERE tenant_id='t_probe'" % TABLE)): try: cur.execute(sql) results.append(("%s 被拒" % verb, False, "居然执行成功了(护栏未生效)")) except Exception as exc: # noqa: BLE001 msg = str(exc) results.append(("%s 被拒" % verb, "APPEND_ONLY" in msg or "45000" in msg or "4500" in msg, msg[:160])) # 清理探针(break-glass,仅验证脚本使用) try: cur.execute("SET @pbl_guard_ctx='%s'" % GUARD_ADMIN_CTX) cur.execute("DELETE FROM `%s` WHERE tenant_id='t_probe'" % TABLE) cur.execute("SET @pbl_guard_ctx=NULL") except Exception as exc: # noqa: BLE001 results.append(("探针清理", False, str(exc)[:120])) ok = True print("[verify-db] env=%s db=%s" % (env_name, db.get("dbname"))) for name, passed, detail in results: ok = ok and passed print(" %s %-16s %s" % ("PASS" if passed else "FAIL", name, detail)) cur.close(); conn.close() return 0 if ok else 1 def main(argv=None): ap = argparse.ArgumentParser(description="pbl_runtime_event DDL 生成器(M11b-1)") ap.add_argument("--print", action="store_true", help="打印建表 DDL") ap.add_argument("--print-guard", action="store_true", help="打印 append-only 段") ap.add_argument("--print-maintain", action="store_true", help="打印分区维护段") ap.add_argument("--emit", nargs="?", const="", default=None, help="生成迁移 SQL 到目录") ap.add_argument("--dbname", default="pbls") ap.add_argument("--verify", action="store_true", help="离线静态自检") ap.add_argument("--apply", action="store_true", help="连库执行(读 env/*.json)") ap.add_argument("--verify-db", action="store_true", help="连库验证分区/append-only") args = ap.parse_args(argv) if args.print: print(create_table_sql()) if args.print_guard: print(append_only_sql()) if args.print_maintain: print(maintenance_sql()) if args.emit is not None: for p in emit(args.emit or None, args.dbname): print("WROTE %s (%d bytes)" % (p, os.path.getsize(p))) if args.verify: ok, msgs = verify() for m in msgs: print(m) print("VERIFY %s" % ("PASS" if ok else "FAIL")) return 0 if ok else 1 if args.apply: try: return apply_ddl() except Exception as exc: # noqa: BLE001 print("APPLY SKIP(环境受限): %s" % str(exc)[:200]) return 2 if getattr(args, "verify_db"): try: return verify_db() except Exception as exc: # noqa: BLE001 print("VERIFY-DB SKIP(环境受限): %s" % str(exc)[:200]) return 2 if not any([args.print, args.print_guard, args.print_maintain, args.emit is not None]): print(__doc__) return 0 if __name__ == "__main__": sys.exit(main())