219 lines
8.7 KiB
Python
219 lines
8.7 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""M11b-2 运行时单事务写入 —— SQL 方言层(纯 SQL 文本 + 占位符,不绑定任何 ORM/驱动)。
|
||
|
||
职责边界:本模块只负责「把 SQL 写对」——方言识别、占位符风格、READ COMMITTED 设置语句、
|
||
悲观锁子句、以及两张表的建表 DDL。真正的事务编排在 :mod:`world_sync.pbl_runtime_tx`。
|
||
|
||
**列集单一真源(M11b-2b-B 起,QC #3 整改)**
|
||
``EVENT_COLUMNS`` / ``STATE_COLUMNS`` 不再手写常量,而是从
|
||
``modules/world_sync/models/pbl_runtime_event.json`` / ``pbl_entity_state.json``
|
||
(spec 四段式内联副本,结构权威源 = pbl_runtime_ext,PM 裁决 B)的 ``fields[]``
|
||
**机械导出**。旧版在本文件里手抄的 ``event_id`` / ``world_id``(state 表) / ``state``
|
||
/ ``updated_by_event`` 等列名在真源中不存在,构成「代码一套、JSON 另一套」的双写分叉,
|
||
是 QC #3 判定的消费方契约断裂根因;现在这类漂移在结构上不可能再发生。
|
||
|
||
设计要点:
|
||
1. 用 DB-API 2.0 规范接口(``paramstyle`` / ``execute`` / ``rowcount`` / ``commit`` /
|
||
``rollback``),因此同一套编排代码既能跑 sqlite(单测、本地)也能跑 MySQL(生产);
|
||
2. 占位符不硬编码:按 ``driver.paramstyle`` 生成(``qmark`` → ``?``,``format`` → ``%s``,
|
||
``pyformat`` → ``%(name)s``),避免在 MySQL 上写出 sqlite 的 ``?``;
|
||
3. ``SELECT ... FOR UPDATE`` 只有 MySQL/PG 支持,sqlite 侧返回空串(sqlite 写锁由事务本身保证),
|
||
并发正确性在两条路径上都由「乐观锁 state_version」兜底;
|
||
4. 不引入任何广播 / 轮询 / 时延优化逻辑(明确排除在 M11b-2 范围之外);
|
||
5. 建表 DDL 由 :func:`world_sync.pbl_runtime_authority.ddl` 按同一份模型渲染,
|
||
本模块不再自带第二套列定义(生产库的分区/触发器 DDL 归 pbl_runtime_ext 所有)。
|
||
"""
|
||
|
||
import sqlite3
|
||
|
||
from . import pbl_runtime_authority as auth
|
||
|
||
__all__ = [
|
||
"EVENT_TABLE",
|
||
"STATE_TABLE",
|
||
"DIALECT_SQLITE",
|
||
"DIALECT_MYSQL",
|
||
"DIALECT_POSTGRES",
|
||
"detect_dialect",
|
||
"placeholder",
|
||
"render_placeholders",
|
||
"begin_stmt",
|
||
"read_committed_stmt",
|
||
"for_update_clause",
|
||
"supports_transaction_control",
|
||
"DDL_STATEMENTS",
|
||
"ddl_for",
|
||
"EVENT_COLUMNS",
|
||
"STATE_COLUMNS",
|
||
"columns_of",
|
||
"percentile",
|
||
]
|
||
|
||
EVENT_TABLE = "pbl_runtime_event"
|
||
STATE_TABLE = "pbl_entity_state"
|
||
|
||
DIALECT_SQLITE = "sqlite"
|
||
DIALECT_MYSQL = "mysql"
|
||
DIALECT_POSTGRES = "postgres"
|
||
|
||
|
||
def columns_of(table):
|
||
"""取某表在模型 JSON(spec 四段式内联副本)里声明的全列集(有序)。
|
||
|
||
fail-closed:模型缺失 / 不是四段式 / fields 为空,一律由
|
||
:mod:`world_sync.pbl_runtime_authority` 抛
|
||
:class:`~world_sync.pbl_runtime_authority.AuthorityResolutionError`,不退回旧常量。
|
||
"""
|
||
return auth.columns(table)
|
||
|
||
|
||
# append-only 事件流列集(= models/pbl_runtime_event.json fields[] 顺序,25 列,含应用层
|
||
# 生成的 id;id 生成方见 world_sync/pbl_runtime_tx.py::new_event_uid 与该文件 summary[0].comment)
|
||
EVENT_COLUMNS = columns_of(EVENT_TABLE)
|
||
|
||
# 实体状态列集(= models/pbl_entity_state.json fields[] 顺序,10 列;定位键是
|
||
# tenant_id + session_id + entity_id,对应唯一键 uk_es,真源无 world_id / state 列)
|
||
STATE_COLUMNS = columns_of(STATE_TABLE)
|
||
|
||
|
||
def detect_dialect(conn_or_module):
|
||
"""从连接/驱动对象推断方言。
|
||
|
||
识别顺序:显式 ``pbl_dialect`` 属性 → DB-API ``paramstyle``+模块名 → sqlite3 连接。
|
||
识别不了抛 :class:`ValueError`(由上层包成 TxConfigError),不猜、不默认 MySQL。
|
||
"""
|
||
explicit = getattr(conn_or_module, "pbl_dialect", None)
|
||
if explicit:
|
||
return str(explicit).lower()
|
||
|
||
if isinstance(conn_or_module, sqlite3.Connection):
|
||
return DIALECT_SQLITE
|
||
|
||
# 连接对象退化到它的模块(sqlite3.Connection 已在上面命中,这里主要给 pymysql/MySQLdb)
|
||
owner = getattr(conn_or_module, "__class__", None)
|
||
mod_name = ""
|
||
cls = getattr(conn_or_module, "__class__", None)
|
||
if cls is not None:
|
||
mod_name = (getattr(cls, "__module__", "") or "").lower()
|
||
if not mod_name:
|
||
mod_name = (getattr(conn_or_module, "__name__", "") or "").lower()
|
||
|
||
if "sqlite" in mod_name:
|
||
return DIALECT_SQLITE
|
||
if "mysql" in mod_name or "pymysql" in mod_name or "mariadb" in mod_name:
|
||
return DIALECT_MYSQL
|
||
if "psycopg" in mod_name or "postgres" in mod_name:
|
||
return DIALECT_POSTGRES
|
||
|
||
# 驱动名推不出来时按 paramstyle 兜底(qmark 是 sqlite/PG 二选一的唯一硬线索,
|
||
# 这里只兜底 sqlite,因为 PG 驱动一律 pyformat;仍未知则报错而不是猜)
|
||
paramstyle = (getattr(conn_or_module, "paramstyle", "") or "").lower()
|
||
if paramstyle == "qmark" and hasattr(conn_or_module, "execute"):
|
||
return DIALECT_SQLITE
|
||
raise ValueError("unsupported db dialect: %r" % (mod_name or conn_or_module,))
|
||
|
||
|
||
def placeholder(dialect, index=0, name=None):
|
||
"""按方言返回一个参数占位符。``index`` 用于位置参数,``name`` 用于命名参数。"""
|
||
if dialect == DIALECT_SQLITE:
|
||
return "?"
|
||
if dialect == DIALECT_MYSQL:
|
||
return "%s"
|
||
if dialect == DIALECT_POSTGRES:
|
||
if name:
|
||
return "%%(%s)s" % name
|
||
return "$%d" % (index + 1)
|
||
raise ValueError("unsupported db dialect: %r" % dialect)
|
||
|
||
|
||
def render_placeholders(dialect, count):
|
||
"""生成 ``count`` 个占位符,逗号连接,直接嵌进 INSERT 的 VALUES 子句。"""
|
||
if count < 0:
|
||
raise ValueError("count must be >= 0")
|
||
return ", ".join(placeholder(dialect, i) for i in range(count))
|
||
|
||
|
||
def begin_stmt(dialect):
|
||
"""显式开启事务的语句。
|
||
|
||
- sqlite: ``BEGIN IMMEDIATE`` —— 立刻拿写锁,避免两事务在升级锁时死锁/脏读;
|
||
- MySQL/PG: 返回 ``None``,由驱动 autocommit=off 语义或上层 ``SET autocommit=0`` 控制,
|
||
多发一个 BEGIN 在部分驱动上会告警,因此不发。
|
||
"""
|
||
if dialect == DIALECT_SQLITE:
|
||
return "BEGIN IMMEDIATE"
|
||
return None
|
||
|
||
|
||
def read_committed_stmt(dialect):
|
||
"""会话级 READ COMMITTED 隔离级别设置(需求点 2)。sqlite 无此概念,返回 ``None``。"""
|
||
if dialect == DIALECT_MYSQL:
|
||
return "SET TRANSACTION ISOLATION LEVEL READ COMMITTED"
|
||
if dialect == DIALECT_POSTGRES:
|
||
return "SET SESSION CHARACTERISTICS AS TRANSACTION ISOLATION LEVEL READ COMMITTED"
|
||
return None
|
||
|
||
|
||
def for_update_clause(dialect):
|
||
"""悲观行锁子句;sqlite 返回空串(单写者模型,由 IMMEDIATE 事务保证)。"""
|
||
if dialect in (DIALECT_MYSQL, DIALECT_POSTGRES):
|
||
return " FOR UPDATE"
|
||
return ""
|
||
|
||
|
||
def supports_transaction_control(dialect):
|
||
"""是否需要本模块显式 commit/rollback(sqlite3 默认 isolation_level 会自己插 BEGIN,
|
||
但我们在建连时关掉它,因此两侧统一由本模块控制)。"""
|
||
return True
|
||
|
||
|
||
def _ddl_sqlite():
|
||
"""sqlite 建表语句列表(单测 / 本地环境用),由模型 JSON 渲染,无第二套列定义。"""
|
||
return [auth.ddl(EVENT_TABLE, DIALECT_SQLITE), auth.ddl(STATE_TABLE, DIALECT_SQLITE)]
|
||
|
||
|
||
def _ddl_mysql():
|
||
"""MySQL 建表语句列表(本地/预发用,InnoDB + utf8mb4),同样由模型 JSON 渲染。
|
||
|
||
生产库真正的 ``pbl_runtime_event`` DDL(按月 RANGE COLUMNS 分区、append-only 存储过程
|
||
与触发器)归 pbl_runtime_ext 的 ``scripts/pbl_runtime_event_ddl.py`` 与
|
||
``scripts/migrations/*.sql`` 所有,本函数不复制那套物理细节。
|
||
"""
|
||
return [auth.ddl(EVENT_TABLE, DIALECT_MYSQL), auth.ddl(STATE_TABLE, DIALECT_MYSQL)]
|
||
|
||
|
||
#: 各方言建表语句集合,供宿主初始化 / 单测建库使用。
|
||
DDL_STATEMENTS = {
|
||
DIALECT_SQLITE: _ddl_sqlite(),
|
||
DIALECT_MYSQL: _ddl_mysql(),
|
||
}
|
||
|
||
|
||
def ddl_for(dialect):
|
||
"""取某方言的建表语句列表;方言不认识直接报错,不静默返回空。"""
|
||
try:
|
||
return list(DDL_STATEMENTS[dialect])
|
||
except KeyError:
|
||
raise ValueError("no DDL for dialect: %r" % (dialect,))
|
||
|
||
|
||
def percentile(values, pct):
|
||
"""线性插值分位数(与 numpy ``percentile(..., method='linear')`` 同口径)。
|
||
|
||
用于 P99 时延断言:``percentile(latencies_ms, 99) <= 200``。
|
||
``values`` 为空返回 ``0.0``(调用方需自行判断样本量是否足够)。
|
||
"""
|
||
xs = sorted(float(v) for v in values)
|
||
n = len(xs)
|
||
if n == 0:
|
||
return 0.0
|
||
if n == 1:
|
||
return xs[0]
|
||
rank = (float(pct) / 100.0) * (n - 1)
|
||
lo = int(rank // 1)
|
||
hi = lo + 1
|
||
if hi >= n:
|
||
return xs[-1]
|
||
frac = rank - lo
|
||
return xs[lo] + (xs[hi] - xs[lo]) * frac
|