1027 lines
48 KiB
Python
1027 lines
48 KiB
Python
#!/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/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<dba> -p <dbname> < 本文件\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())
|
||
|
||
|
||
_BLOCK_OPEN = re.compile(r"^BEGIN\b", re.I)
|
||
_BLOCK_CLOSE = re.compile(r"^END\s*;?\s*$", re.I)
|
||
_SUB_CLOSE = re.compile(r"^END\s+(IF|LOOP|CASE)\b", re.I)
|
||
_STMT_HEAD = re.compile(
|
||
r"^(DROP|CREATE|GRANT|REVOKE|SET|SELECT|CALL|ALTER|WITH|INSERT|UPDATE|DELETE|START)\b", re.I)
|
||
# 一个 $$ 段内不应再出现的「另一条 DDL 语句头」——出现即说明两条语句被黏在了一起
|
||
_INLINE_DDL_HEAD = re.compile(
|
||
r"^\s*(DROP|CREATE)\s+(PROCEDURE|TRIGGER|EVENT|TABLE|FUNCTION)\b", re.I)
|
||
# 裸关键字残缺行:`CREATE$$` / `DROP$$`(QC #12 的致命产物形态)
|
||
_BARE_TOKEN = re.compile(r"^(CREATE|DROP|ALTER|GRANT|REVOKE)\s*\$\$\s*$", re.I)
|
||
|
||
|
||
def _with_delimiter(body, delim="$$"):
|
||
"""把过程/触发器/EVENT 正文按语句切分,逐条以 $$ 收尾,交给 mysql CLI(DELIMITER)执行。
|
||
|
||
QC #12 修复(致命缺陷:迁移资产整体不可执行):
|
||
旧实现 `re.split(r"(?<=;)\n(?=(DROP|CREATE|...))", body)` 的前瞻里带**捕获组**,
|
||
re.split 会把捕获到的文本('CREATE' / 'DROP')作为独立元素插进结果列表;这些孤立 token
|
||
随后被当成一条完整语句,拼成裸 `CREATE$$` / `DROP$$` 行写进 .sql。在 `DELIMITER $$` 下
|
||
服务端把 `CREATE` 当独立语句解析 → ERROR 1064,02/03/canonical 三个文件全部无法执行。
|
||
另外旧实现把 `DROP PROCEDURE/TRIGGER/EVENT` 留在 `;` 结尾,而 DELIMITER 已改成 `$$`,
|
||
客户端只认 `$$` 为语句结束符,`;` 结尾的 DROP 会与下一条 CREATE 黏成一条语句。
|
||
现改为按行扫描 + BEGIN/END 嵌套深度判定语句边界(不再使用带捕获组的 re.split):
|
||
· 每条语句去掉结尾 `;` 后以 `$$` 结束(含 DROP ... IF EXISTS);
|
||
· `--` 注释行原样保留、不加 `$$`(否则 `$$` 被行注释吞掉,语句永不结束);
|
||
· 结构上只可能产出以完整关键字开头的语句段,不可能产出裸 `CREATE`/`DROP`。
|
||
"""
|
||
out = []
|
||
buf = []
|
||
cmts = []
|
||
depth = [0]
|
||
|
||
def flush():
|
||
if not buf:
|
||
if cmts:
|
||
out.append("\n".join(cmts))
|
||
del cmts[:]
|
||
return
|
||
stmt = re.sub(r";\s*$", "", "\n".join(buf).strip())
|
||
head = ("\n".join(cmts) + "\n") if cmts else ""
|
||
out.append(head + stmt + delim)
|
||
del buf[:]
|
||
del cmts[:]
|
||
|
||
for raw in body.replace("\r\n", "\n").split("\n"):
|
||
line = raw.rstrip()
|
||
s = line.strip()
|
||
if not s:
|
||
continue
|
||
if not buf:
|
||
if s.startswith("--") or s.startswith("#") or not _STMT_HEAD.match(s):
|
||
cmts.append(line)
|
||
continue
|
||
if _BLOCK_OPEN.match(s):
|
||
depth[0] += 1
|
||
elif _BLOCK_CLOSE.match(s) and depth[0] > 0:
|
||
depth[0] -= 1
|
||
buf.append(line)
|
||
if depth[0] == 0 and s.endswith(";") and not _SUB_CLOSE.match(s):
|
||
flush()
|
||
flush()
|
||
return "\n\n".join(out) + "\n"
|
||
|
||
|
||
def _delimiter_integrity_errors(text, label=""):
|
||
"""DELIMITER 语句完整性静态断言(QC #16:PASS 要能证明「可执行」,而非只查关键字存在)。
|
||
|
||
断言项:
|
||
1) `DELIMITER x` 与 `DELIMITER ;` 严格配对,文件结束时不残留未闭合块;
|
||
2) 每个 `$$` 段首个代码行必须是**完整**语句头(关键字 + 对象类型),
|
||
禁止 `CREATE$$` / `DROP$$` 这类孤立残缺 token(服务端 ERROR 1064);
|
||
3) 一个 `$$` 段内不得再出现第二条 DDL 语句头(防止 `;` 结尾语句与下一条黏连);
|
||
4) 未以 `$$` 收尾的残段里若含可执行语句(纯注释尾巴放行)即判失败。
|
||
返回缺陷描述列表;空列表 = 结构通过。
|
||
"""
|
||
errs = []
|
||
pfx = ("[%s] " % label) if label else ""
|
||
lines = text.split("\n")
|
||
delim = [None]
|
||
seg = [[]]
|
||
|
||
def code_of(ls):
|
||
return [x.strip() for x in ls if x.strip() and not x.strip().startswith("--")
|
||
and not x.strip().startswith("#")]
|
||
|
||
def close_segment():
|
||
code = code_of(seg[0])
|
||
if not code:
|
||
seg[0] = []
|
||
return
|
||
d = delim[0] or "$$"
|
||
first_words = code[0].replace(d, "").split()
|
||
if len(first_words) < 2:
|
||
errs.append(pfx + "残缺语句 token(服务端按独立语句解析 → ERROR 1064): %r" % code[0])
|
||
for x in code[1:]:
|
||
if _INLINE_DDL_HEAD.match(x):
|
||
errs.append(pfx + "一个 %s 段内黏连了多条 DDL 语句: %r" % (d, x.strip()))
|
||
seg[0] = []
|
||
|
||
for i, ln in enumerate(lines, 1):
|
||
s = ln.strip()
|
||
if _BARE_TOKEN.match(s):
|
||
errs.append(pfx + "L%d 裸关键字$$ 残缺行: %r" % (i, s))
|
||
if s.upper().startswith("DELIMITER"):
|
||
parts = s.split(None, 1)
|
||
new = parts[1].strip() if len(parts) > 1 else ";"
|
||
if seg[0]:
|
||
close_segment()
|
||
if delim[0] is None:
|
||
if new == ";":
|
||
errs.append(pfx + "L%d 出现没有配对的 `DELIMITER ;`" % i)
|
||
else:
|
||
delim[0] = new
|
||
else:
|
||
if new != ";":
|
||
errs.append(pfx + "L%d DELIMITER %s 未闭合就切换到 %s"
|
||
% (i, delim[0], new))
|
||
delim[0] = new
|
||
else:
|
||
delim[0] = None
|
||
continue
|
||
if delim[0]:
|
||
seg[0].append(ln)
|
||
if s.endswith(delim[0]):
|
||
close_segment()
|
||
if seg[0]:
|
||
errs.append(pfx + "DELIMITER 块末尾存在未以 %s 收尾的可执行残段: %r"
|
||
% (delim[0] or "$$", code_of(seg[0])[0][:60]))
|
||
seg[0] = []
|
||
if delim[0] is not None:
|
||
errs.append(pfx + "文件结束时 DELIMITER %s 仍未闭合" % delim[0])
|
||
return errs
|
||
|
||
|
||
# ---------------------------------------------------------------- 迁移文件产出
|
||
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 + 分区存在性判断")
|
||
|
||
# 8. DELIMITER 语句完整性(QC #16:PASS 必须能证明「可执行」,不能只查关键字存在性)
|
||
artifacts = [
|
||
("01", create_table_sql()),
|
||
("02", MIGRATION_FILES[1][1]("pbls")),
|
||
("03", MIGRATION_FILES[2][1]("pbls")),
|
||
("canonical", full_script()),
|
||
]
|
||
struct_errs = []
|
||
for label, text in artifacts:
|
||
struct_errs.extend(_delimiter_integrity_errors(text, label))
|
||
if struct_errs:
|
||
ok = False
|
||
for e in struct_errs:
|
||
msgs.append("FAIL DELIMITER 语句完整性: " + e)
|
||
else:
|
||
msgs.append("PASS DELIMITER 语句完整性: 4 个产物 DELIMITER 配对正确、%d 个 $$ 段首关键字完整、"
|
||
"无裸 CREATE$$/DROP$$ 残缺行"
|
||
% sum(len(re.findall(r"\$\$\s*$", t, re.M)) for _, t in artifacts))
|
||
|
||
# 9. 幂等/结构关键字逐产物覆盖(QC #17:完整陈述,不再跳号回避缺陷)
|
||
need_map = [
|
||
("01", ("CREATE TABLE IF NOT EXISTS", "PARTITION BY RANGE COLUMNS",
|
||
"VALUES LESS THAN (MAXVALUE)")),
|
||
("02", ("DROP PROCEDURE IF EXISTS", "DROP TRIGGER IF EXISTS",
|
||
"BEFORE UPDATE", "BEFORE DELETE")),
|
||
("03", ("DROP PROCEDURE IF EXISTS", "DROP EVENT IF EXISTS", "CREATE EVENT")),
|
||
]
|
||
by_label = dict(artifacts)
|
||
miss_all = []
|
||
for label, needs in need_map:
|
||
miss_all.extend("%s:%s" % (label, n) for n in needs if n not in by_label[label])
|
||
if miss_all:
|
||
ok = False
|
||
msgs.append("FAIL 幂等/结构关键字缺失: " + ", ".join(miss_all))
|
||
else:
|
||
msgs.append("PASS 幂等关键字逐文件覆盖: 01 建表分区 / 02 过程+双触发器 / 03 双过程+EVENT")
|
||
return ok, msgs
|
||
|
||
|
||
# ---------------------------------------------------------------- 连库执行 / 验证
|
||
def _db_conf():
|
||
"""读项目唯一事实源 env/test.json(禁止硬编码连接串)。
|
||
|
||
沙箱库开关(harness,QC #1 配套的验证隔离能力):环境变量 PBL_DDL_DB_HOST /
|
||
PBL_DDL_DB_PORT **只覆盖 host/port**,供 CI / 验证 harness 把连库分支指向一次性沙箱库
|
||
——破坏性判定(DROP PARTITION、探针 INSERT、guard break-glass)不能拿共享 test/prod 库做实验。
|
||
· user / password / dbname 一律仍只来自 env/*.json:本函数不接受任何环境变量覆盖,
|
||
也不内置默认主机、默认端口、默认账号(禁止硬编码);
|
||
· 两个变量都不设置时,返回值与引入该开关之前逐字一致(零行为变化);
|
||
· 覆盖生效时 env_name 追加 "+harness" 后缀,日志与 --verify-db 输出可一眼看出连的是沙箱库。
|
||
"""
|
||
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"):
|
||
host = (os.environ.get("PBL_DDL_DB_HOST") or "").strip()
|
||
port = (os.environ.get("PBL_DDL_DB_PORT") or "").strip()
|
||
if host:
|
||
db = dict(db, host=host)
|
||
env_name = "%s+harness" % env_name
|
||
if port:
|
||
if not port.isdigit():
|
||
raise RuntimeError("PBL_DDL_DB_PORT 必须是数字端口号: %r" % port)
|
||
db = dict(db, port=int(port))
|
||
if "+harness" not in env_name:
|
||
env_name = "%s+harness" % env_name
|
||
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 _norm_part(name):
|
||
"""分区名统一口径(QC #1 附带缺陷:两处来源比对方式不一致)。
|
||
|
||
information_schema.PARTITIONS 的 PARTITION_NAME 与 EXPLAIN FORMAT=JSON 计划里的
|
||
partitions 元素,跨引擎/版本可能带反引号、双引号、空白或大小写差异。两边都过这个
|
||
归一化再比对,避免「分区确实存在但判定说不属于清单」的假失败。
|
||
"""
|
||
if not name:
|
||
return ""
|
||
text = name.decode("utf-8", "replace") if isinstance(name, (bytes, bytearray)) else str(name)
|
||
return text.strip().strip("`").strip('"').strip("'").lower()
|
||
|
||
|
||
def _render(sql, args=None, placeholders=None):
|
||
"""SQL 的 % 转义统一出口:可选渲染标识符 + fail-fast 自检(QC #1 的根因处置)。
|
||
|
||
约定:模板中由 Python 侧填入的标识符(表名/分区名)写 %s;要留给 DBAPI 绑参的
|
||
百分号占位符一律写 %%s,标识符渲染后即变回 %s。
|
||
· args 不为 None:先执行 sql % args 完成标识符渲染;
|
||
· args 为 None:视为调用方已完成渲染(模板就地 % 过),本函数只做转义自检。
|
||
自检两条,任一不满足立即抛 RuntimeError(而不是丢给驱动报
|
||
"unsupported format character",那会在 verify_db 里被 except 吞成一条模糊失败):
|
||
· 结果仍残留 '%%' → 转义写错;
|
||
· '%s' 个数与 placeholders 不符 → 漏写/多写绑参占位符。
|
||
渲染单独成行、括号显式定界,不再依赖「相邻字面量先拼接、再整体 %」的隐式优先级。
|
||
"""
|
||
if args is not None:
|
||
sql = sql % tuple(args)
|
||
if "%%" in sql:
|
||
raise RuntimeError("SQL 渲染后仍残留 '%%',转义写法有误: " + sql[:200])
|
||
if placeholders is not None and sql.count("%s") != placeholders:
|
||
raise RuntimeError("SQL 渲染后 DBAPI 占位符应为 %d 个,实际 %d 个: %s"
|
||
% (placeholders, sql.count("%s"), sql[:200]))
|
||
return sql
|
||
|
||
|
||
def _explain_partitions(cur, sql, params=None):
|
||
"""EXPLAIN FORMAT=JSON 取执行计划实际扫描的分区清单(QC #6)。
|
||
|
||
不再用旧的 PARTITIONS 扩展语法(EXPLAIN 后跟 PARTITIONS 关键字那种):那种写法把分区列固定在
|
||
结果集第 4 列(下标 3),列序随 MariaDB/MySQL 版本变化,且该扩展语法在新版 MySQL 已废弃。
|
||
JSON 计划里的 partitions 字段两版一致。
|
||
取不到(引擎不支持/解析失败)返回 [],由调用方判失败,不猜。
|
||
"""
|
||
# QC #1(转义):传进来的 sql 必须已完成标识符渲染,只剩 DBAPI 占位符 %s;残留 '%%' 说明
|
||
# 调用方转义写错,直接抛出让它响,别丢给驱动报 "unsupported format character"。
|
||
if "%%" in sql:
|
||
raise RuntimeError("EXPLAIN SQL 仍残留 '%%',转义写法有误: " + sql[:200])
|
||
args = tuple(params) if params else None
|
||
if args is not None and sql.count("%s") != len(args):
|
||
raise RuntimeError("EXPLAIN SQL 占位符 %d 个与参数 %d 个不符: %s"
|
||
% (sql.count("%s"), len(args), sql[:200]))
|
||
try:
|
||
cur.execute("EXPLAIN FORMAT=JSON " + sql, args)
|
||
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:
|
||
name = _norm_part(p)
|
||
if name and name not in found:
|
||
found.append(name)
|
||
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) 反查:探针行是否物理落在该分区。"""
|
||
sql = _render("SELECT id FROM `%s` PARTITION (`%s`)"
|
||
" WHERE tenant_id='t_probe' AND event_uid=%%s" % (TABLE, _norm_part(part)),
|
||
placeholders=1)
|
||
try:
|
||
cur.execute(sql, (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 个分区)。
|
||
"""
|
||
# 1) 逐个 PARTITION(p) 物理反查(复用 _row_in_partition,SQL 只在其内拼装一次,
|
||
# 避免同一份转义逻辑在两处各写一遍、改一处漏一处)
|
||
for part in part_names:
|
||
if _row_in_partition(cur, part, uid):
|
||
return _norm_part(part)
|
||
# 2) 反查不中才退回 EXPLAIN FORMAT=JSON;标识符渲染与转义自检统一走 _render
|
||
probe_sql = _render("SELECT id FROM `%s` WHERE tenant_id='t_probe'"
|
||
" AND event_uid=%%s AND created_at=%%s" % TABLE,
|
||
placeholders=2)
|
||
planned = _explain_partitions(cur, probe_sql, (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)))
|
||
# 分区名比对统一走归一化口径(information_schema 与 EXPLAIN JSON 计划两来源一致)
|
||
part_norm = set(_norm_part(n) for n in 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"))
|
||
ins_sql = _render("INSERT INTO `%s` (tenant_id,event_uid,event_type,scene_id,seq,created_at)"
|
||
" VALUES ('t_probe',%%s,'probe',0,%%s,%%s)" % TABLE,
|
||
placeholders=3)
|
||
try:
|
||
for i, (uid, ts) in enumerate(probes):
|
||
cur.execute(ins_sql, (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(_norm_part(x) in part_norm
|
||
for p in located for x in str(p).split(","))
|
||
results.append(("跨月落不同分区", distinct and in_list,
|
||
" | ".join("%s -> %s" % (uid, p)
|
||
for (uid, _), p in zip(probes, located))))
|
||
|
||
# 分区裁剪:created_at 区间只应命中该区间所属的少数分区,而非全部分区。
|
||
# 旧写法(PARTITIONS 扩展语法 + 取结果集下标 3 那一列):列序跨版本不稳定,新版 MySQL 已废弃。
|
||
prune_sql = _render("SELECT id FROM `%s` WHERE created_at>='2026-01-01'"
|
||
" AND created_at<'2026-02-01'" % TABLE,
|
||
placeholders=0)
|
||
pruned = _explain_partitions(cur, prune_sql)
|
||
results.append(("分区裁剪",
|
||
bool(pruned) and all(p in part_norm 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())
|