pbl_runtime_ext/scripts/pbl_runtime_event_ddl.py
agent.develop 907603653d revert(M11b-1c): 还原越界生成的 scripts/pbl_runtime_event_ddl.py 至基线 4a7c7b9(blob 7b9c991)
M11b-1c 任务书边界=不改生成器、DDL 侧为权威派生源;309d374 与引擎代提交 7e7477f
一并改动了 scripts/pbl_runtime_event_ddl.py(_row_in_partition/_locate_row_partition/
verify_db 的 % 转义重构 + PBL_DDL_DB_HOST/PORT harness 覆盖),触发 QC #1/#2 越界硬门禁。
本提交仅还原该 .py,models/pbl_runtime_event.json 本体一字未动。
2026-09-19 21:54:07 +08:00

796 lines
37 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/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<YYYYMM> +
末尾 ``pmax VALUES LESS THAN (MAXVALUE)`` 兜底,预建 6 个月。
3. append-only 双保险:
(a) 权限层 —— 迁移脚本 REVOKE app 账号的 UPDATE/DELETEdeploy 按需执行);
(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() 会被
触发器拒绝 —— 因此本模块所有写路径必须走 INSERTtx_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.pyM11a 单事务写路径)使用的列
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 PARTITIONappend-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"
"-- 目标引擎 mariadbprojects/pbls/env/test.json: db.engine=mariadb\n"
"-- 执行mysql --protocol=tcp -u<dba> -p <dbname> < 本文件\n"
"-- 说明:触发器含多语句体,须用 mysql CLIDELIMITER执行\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 缺字段 %sevent_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 分区键用了 TIMESTAMPMySQL 不允许)")
else:
msgs.append("PASS 分区键类型 DATETIME非 TIMESTAMP")
# 5. append-onlyDDL/迁移全文不得出现对本表的 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-onlyUPDATE / 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())