deliver: 交付收口(引擎代为提交)
This commit is contained in:
parent
2ee0d74dad
commit
3f0ee54c1c
@ -1,6 +1,6 @@
|
|||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
"""validate_models_json —— database-table-definition-spec 机械校验器(M11b-2b-A 修复版)。
|
"""validate_models_json —— database-table-definition-spec 机械校验器(M11b-2b-A 修复版 rev3)。
|
||||||
|
|
||||||
QC #4/#7/#9/#10 要求:表定义 JSON 必须是 spec 规定的四段式(summary/fields/indexes/
|
QC #4/#7/#9/#10 要求:表定义 JSON 必须是 spec 规定的四段式(summary/fields/indexes/
|
||||||
codes),类型必须是 spec 声明的抽象类型,主键 id 必须是 str/length>=32,索引必须用
|
codes),类型必须是 spec 声明的抽象类型,主键 id 必须是 str/length>=32,索引必须用
|
||||||
@ -17,12 +17,25 @@ codes),类型必须是 spec 声明的抽象类型,主键 id 必须是 str/
|
|||||||
QC #5 引用式 key/schema:spec 不承认引用式结构,直接判 ERROR 并提示改为内联四段式;
|
QC #5 引用式 key/schema:spec 不承认引用式结构,直接判 ERROR 并提示改为内联四段式;
|
||||||
COMPLIANCE_DEVIATIONS 中针对 pbl_runtime_event/pbl_entity_state 的 4 条失效登记项
|
COMPLIANCE_DEVIATIONS 中针对 pbl_runtime_event/pbl_entity_state 的 4 条失效登记项
|
||||||
一并删除(引用式在根键检查处即 FAIL,根本到不了 fields 校验,登记形同虚设)。
|
一并删除(引用式在根键检查处即 FAIL,根本到不了 fields 校验,登记形同虚设)。
|
||||||
|
QC #1(本轮)放行机制对 auto_increment 生效:字段段「spec 未定义键」检查原先无条件报错,
|
||||||
|
导致 (表名,"id","auto-increment") 即使已登记、即使传了 --allow-registered-deviation,
|
||||||
|
仍残留 `fields[N](id) 含 spec 未定义键: auto_increment` 而整体 FAIL —— id-not-str32
|
||||||
|
可豁免、auto_increment 不可豁免,机制自相矛盾。现在该检查改为**偏离感知**:
|
||||||
|
id 上的 auto_increment 键统一交由 id 分支按 strict/allow 两态处理(已登记 + 传开关
|
||||||
|
→ 降级为 DEVIATION(registered) 明细,不计入 errors;未登记或未传开关 → 照旧 ERROR);
|
||||||
|
其它字段上的未定义键也支持通用登记键 (表名, 列名, "unknown-key:<键名>"),同样两态。
|
||||||
|
|
||||||
已登记偏离(唯一允许的偏离,需 --allow-registered-deviation 显式声明才放行):
|
偏离放行语义(唯一允许的偏离,需 --allow-registered-deviation 显式声明才放行):
|
||||||
主键 id 非 str(32) 或 id 带 auto_increment,且已登记在 COMPLIANCE_DEVIATIONS 中
|
1) (表名, "id", "id-not-str32") —— 主键 id 非 str 或 length<32;
|
||||||
(典型场景:分区表物理主键需要 long + auto_increment)。未传该开关时这类项仍判 ERROR,
|
2) (表名, "id", "auto-increment") —— id 带 spec 未定义键 auto_increment;
|
||||||
并打印「需 --allow-registered-deviation 放行」提示;传入时逐条打印 DEVIATION(registered)
|
3) (表名, 列名, "unknown-key:<K>") —— 任意字段带 spec 未定义键 <K>(通用形式)。
|
||||||
明细,不静默放过。当前登记表为空(引用式失效项已删除),后续若真源确需偏离再登记。
|
未传该开关时上述项仍判 ERROR,并打印「需 --allow-registered-deviation 放行」提示;传入时
|
||||||
|
逐条打印 DEVIATION(registered) 明细(含登记说明),不静默放过。
|
||||||
|
注意:只有**登记过**的才可能放行;登记表为空时 --allow-registered-deviation 不放行任何项。
|
||||||
|
典型场景:分区表物理主键需要 long + auto_increment(MySQL 1503:分区键必须进主键)。
|
||||||
|
若某表确需此类写法,须先在其模块任务里把 (表名,"id","id-not-str32") 与
|
||||||
|
(表名,"id","auto-increment") 登记进 COMPLIANCE_DEVIATIONS,再配合开关使用 —— 两者缺一不可,
|
||||||
|
只登记不传开关仍 FAIL,只传开关不登记也 FAIL(本文件末尾自测 5 给出四态实测输出)。
|
||||||
|
|
||||||
用法:
|
用法:
|
||||||
python validate_models_json.py <models目录1> [<目录2> ...]
|
python validate_models_json.py <models目录1> [<目录2> ...]
|
||||||
@ -81,6 +94,20 @@ def _is_reference_style(data):
|
|||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def _unknown_key_deviation_key(table, field_name, key_name):
|
||||||
|
"""字段段 spec 未定义键对应的登记键(QC #1)。
|
||||||
|
|
||||||
|
id 上的 auto_increment 复用 (表名,"id","auto-increment"),与 id 分支同一登记项,
|
||||||
|
避免「登记了 auto-increment 却因 unknown-key 分支无条件报错」的放行失效缺陷;
|
||||||
|
其余未定义键使用通用形式 (表名, 列名, "unknown-key:<键名>")。
|
||||||
|
"""
|
||||||
|
if table is None:
|
||||||
|
return None
|
||||||
|
if field_name == "id" and key_name == "auto_increment":
|
||||||
|
return (table, "id", "auto-increment")
|
||||||
|
return (table, field_name, "unknown-key:%s" % key_name)
|
||||||
|
|
||||||
|
|
||||||
def check_type_value(type_value):
|
def check_type_value(type_value):
|
||||||
"""返回该 type 写法的问题列表(spec 只允许抽象类型)。
|
"""返回该 type 写法的问题列表(spec 只允许抽象类型)。
|
||||||
|
|
||||||
@ -181,9 +208,26 @@ def validate_model(path, strict_id=True):
|
|||||||
names.append(name)
|
names.append(name)
|
||||||
if not f.get("title"):
|
if not f.get("title"):
|
||||||
errors.append("%s.title 必填(DDL COMMENT 用)" % where)
|
errors.append("%s.title 必填(DDL COMMENT 用)" % where)
|
||||||
unknown = [k for k in f if k not in FIELD_KEYS]
|
# QC #1:spec 未定义键检查改为偏离感知 —— 已登记且传入 --allow-registered-deviation
|
||||||
|
# 的键不再计入 errors(id.auto_increment 交由下方 id 分支统一打印 DEVIATION 明细,
|
||||||
|
# 避免同一偏离重复计数;其它键在此直接降级为 DEVIATION(registered))。
|
||||||
|
unknown = []
|
||||||
|
for k in sorted([k for k in f if k not in FIELD_KEYS]):
|
||||||
|
dkey = _unknown_key_deviation_key(table, name, k)
|
||||||
|
if k == "auto_increment" and name == "id" and dkey in COMPLIANCE_DEVIATIONS:
|
||||||
|
continue # strict/allow 两态由 id 分支处理(同一登记项,只报一次)
|
||||||
|
if dkey is not None and dkey in COMPLIANCE_DEVIATIONS:
|
||||||
|
if strict_id:
|
||||||
|
# 已登记但未传开关:照旧 ERROR,并明确提示如何放行(与 id 分支同一语义)
|
||||||
|
errors.append("%s 含 spec 未定义键: %s —— 该项已登记偏离,%s"
|
||||||
|
% (where, k, DEVIATION_HINT))
|
||||||
|
continue
|
||||||
|
deviations.append("%s: 含 spec 未定义键 %s —— 已登记偏离:%s"
|
||||||
|
% (where, k, COMPLIANCE_DEVIATIONS[dkey]))
|
||||||
|
continue
|
||||||
|
unknown.append(k)
|
||||||
if unknown:
|
if unknown:
|
||||||
errors.append("%s 含 spec 未定义键: %s" % (where, ", ".join(sorted(unknown))))
|
errors.append("%s 含 spec 未定义键: %s" % (where, ", ".join(unknown)))
|
||||||
t = f.get("type")
|
t = f.get("type")
|
||||||
errors.extend("%s: %s" % (where, m) for m in check_type_value(t))
|
errors.extend("%s: %s" % (where, m) for m in check_type_value(t))
|
||||||
req = ABSTRACT_TYPES.get(str(t).strip().lower(), ())
|
req = ABSTRACT_TYPES.get(str(t).strip().lower(), ())
|
||||||
@ -208,14 +252,15 @@ def validate_model(path, strict_id=True):
|
|||||||
else:
|
else:
|
||||||
errors.append("%s: 主键 id 必须 str 且 length>=32(spec 规则),实测 type=%r length=%r"
|
errors.append("%s: 主键 id 必须 str 且 length>=32(spec 规则),实测 type=%r length=%r"
|
||||||
% (where, t, f.get("length")))
|
% (where, t, f.get("length")))
|
||||||
if f.get("auto_increment"):
|
if "auto_increment" in f:
|
||||||
key = (table, "id", "auto-increment")
|
key = (table, "id", "auto-increment")
|
||||||
if key in COMPLIANCE_DEVIATIONS:
|
if key in COMPLIANCE_DEVIATIONS:
|
||||||
if strict_id:
|
if strict_id:
|
||||||
errors.append("%s: id auto_increment —— 该项已登记偏离,%s"
|
errors.append("%s: id 带 spec 未定义键 auto_increment(实测 %r)"
|
||||||
% (where, DEVIATION_HINT))
|
" —— 该项已登记偏离,%s"
|
||||||
|
% (where, f.get("auto_increment"), DEVIATION_HINT))
|
||||||
else:
|
else:
|
||||||
deviations.append("%s: id auto_increment —— 已登记偏离:%s"
|
deviations.append("%s: id 带 spec 未定义键 auto_increment —— 已登记偏离:%s"
|
||||||
% (where, COMPLIANCE_DEVIATIONS[key]))
|
% (where, COMPLIANCE_DEVIATIONS[key]))
|
||||||
else:
|
else:
|
||||||
errors.append("%s: auto_increment 未在偏离登记表中(spec 无此键)" % where)
|
errors.append("%s: auto_increment 未在偏离登记表中(spec 无此键)" % where)
|
||||||
|
|||||||
@ -1,15 +1,27 @@
|
|||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
"""world_sync —— 世界运行时同步模块(PBL scense 侧)。
|
"""world_sync —— 世界运行时同步模块(PBL scense 侧)包入口。
|
||||||
|
|
||||||
M11b-2 增量:单事务内「事件 + 实体状态」原子落库。
|
**单一真源原则(响应 QC #2 / QC #3)**
|
||||||
对外只暴露写入/读取两个函数与异常族;**不含广播、不含轮询兜底、不含时延优化**
|
=====================================
|
||||||
|
本文件只做「转发 + 导出」,**不再自带第二份 ``load_world_sync`` 实现**。
|
||||||
|
唯一实现位于 :mod:`world_sync.init`(宿主挂载层),它同时完成:
|
||||||
|
|
||||||
|
- ① 实现:``init.py`` 内的 ``load_world_sync(server_env=None, conn_factory=None)``;
|
||||||
|
- ② 导出:本文件 ``from .init import load_world_sync, ...``(宿主 ``from world_sync import load_world_sync``
|
||||||
|
取到的就是 ①,不存在第二份);
|
||||||
|
- ③ env 注册:``init.py`` 内 ``env.xxx = xxx`` 的三处同步登记。
|
||||||
|
|
||||||
|
M11b-2 能力:单事务内「事件 + 实体状态」原子落库。
|
||||||
|
对外只暴露写入/读取两个契约函数与异常族;**不含广播、不含轮询兜底、不含时延优化**
|
||||||
(这些明确在 M11b-2 范围之外,由后续里程碑承担)。
|
(这些明确在 M11b-2 范围之外,由后续里程碑承担)。
|
||||||
|
|
||||||
宿主接线(遵循 project-directory-spec 8.1:模块不读部署环境凭据文件)::
|
宿主接线(遵循 project-directory-spec 8.1:模块不读部署环境凭据文件)::
|
||||||
|
|
||||||
# apps/scense/app/scense.py 的 init() 里
|
# apps/scense/app/scense.py 的 init() 里
|
||||||
from world_sync import load_world_sync
|
from world_sync import load_world_sync
|
||||||
load_world_sync(server_env) # 登记 pbl_runtime_conn_factory 等
|
load_world_sync() # 仅登记 CRUD + 表注册 + 契约
|
||||||
|
# 或宿主显式注入连接工厂(推荐,模块内部因此完全不需要知道凭据来源):
|
||||||
|
load_world_sync(server_env=env, conn_factory=my_factory)
|
||||||
|
|
||||||
# 业务侧写入
|
# 业务侧写入
|
||||||
from world_sync import write_event_with_state
|
from world_sync import write_event_with_state
|
||||||
@ -21,64 +33,152 @@ M11b-2 增量:单事务内「事件 + 实体状态」原子落库。
|
|||||||
|
|
||||||
``conn_factory`` 解析优先级(见 :mod:`world_sync.pbl_runtime_tx_env`):
|
``conn_factory`` 解析优先级(见 :mod:`world_sync.pbl_runtime_tx_env`):
|
||||||
显式注入 → ServerEnv 登记 → 环境变量 ``PBL_RUNTIME_DB_URL`` → 模块自有 ``conf/db.json``;
|
显式注入 → ServerEnv 登记 → 环境变量 ``PBL_RUNTIME_DB_URL`` → 模块自有 ``conf/db.json``;
|
||||||
四级全空则 fail-closed 抛 :class:`ConnFactoryNotConfigured`。
|
四级全空则 fail-closed 抛 :class:`ConnFactoryNotConfigured`(别名 :class:`ConnectionUnavailable`)。
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
# --- ② 导出:实现层(init.py)为唯一真源,本文件不重复定义任何函数 ---
|
||||||
|
from .init import (
|
||||||
|
# 挂载入口(唯一实现)
|
||||||
|
load_world_sync,
|
||||||
|
# 既有 world_sync CRUD
|
||||||
|
list_world_syncs,
|
||||||
|
get_world_sync,
|
||||||
|
create_world_sync,
|
||||||
|
update_world_sync,
|
||||||
|
delete_world_sync,
|
||||||
|
execute_world_sync,
|
||||||
|
# M11b-1 表注册(只读元数据)
|
||||||
|
MODULE_TABLES,
|
||||||
|
module_tables,
|
||||||
|
get_table_schema,
|
||||||
|
load_table_model,
|
||||||
|
pbl_runtime_event_schema,
|
||||||
|
pbl_runtime_event_partitions,
|
||||||
|
# M11b-2 运行时写入契约(配置 / 预演 / 诊断)
|
||||||
|
configure_runtime_writer,
|
||||||
|
get_runtime_writer,
|
||||||
|
build_runtime_write_plan,
|
||||||
|
runtime_write_percentile,
|
||||||
|
runtime_write_dialect_name,
|
||||||
|
runtime_write_ddl,
|
||||||
|
# 模块元信息
|
||||||
|
MODULE_NAME,
|
||||||
|
MODULE_VERSION,
|
||||||
|
SERVER_ENV_CONN_FACTORY_KEYS,
|
||||||
|
# 宿主基础包可用性(自检用,不建连)
|
||||||
|
HOST_DEPS_AVAILABLE,
|
||||||
|
HOST_DEPS_ERROR,
|
||||||
|
HostDepsNotAvailable,
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- M11b-2 契约函数与异常族(真源:pbl_runtime_tx* 系列,init.py 亦登记同名) ---
|
||||||
from .pbl_runtime_errors import (
|
from .pbl_runtime_errors import (
|
||||||
ConcurrentStateConflict,
|
ConcurrentStateConflict,
|
||||||
ConnFactoryNotConfigured,
|
ConnFactoryNotConfigured,
|
||||||
|
ConnectionUnavailable,
|
||||||
EventWriteError,
|
EventWriteError,
|
||||||
|
InvalidWriteRequest,
|
||||||
PblRuntimeError,
|
PblRuntimeError,
|
||||||
|
RuntimeWriteError,
|
||||||
StateWriteError,
|
StateWriteError,
|
||||||
|
TransactionAborted,
|
||||||
TxAbortError,
|
TxAbortError,
|
||||||
TxConfigError,
|
TxConfigError,
|
||||||
)
|
)
|
||||||
from .pbl_runtime_sql import (
|
from .pbl_runtime_sql import (
|
||||||
|
DIALECT_MYSQL,
|
||||||
|
DIALECT_SQLITE,
|
||||||
EVENT_TABLE,
|
EVENT_TABLE,
|
||||||
STATE_TABLE,
|
STATE_TABLE,
|
||||||
|
ddl_for,
|
||||||
detect_dialect,
|
detect_dialect,
|
||||||
percentile,
|
percentile,
|
||||||
)
|
)
|
||||||
from .pbl_runtime_tx import (
|
from .pbl_runtime_tx import (
|
||||||
|
DEFAULT_EVENT_TYPE,
|
||||||
|
DEFAULT_SOURCE,
|
||||||
|
EVENT_REQUIRED,
|
||||||
|
STATE_KEY_COLUMNS,
|
||||||
new_event_id,
|
new_event_id,
|
||||||
read_entity_state as _tx_read_entity_state,
|
new_event_uid,
|
||||||
utc_now_text,
|
utc_now_text,
|
||||||
write_event_with_state as _tx_write_event_with_state,
|
|
||||||
)
|
)
|
||||||
from .pbl_runtime_tx_env import (
|
from .pbl_runtime_tx_env import (
|
||||||
|
SERVER_ENV_KEYS,
|
||||||
|
build_mysql_conn_factory,
|
||||||
build_sqlite_conn_factory,
|
build_sqlite_conn_factory,
|
||||||
|
load_db_config,
|
||||||
|
module_conf_dir,
|
||||||
resolve_conn_factory,
|
resolve_conn_factory,
|
||||||
)
|
)
|
||||||
from .pbl_runtime_tx_sqlor import (
|
from .pbl_runtime_tx_sqlor import (
|
||||||
|
ENV_NAME_KEYS,
|
||||||
|
POOL_KEYS,
|
||||||
|
build_conn_factory_from_pools,
|
||||||
|
host_conn_factory,
|
||||||
read_entity_state,
|
read_entity_state,
|
||||||
write_event_with_state,
|
write_event_with_state,
|
||||||
)
|
)
|
||||||
|
|
||||||
MODULE_NAME = "world_sync"
|
#: 契约解析优先级键(与 pbl_runtime_tx_env.SERVER_ENV_KEYS 一致,单一来源)
|
||||||
MODULE_VERSION = "0.2.0-m11b2"
|
assert tuple(SERVER_ENV_CONN_FACTORY_KEYS) == tuple(SERVER_ENV_KEYS), \
|
||||||
|
"world_sync: SERVER_ENV_CONN_FACTORY_KEYS 与 pbl_runtime_tx_env.SERVER_ENV_KEYS 分叉"
|
||||||
#: 宿主 ServerEnv 上登记的 conn_factory 键(pbl_runtime_tx_env.SERVER_ENV_KEYS 与此一致)
|
|
||||||
SERVER_ENV_CONN_FACTORY_KEYS = ("pbl_runtime_conn_factory", "world_sync_conn_factory")
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
# 写入 / 读取(M11b-2 核心)
|
# 挂载入口(唯一实现位于 init.py)
|
||||||
"write_event_with_state",
|
|
||||||
"read_entity_state",
|
|
||||||
# 接线
|
|
||||||
"load_world_sync",
|
"load_world_sync",
|
||||||
"MODULE_NAME",
|
"MODULE_NAME",
|
||||||
"MODULE_VERSION",
|
"MODULE_VERSION",
|
||||||
|
"MODULE_TABLES",
|
||||||
"SERVER_ENV_CONN_FACTORY_KEYS",
|
"SERVER_ENV_CONN_FACTORY_KEYS",
|
||||||
|
"SERVER_ENV_KEYS",
|
||||||
|
# 既有 CRUD
|
||||||
|
"list_world_syncs",
|
||||||
|
"get_world_sync",
|
||||||
|
"create_world_sync",
|
||||||
|
"update_world_sync",
|
||||||
|
"delete_world_sync",
|
||||||
|
"execute_world_sync",
|
||||||
|
# M11b-1 表注册(只读)
|
||||||
|
"module_tables",
|
||||||
|
"get_table_schema",
|
||||||
|
"load_table_model",
|
||||||
|
"pbl_runtime_event_schema",
|
||||||
|
"pbl_runtime_event_partitions",
|
||||||
|
# M11b-2 写入 / 读取契约(核心)
|
||||||
|
"write_event_with_state",
|
||||||
|
"read_entity_state",
|
||||||
|
"configure_runtime_writer",
|
||||||
|
"get_runtime_writer",
|
||||||
|
"build_runtime_write_plan",
|
||||||
|
"runtime_write_percentile",
|
||||||
|
"runtime_write_dialect_name",
|
||||||
|
"runtime_write_ddl",
|
||||||
# 工厂与方言
|
# 工厂与方言
|
||||||
"resolve_conn_factory",
|
"resolve_conn_factory",
|
||||||
"build_sqlite_conn_factory",
|
"build_sqlite_conn_factory",
|
||||||
|
"build_mysql_conn_factory",
|
||||||
|
"build_conn_factory_from_pools",
|
||||||
|
"host_conn_factory",
|
||||||
|
"load_db_config",
|
||||||
|
"module_conf_dir",
|
||||||
"detect_dialect",
|
"detect_dialect",
|
||||||
|
"ddl_for",
|
||||||
"percentile",
|
"percentile",
|
||||||
"EVENT_TABLE",
|
"EVENT_TABLE",
|
||||||
"STATE_TABLE",
|
"STATE_TABLE",
|
||||||
|
"DIALECT_MYSQL",
|
||||||
|
"DIALECT_SQLITE",
|
||||||
|
"DEFAULT_EVENT_TYPE",
|
||||||
|
"DEFAULT_SOURCE",
|
||||||
|
"EVENT_REQUIRED",
|
||||||
|
"STATE_KEY_COLUMNS",
|
||||||
|
"POOL_KEYS",
|
||||||
|
"ENV_NAME_KEYS",
|
||||||
"new_event_id",
|
"new_event_id",
|
||||||
|
"new_event_uid",
|
||||||
"utc_now_text",
|
"utc_now_text",
|
||||||
# 异常族
|
# 异常族(含 M11b-2a 语义别名,同一对象)
|
||||||
"PblRuntimeError",
|
"PblRuntimeError",
|
||||||
"TxConfigError",
|
"TxConfigError",
|
||||||
"ConnFactoryNotConfigured",
|
"ConnFactoryNotConfigured",
|
||||||
@ -86,34 +186,12 @@ __all__ = [
|
|||||||
"StateWriteError",
|
"StateWriteError",
|
||||||
"ConcurrentStateConflict",
|
"ConcurrentStateConflict",
|
||||||
"TxAbortError",
|
"TxAbortError",
|
||||||
|
"RuntimeWriteError",
|
||||||
|
"InvalidWriteRequest",
|
||||||
|
"ConnectionUnavailable",
|
||||||
|
"TransactionAborted",
|
||||||
|
# 宿主基础包可用性诊断
|
||||||
|
"HOST_DEPS_AVAILABLE",
|
||||||
|
"HOST_DEPS_ERROR",
|
||||||
|
"HostDepsNotAvailable",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
def load_world_sync(server_env=None, conn_factory=None):
|
|
||||||
"""挂载 world_sync:向 ServerEnv 登记本模块能力与(可选)conn_factory。
|
|
||||||
|
|
||||||
幂等:重复调用只覆盖同名键,不产生副作用。返回本模块 ``__all__`` 便于宿主自检。
|
|
||||||
|
|
||||||
:param server_env: 宿主 ServerEnv 实例(None 时尝试取全局单例)
|
|
||||||
:param conn_factory: 宿主连接池工厂;提供则登记到 ``pbl_runtime_conn_factory``,
|
|
||||||
使模块内部无需知道任何部署环境凭据来源(fail-closed 而非自行找配置文件)。
|
|
||||||
"""
|
|
||||||
if server_env is None:
|
|
||||||
try:
|
|
||||||
from appPublic.serverEnv import ServerEnv
|
|
||||||
server_env = ServerEnv()
|
|
||||||
except Exception:
|
|
||||||
server_env = None
|
|
||||||
|
|
||||||
if server_env is not None:
|
|
||||||
if conn_factory is not None:
|
|
||||||
setattr(server_env, SERVER_ENV_CONN_FACTORY_KEYS[0], conn_factory)
|
|
||||||
setattr(server_env, "world_sync_module", MODULE_NAME)
|
|
||||||
setattr(server_env, "world_sync_version", MODULE_VERSION)
|
|
||||||
setter = getattr(server_env, "set", None)
|
|
||||||
if callable(setter):
|
|
||||||
try:
|
|
||||||
setter("world_sync_version", MODULE_VERSION)
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
return list(__all__)
|
|
||||||
|
|||||||
@ -23,10 +23,79 @@ M11b-2a 迁移修正(响应 QC #1):本文件此前导入的
|
|||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
|
|
||||||
from appPublic.uniqueID import getID
|
# 宿主基础包(appPublic / sqlor / ahserver):M11b-2a 起改为「容错导入 + 显式记录」——
|
||||||
from appPublic.timeUtils import curDateString, timestampstr
|
# 目的是让「三处接线」自检(import 本包 + 调 load_world_sync())不依赖宿主部署环境,
|
||||||
from sqlor.dbpools import DBPools
|
# 而不是吞错误:基础包缺失时记入 HOST_DEPS_ERROR,任何真正访问数据库的函数在调用点
|
||||||
from ahserver.serverenv import ServerEnv
|
# 通过 _require_host_deps() fail-fast 抛出带原因的 RuntimeError(不静默返回 None)。
|
||||||
|
HOST_DEPS_ERROR = None
|
||||||
|
try:
|
||||||
|
from appPublic.uniqueID import getID
|
||||||
|
from appPublic.timeUtils import curDateString, timestampstr
|
||||||
|
from sqlor.dbpools import DBPools
|
||||||
|
from ahserver.serverenv import ServerEnv
|
||||||
|
HOST_DEPS_AVAILABLE = True
|
||||||
|
except ImportError as _exc: # pragma: no cover - 仅在无宿主基础包的开发/自检环境触发
|
||||||
|
HOST_DEPS_AVAILABLE = False
|
||||||
|
HOST_DEPS_ERROR = _exc
|
||||||
|
getID = curDateString = timestampstr = None
|
||||||
|
DBPools = ServerEnv = None
|
||||||
|
|
||||||
|
|
||||||
|
class HostDepsNotAvailable(RuntimeError):
|
||||||
|
"""宿主基础包(appPublic/sqlor/ahserver)不可用时,访问数据库的入口 fail-fast。"""
|
||||||
|
|
||||||
|
def __init__(self, op):
|
||||||
|
super().__init__(
|
||||||
|
'%s requires host foundation packages (appPublic/sqlor/ahserver); '
|
||||||
|
'import error was: %r' % (op, HOST_DEPS_ERROR))
|
||||||
|
self.op = op
|
||||||
|
|
||||||
|
|
||||||
|
def _require_host_deps(op):
|
||||||
|
if not HOST_DEPS_AVAILABLE:
|
||||||
|
raise HostDepsNotAvailable(op)
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
class _StandaloneEnv(object):
|
||||||
|
"""无宿主 ServerEnv 时的登记载体(仅供接线自检,不建连、不接库)。"""
|
||||||
|
|
||||||
|
def __init__(self):
|
||||||
|
self._kv = {}
|
||||||
|
|
||||||
|
def __setattr__(self, k, v):
|
||||||
|
object.__setattr__(self, k, v)
|
||||||
|
|
||||||
|
def set(self, k, v):
|
||||||
|
self._kv[k] = v
|
||||||
|
|
||||||
|
def get(self, k, default=None):
|
||||||
|
return getattr(self, k, self._kv.get(k, default))
|
||||||
|
|
||||||
|
def get_module_dbname(self, m):
|
||||||
|
raise HostDepsNotAvailable('get_module_dbname(%s)' % m)
|
||||||
|
|
||||||
|
|
||||||
|
def _resolve_env(server_env=None):
|
||||||
|
"""解析登记目标 ServerEnv:显式注入优先,其次宿主全局单例,最后 standalone 载体。"""
|
||||||
|
if server_env is not None:
|
||||||
|
return server_env
|
||||||
|
if ServerEnv is not None:
|
||||||
|
try:
|
||||||
|
return ServerEnv()
|
||||||
|
except Exception as exc: # 宿主单例不可用:记录后降级,不中断接线登记
|
||||||
|
standalone_env_error(str(exc))
|
||||||
|
return _StandaloneEnv()
|
||||||
|
|
||||||
|
|
||||||
|
def standalone_env_error(message):
|
||||||
|
"""记录 ServerEnv 降级原因(幂等,不抛错)。"""
|
||||||
|
global STANDALONE_ENV_NOTICES
|
||||||
|
STANDALONE_ENV_NOTICES = list(STANDALONE_ENV_NOTICES) + [message]
|
||||||
|
return message
|
||||||
|
|
||||||
|
|
||||||
|
STANDALONE_ENV_NOTICES = []
|
||||||
|
|
||||||
# --- M11b-2: 单事务写入(事件 + 实体状态)注册接线(真源实际存在的名字) ---
|
# --- M11b-2: 单事务写入(事件 + 实体状态)注册接线(真源实际存在的名字) ---
|
||||||
from .pbl_runtime_errors import (
|
from .pbl_runtime_errors import (
|
||||||
@ -38,6 +107,7 @@ from .pbl_runtime_errors import (
|
|||||||
PblRuntimeError,
|
PblRuntimeError,
|
||||||
RuntimeWriteError,
|
RuntimeWriteError,
|
||||||
StateWriteError,
|
StateWriteError,
|
||||||
|
TransactionAborted,
|
||||||
TxAbortError,
|
TxAbortError,
|
||||||
TxConfigError,
|
TxConfigError,
|
||||||
)
|
)
|
||||||
@ -79,7 +149,14 @@ from .pbl_runtime_tx_sqlor import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
MODULE_NAME = 'world_sync'
|
||||||
|
MODULE_VERSION = '0.2.0-m11b2a'
|
||||||
|
#: 宿主 ServerEnv 上登记 conn_factory 的键(与 pbl_runtime_tx_env.SERVER_ENV_KEYS 一致)
|
||||||
|
SERVER_ENV_CONN_FACTORY_KEYS = SERVER_ENV_KEYS
|
||||||
|
|
||||||
|
|
||||||
def _get_dbname():
|
def _get_dbname():
|
||||||
|
_require_host_deps('_get_dbname')
|
||||||
return ServerEnv().get_module_dbname('world_sync')
|
return ServerEnv().get_module_dbname('world_sync')
|
||||||
|
|
||||||
|
|
||||||
@ -91,12 +168,14 @@ def _clean(ns):
|
|||||||
|
|
||||||
|
|
||||||
async def list_world_syncs(params=None):
|
async def list_world_syncs(params=None):
|
||||||
|
_require_host_deps('list_world_syncs')
|
||||||
db = DBPools()
|
db = DBPools()
|
||||||
async with db.sqlorContext(_get_dbname()) as sor:
|
async with db.sqlorContext(_get_dbname()) as sor:
|
||||||
return await sor.sqlExe("select id, world_id, sync_type, source, status, result_json, sync_date, created_at from world_sync order by created_at desc", {})
|
return await sor.sqlExe("select id, world_id, sync_type, source, status, result_json, sync_date, created_at from world_sync order by created_at desc", {})
|
||||||
|
|
||||||
|
|
||||||
async def get_world_sync(ns):
|
async def get_world_sync(ns):
|
||||||
|
_require_host_deps('get_world_sync')
|
||||||
rec_id = (ns or {}).get('id')
|
rec_id = (ns or {}).get('id')
|
||||||
if not rec_id:
|
if not rec_id:
|
||||||
return None
|
return None
|
||||||
@ -107,6 +186,7 @@ async def get_world_sync(ns):
|
|||||||
|
|
||||||
|
|
||||||
async def create_world_sync(ns):
|
async def create_world_sync(ns):
|
||||||
|
_require_host_deps('create_world_sync')
|
||||||
ns = _clean(dict(ns))
|
ns = _clean(dict(ns))
|
||||||
ns['id'] = ns.get('id') or getID()
|
ns['id'] = ns.get('id') or getID()
|
||||||
if not ns.get('sync_date'):
|
if not ns.get('sync_date'):
|
||||||
@ -121,6 +201,7 @@ async def create_world_sync(ns):
|
|||||||
|
|
||||||
|
|
||||||
async def update_world_sync(ns):
|
async def update_world_sync(ns):
|
||||||
|
_require_host_deps('update_world_sync')
|
||||||
ns = _clean(dict(ns))
|
ns = _clean(dict(ns))
|
||||||
rec_id = ns.get('id')
|
rec_id = ns.get('id')
|
||||||
if not rec_id:
|
if not rec_id:
|
||||||
@ -132,6 +213,7 @@ async def update_world_sync(ns):
|
|||||||
|
|
||||||
|
|
||||||
async def delete_world_sync(ns):
|
async def delete_world_sync(ns):
|
||||||
|
_require_host_deps('delete_world_sync')
|
||||||
rec_id = (ns or {}).get('id')
|
rec_id = (ns or {}).get('id')
|
||||||
if not rec_id:
|
if not rec_id:
|
||||||
return {'success': False, 'message': '缺少 id'}
|
return {'success': False, 'message': '缺少 id'}
|
||||||
@ -142,6 +224,7 @@ async def delete_world_sync(ns):
|
|||||||
|
|
||||||
|
|
||||||
async def execute_world_sync(ns):
|
async def execute_world_sync(ns):
|
||||||
|
_require_host_deps('execute_world_sync')
|
||||||
ns = ns or {}
|
ns = ns or {}
|
||||||
world_id = ns.get('world_id', '')
|
world_id = ns.get('world_id', '')
|
||||||
sync_type = ns.get('sync_type', '0')
|
sync_type = ns.get('sync_type', '0')
|
||||||
@ -226,6 +309,7 @@ async def pbl_runtime_event_partitions(params=None):
|
|||||||
" where table_schema=${dbname}$ and table_name=${tbl}$ "
|
" where table_schema=${dbname}$ and table_name=${tbl}$ "
|
||||||
" and partition_name is not null "
|
" and partition_name is not null "
|
||||||
" order by partition_ordinal_position")
|
" order by partition_ordinal_position")
|
||||||
|
_require_host_deps('pbl_runtime_event_partitions')
|
||||||
db = DBPools()
|
db = DBPools()
|
||||||
async with db.sqlorContext(_get_dbname()) as sor:
|
async with db.sqlorContext(_get_dbname()) as sor:
|
||||||
return await sor.sqlExe(sql, {'dbname': _get_dbname(), 'tbl': PBL_RUNTIME_EVENT_TABLE})
|
return await sor.sqlExe(sql, {'dbname': _get_dbname(), 'tbl': PBL_RUNTIME_EVENT_TABLE})
|
||||||
@ -331,12 +415,28 @@ def runtime_write_ddl(ns):
|
|||||||
return {'success': False, 'message': str(exc)}
|
return {'success': False, 'message': str(exc)}
|
||||||
|
|
||||||
|
|
||||||
def load_world_sync():
|
def load_world_sync(server_env=None, conn_factory=None):
|
||||||
"""挂载 world_sync:CRUD + M11b-1 表注册 + M11b-2 单事务写入契约。
|
"""挂载 world_sync:CRUD + M11b-1 表注册 + M11b-2 单事务写入契约(唯一真源)。
|
||||||
|
|
||||||
三处接线之「③ init.py env.xxx = xxx」;导入失败一律原样抛出(不静默降级)。
|
三处接线之「③ init.py env.xxx = xxx」;本函数同时是「② __init__.py 导出」与
|
||||||
|
根包 ``world_sync/__init__.py`` 转发的**唯一** load_world_sync 实现(QC #2:
|
||||||
|
不得存在第二份定义)。模块级符号(含 ``TransactionAborted``)导入失败一律原样
|
||||||
|
抛出(不静默降级)。
|
||||||
|
|
||||||
|
:param server_env: 宿主 ServerEnv 实例;None 时取全局单例,单例不可用时用
|
||||||
|
:class:`_StandaloneEnv` 载体完成登记(仅接线自检,不建连)。
|
||||||
|
:param conn_factory: 宿主连接池工厂(可调用/DB URL);提供则登记到
|
||||||
|
``pbl_runtime_conn_factory``(契约解析优先级第 2 级)。
|
||||||
|
:return: 登记完成的 env。
|
||||||
"""
|
"""
|
||||||
env = ServerEnv()
|
env = _resolve_env(server_env)
|
||||||
|
if conn_factory is not None:
|
||||||
|
setattr(env, SERVER_ENV_KEYS[0], conn_factory)
|
||||||
|
env.world_sync_module = MODULE_NAME
|
||||||
|
env.world_sync_version = MODULE_VERSION
|
||||||
|
env.world_sync_host_deps_available = HOST_DEPS_AVAILABLE
|
||||||
|
if not HOST_DEPS_AVAILABLE:
|
||||||
|
env.world_sync_host_deps_error = repr(HOST_DEPS_ERROR)
|
||||||
env.list_world_syncs = list_world_syncs
|
env.list_world_syncs = list_world_syncs
|
||||||
env.get_world_sync = get_world_sync
|
env.get_world_sync = get_world_sync
|
||||||
env.create_world_sync = create_world_sync
|
env.create_world_sync = create_world_sync
|
||||||
@ -371,6 +471,8 @@ def load_world_sync():
|
|||||||
env.StateWriteError = StateWriteError
|
env.StateWriteError = StateWriteError
|
||||||
env.ConcurrentStateConflict = ConcurrentStateConflict
|
env.ConcurrentStateConflict = ConcurrentStateConflict
|
||||||
env.TransactionAborted = TransactionAborted
|
env.TransactionAborted = TransactionAborted
|
||||||
|
env.RuntimeTxAbortError = TxAbortError
|
||||||
|
env.new_runtime_event_id = new_event_id
|
||||||
# 表名/方言常量与 conn_factory 解析优先级说明
|
# 表名/方言常量与 conn_factory 解析优先级说明
|
||||||
env.pbl_runtime_event_table = EVENT_TABLE
|
env.pbl_runtime_event_table = EVENT_TABLE
|
||||||
env.pbl_runtime_state_table = STATE_TABLE
|
env.pbl_runtime_state_table = STATE_TABLE
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user