From d1f5c7cb42180b18309a3028eba77b0d0f2a16de Mon Sep 17 00:00:00 2001 From: "agent.develop" Date: Sat, 19 Sep 2026 21:07:06 +0800 Subject: [PATCH] =?UTF-8?q?deliver:=20=E4=BA=A4=E4=BB=98=E6=94=B6=E5=8F=A3?= =?UTF-8?q?=EF=BC=88=E5=BC=95=E6=93=8E=E4=BB=A3=E4=B8=BA=E6=8F=90=E4=BA=A4?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/pbl_runtime_event_ddl.py | 716 +++++++++++++++++++++++++++++++ 1 file changed, 716 insertions(+) create mode 100644 scripts/pbl_runtime_event_ddl.py diff --git a/scripts/pbl_runtime_event_ddl.py b/scripts/pbl_runtime_event_ddl.py new file mode 100644 index 0000000..b7472a4 --- /dev/null +++ b/scripts/pbl_runtime_event_ddl.py @@ -0,0 +1,716 @@ +#!/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。""" + months = int(months or MONTHS_PRECREATE) + return """-- 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)s' + AND PARTITION_NAME = v_name; + IF v_exists = 0 THEN + SET @pbl_part_sql = CONCAT( + 'ALTER TABLE `%(table)s` REORGANIZE PARTITION `%(pmax)s` INTO (', + 'PARTITION `', v_name, '` VALUES LESS THAN (DATE ''', v_upper, '''), ', + 'PARTITION `%(pmax)s` 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)s' + AND PARTITION_NAME IS NOT NULL + AND PARTITION_NAME <> '%(pmax)s' + 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)s` 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)s_partition_maintain`; +CREATE EVENT IF NOT EXISTS `ev_%(table)s_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)s 个月分区' +DO CALL pbl_runtime_event_ensure_partitions(%(months)s); +""" % {"table": TABLE, "pmax": PARTITION_MAXVALUE, "months": months, + "today": (today or date.today()).isoformat()} + + +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 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)) + + cur.execute("SELECT PARTITION_NAME FROM information_schema.PARTITIONS " + "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=%s " + "AND PARTITION_NAME IS NOT NULL ORDER BY PARTITION_ORDINAL_POSITION", (TABLE,)) + names = [r[0] for r in cur.fetchall()] + results.append(("MAXVALUE 兜底", PARTITION_MAXVALUE in names, ",".join(names))) + + # 插入两个不同月份 → 落不同分区 + probe = "m11bprobe%s" % date.today().strftime("%H%M%S") + try: + cur.execute("INSERT INTO `%s` (tenant_id,event_uid,event_type,scene_id,seq,created_at)" + " VALUES ('t_probe','%s_a','probe',0,1,'2026-01-05 00:00:00')" % (TABLE, probe)) + cur.execute("INSERT INTO `%s` (tenant_id,event_uid,event_type,scene_id,seq,created_at)" + " VALUES ('t_probe','%s_b','probe',0,2,'2027-03-05 00:00:00')" % (TABLE, probe)) + except Exception as exc: # noqa: BLE001 + results.append(("跨月插入", False, str(exc)[:160])) + else: + cur.execute("SELECT PARTITION_NAME FROM (SELECT '%s_a' k" + " UNION ALL SELECT '%s_b') x" % (probe, probe)) + parts = set() + for k in ("%s_a" % probe, "%s_b" % probe): + cur.execute("SELECT PARTITION_NAME FROM information_schema.PARTITIONS" + " WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=%s" + " AND PARTITION_DESCRIPTION <> 'MAXVALUE'" + " AND PARTITION_DESCRIPTION IS NOT NULL", (TABLE,)) + # 精确定位:EXPLAIN 每行 + located = [] + for uid in ("%s_a" % probe, "%s_b" % probe): + cur.execute("EXPLAIN PARTITIONS SELECT id FROM `%s` WHERE event_uid=%%s" % TABLE, (uid,)) + row = cur.fetchone() + located.append(str(row[3] if row and len(row) > 3 else row)) + results.append(("跨月落不同分区", len(set(located)) == 2, " | ".join(located))) + + # EXPLAIN 分区裁剪 + cur.execute("EXPLAIN PARTITIONS SELECT id FROM `%s` WHERE created_at>='2026-01-01'" + " AND created_at<'2026-02-01'" % TABLE) + row = cur.fetchone() + pruned = str(row[3] if row and len(row) > 3 else row) + results.append(("分区裁剪", pruned.count(",") == 0 and "p" in pruned, "EXPLAIN partitions=%s" % pruned)) + + # 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())