deliver: 交付收口(引擎代为提交)

This commit is contained in:
agent.develop 2026-09-20 17:59:38 +08:00
parent cf9baea7e9
commit b514c1bdd9
2 changed files with 206 additions and 62 deletions

View File

@ -1,3 +1,25 @@
# -*- coding: utf-8 -*-
"""world_sync 宿主挂载层(module-development-spec:init.py 三处同步注册)。
职责:
1. world_sync 业务表 CRUD(list/get/create/update/delete/execute);
2. M11b-1 表定义注册(pbl_runtime_event / pbl_entity_state schema 元数据,零写入);
3. M11b-2 单事务「事件 + 实体状态」原子写入契约登记到 ServerEnv
(``env.write_event_with_state`` / ``env.read_entity_state`` / ``env.pbl_runtime_conn_factory``)。
铁律:
- 不硬编码库名:取库名走 ``ServerEnv().get_module_dbname('world_sync')``;
- 导入期不建立任何数据库连接(连接工厂缺失时由写入调用点 fail-closed 抛异常);
- 不吞错误:本文件所有 import 均为**硬依赖**,任一失败即抛(QC 退回意见 #1/#2)。
M11b-2a 迁移修正(响应 QC #1):本文件此前导入的
``RuntimeWriteError/InvalidWriteRequest/ConnectionUnavailable/TransactionAborted`` 与
``RuntimeStateEventWriter/build_plan/configure_writer/get_writer/write_event_with_state_async``
以及 ``pbl_runtime_sql.MYSQL/SQLITE`` 在真源模块中并不存在(照抄 apps 侧旧 API 草案),
现改为只导入迁移后的真源实际存在的名字;旧语义名以「同一类对象的别名」形式保留
(见 pbl_runtime_errors.py 末尾别名段),不存在双真源。
"""
import json import json
import os import os
@ -6,33 +28,54 @@ from appPublic.timeUtils import curDateString, timestampstr
from sqlor.dbpools import DBPools from sqlor.dbpools import DBPools
from ahserver.serverenv import ServerEnv from ahserver.serverenv import ServerEnv
# --- M11b-2: 单事务写入(事件 + 实体状态)注册接线 --- # --- M11b-2: 单事务写入(事件 + 实体状态)注册接线(真源实际存在的名字) ---
from .pbl_runtime_errors import ( from .pbl_runtime_errors import (
RuntimeWriteError,
InvalidWriteRequest,
ConnectionUnavailable,
StateWriteError,
EventWriteError,
ConcurrentStateConflict, ConcurrentStateConflict,
TransactionAborted, ConnFactoryNotConfigured,
ConnectionUnavailable,
EventWriteError,
InvalidWriteRequest,
PblRuntimeError,
RuntimeWriteError,
StateWriteError,
TxAbortError,
TxConfigError,
) )
from .pbl_runtime_sql import STATE_TABLE, MYSQL, SQLITE from .pbl_runtime_sql import (
from .pbl_runtime_tx import ( DIALECT_MYSQL,
RuntimeStateEventWriter, DIALECT_SQLITE,
build_plan, EVENT_TABLE,
configure_writer, STATE_TABLE,
get_writer, ddl_for,
detect_dialect,
percentile, percentile,
write_event_with_state, )
write_event_with_state_async, from .pbl_runtime_tx import (
DEFAULT_EVENT_TYPE,
DEFAULT_SOURCE,
EVENT_REQUIRED,
STATE_KEY_COLUMNS,
new_event_id,
new_event_uid,
read_entity_state as read_entity_state_via_tx,
utc_now_text,
write_event_with_state as write_event_with_state_via_tx,
)
from .pbl_runtime_tx_env import (
SERVER_ENV_KEYS,
build_mysql_conn_factory,
build_sqlite_conn_factory,
load_db_config,
module_conf_dir,
resolve_conn_factory,
) )
from .pbl_runtime_tx_sqlor import ( from .pbl_runtime_tx_sqlor import (
build_connection_factory, ENV_NAME_KEYS,
build_writer, POOL_KEYS,
configure_default_writer, build_conn_factory_from_pools,
detect_driver, host_conn_factory,
load_db_conf, read_entity_state,
resolve_dialect, write_event_with_state,
) )
@ -190,34 +233,70 @@ async def pbl_runtime_event_partitions(params=None):
# =========================================================================== # ===========================================================================
# M11b-2: 单事务写入(事件 + 实体状态原子落库)—— 运行时接线 # M11b-2: 单事务写入(事件 + 实体状态原子落库)—— 运行时接线
# 范围铁律:本段只做「契约函数登记 + 默认写入器配置」, # 范围铁律:本段只做「契约函数登记 + 连接工厂登记」,
# 不含广播(属 M11b-3)、不含轮询兜底、不含时延优化。 # 不含广播(属 M11b-3)、不含轮询兜底、不含时延优化。
# 写入器采用「懒失败」:init 阶段不建连,真正写入时才取连接; # 写入器为**无状态函数式 API**(每次调用按三级优先级解析 conn_factory),
# 参数不全则抛 ConnectionUnavailable(fail-closed,不静默降级、不吞错误)。 # init 阶段不建连;参数不全则调用点 fail-closed 抛 ConnectionUnavailable。
# =========================================================================== # ===========================================================================
def configure_runtime_writer(conn_factory=None, dialect=None, conf=None, def configure_runtime_writer(conn_factory=None, server_env=None, dialect=None):
driver=None, before_commit=None, """把宿主连接工厂登记到 ServerEnv(契约解析优先级第 2 级),返回登记后的工厂。
close_connection=True):
"""应用侧/测试侧配置默认单事务写入器,返回写入器实例。
与 :func:`pbl_runtime_tx_sqlor.configure_default_writer` 等价,只是暴露给 真源实现(pbl_runtime_tx / pbl_runtime_tx_env)是函数式 API,不持有全局 writer
ServerEnv 的名字更贴近业务语义(上层 dspy 只需 ``env.configure_runtime_writer``)。 单例;因此「配置写入器」在 M11b-2a 后的语义 = 登记 conn_factory + 记录方言名。
``dialect`` 可传字符串('mysql'/'sqlite')或方言对象。
:param conn_factory: 可调用对象(返回 DB-API 连接)或 DB URL 字符串
:param server_env: 目标 ServerEnv(None 时取全局单例)
:param dialect: 可选方言名('mysql'/'sqlite'),仅登记供诊断用
:raises InvalidWriteRequest: conn_factory 不可调用且不是字符串
""" """
return configure_default_writer( if conn_factory is None:
conn_factory=conn_factory, dialect=dialect, conf=conf, driver=driver, raise ConnectionUnavailable(
before_commit=before_commit, close_connection=close_connection) "configure_runtime_writer requires an explicit conn_factory "
"(callable or DB URL); refusing to register nothing",
server_env_keys=list(SERVER_ENV_KEYS))
if not (callable(conn_factory) or isinstance(conn_factory, str)):
raise InvalidWriteRequest(
"conn_factory must be callable or a DB URL string",
got=type(conn_factory).__name__)
env = server_env if server_env is not None else ServerEnv()
setattr(env, SERVER_ENV_KEYS[0], conn_factory)
if dialect:
setattr(env, 'pbl_runtime_dialect', str(dialect).lower())
setattr(env, 'pbl_runtime_writer_ready', True)
try:
setattr(env, 'pbl_runtime_module_conf_dir', module_conf_dir())
except Exception as exc: # 仅诊断信息,不影响契约可用性
setattr(env, 'pbl_runtime_module_conf_dir_error', str(exc))
return conn_factory
def get_runtime_writer(): def get_runtime_writer(conn_factory=None):
"""取当前默认写入器;未配置时抛 ConnectionUnavailable。""" """解析当前可用的连接工厂(三级优先级);全空抛 ConnectionUnavailable(fail-closed)。"""
return get_writer() return resolve_conn_factory(conn_factory)
def build_runtime_write_plan(event_data, state_update, schemas=None, now=None): def build_runtime_write_plan(event_data, state_update, event_uid=None, source=None, now=None):
"""把 (event_data, state_update) 装配成 WritePlan(纯函数,不落库,便于校验/调试)。""" """预演一次写入的标识装配(不落库):返回 event_uid / event_type / source / 时间戳。
return build_plan(event_data, state_update, schemas=schemas, now=now)
真源没有独立的 plan 对象(写入即执行),本函数只提供「调用前自检」能力,
便于接口层与自测脚本在不碰数据库的情况下确认参数装配是否完整。
"""
event_data = dict(event_data or {})
state_update = dict(state_update or {})
uid = event_uid or event_data.get('event_uid') or event_data.get('event_id') or new_event_uid()
return {
'event_uid': uid,
'event_type': event_data.get('event_type') or DEFAULT_EVENT_TYPE,
'source': source or event_data.get('source') or DEFAULT_SOURCE,
'tenant_id': event_data.get('tenant_id') or state_update.get('tenant_id'),
'session_id': event_data.get('session_id') or state_update.get('session_id'),
'entity_id': event_data.get('entity_id') or state_update.get('entity_id'),
'event_required': list(EVENT_REQUIRED),
'state_key_columns': list(STATE_KEY_COLUMNS),
'tables': {'event': EVENT_TABLE, 'state': STATE_TABLE},
'created_at': utc_now_text(now) if now else utc_now_text(),
}
def runtime_write_percentile(values, p=99): def runtime_write_percentile(values, p=99):
@ -226,14 +305,37 @@ def runtime_write_percentile(values, p=99):
def runtime_write_dialect_name(params=None): def runtime_write_dialect_name(params=None):
"""当前默认写入器使用的方言名(未配置时返回 None,不抛错,只读语义)。""" """当前登记的方言名(未登记时返回 None,不建连、不抛错,纯只读语义)。"""
try: try:
return get_writer().dialect.name env = ServerEnv()
except RuntimeWriteError: except Exception:
return None return None
name = getattr(env, 'pbl_runtime_dialect', None)
if name:
return name
getter = getattr(env, 'get', None)
if callable(getter):
try:
return getter('pbl_runtime_dialect', None)
except Exception:
return None
return None
def runtime_write_ddl(ns):
"""按方言返回建表 DDL(只读文本,供部署/自测建表,不执行)。"""
dialect = (ns or {}).get('dialect') or DIALECT_SQLITE
try:
return {'success': True, 'dialect': dialect, 'ddl': ddl_for(dialect)}
except Exception as exc:
return {'success': False, 'message': str(exc)}
def load_world_sync(): def load_world_sync():
"""挂载 world_sync:CRUD + M11b-1 表注册 + M11b-2 单事务写入契约。
三处接线之「③ init.py env.xxx = xxx」;导入失败一律原样抛出(不静默降级)。
"""
env = ServerEnv() env = ServerEnv()
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
@ -250,43 +352,58 @@ def load_world_sync():
env.pbl_runtime_event_partitions = pbl_runtime_event_partitions env.pbl_runtime_event_partitions = pbl_runtime_event_partitions
# --- M11b-2: 单事务写入契约登记(事件 + 实体状态) --- # --- M11b-2: 单事务写入契约登记(事件 + 实体状态) ---
env.write_event_with_state = write_event_with_state env.write_event_with_state = write_event_with_state
env.write_event_with_state_async = write_event_with_state_async env.read_entity_state = read_entity_state
env.configure_runtime_writer = configure_runtime_writer env.configure_runtime_writer = configure_runtime_writer
env.get_runtime_writer = get_runtime_writer env.get_runtime_writer = get_runtime_writer
env.build_runtime_write_plan = build_runtime_write_plan env.build_runtime_write_plan = build_runtime_write_plan
env.runtime_write_percentile = runtime_write_percentile env.runtime_write_percentile = runtime_write_percentile
env.runtime_write_dialect_name = runtime_write_dialect_name env.runtime_write_dialect_name = runtime_write_dialect_name
env.RuntimeStateEventWriter = RuntimeStateEventWriter env.runtime_write_ddl = runtime_write_ddl
env.new_runtime_event_uid = new_event_uid
env.utc_now_text = utc_now_text
# 异常族(fail-closed,供接口层 except)
env.RuntimeWriteError = RuntimeWriteError env.RuntimeWriteError = RuntimeWriteError
env.PblRuntimeError = PblRuntimeError
env.TxConfigError = TxConfigError
env.ConnectionUnavailable = ConnectionUnavailable
env.InvalidWriteRequest = InvalidWriteRequest
env.EventWriteError = EventWriteError
env.StateWriteError = StateWriteError
env.ConcurrentStateConflict = ConcurrentStateConflict
env.TransactionAborted = TransactionAborted
# 表名/方言常量与 conn_factory 解析优先级说明
env.pbl_runtime_event_table = EVENT_TABLE
env.pbl_runtime_state_table = STATE_TABLE env.pbl_runtime_state_table = STATE_TABLE
env.pbl_runtime_tx_dialects = {'mysql': MYSQL, 'sqlite': SQLITE} env.pbl_runtime_tx_dialects = {'mysql': DIALECT_MYSQL, 'sqlite': DIALECT_SQLITE}
env.pbl_runtime_conn_factory_keys = list(SERVER_ENV_KEYS)
env.pbl_runtime_pool_keys = list(POOL_KEYS)
env.pbl_runtime_env_name_keys = list(ENV_NAME_KEYS)
_configure_runtime_writer_from_env(env) _configure_runtime_writer_from_env(env)
_register_tables_to_env(env) _register_tables_to_env(env)
return env return env
def _configure_runtime_writer_from_env(env): def _configure_runtime_writer_from_env(env):
"""按环境注入 conn_factory 并配置默认写入器(懒失败,init 不建连)。 """按环境尽力登记 conn_factory(懒失败,init 阶段绝不建连)。
驱动来源:环境变量 WORLD_SYNC_TX_DRIVER / PBL_RUNTIME_TX_DRIVER / DB_DRIVER 顺序:宿主连接池(sqlor pools)→ 环境变量 PBL_RUNTIME_DB_URL → 模块 conf/db.json。
(sqlite 仅本地/测试);缺省 pymysql + MySQL 方言(生产口径)。 任何「配置不全」都只记录原因并返回 None —— 这不是吞错误:写入契约一旦被调用
任何异常都不吞:记录到 env 属性并原样抛出,让应用启动失败可见。 即 fail-closed 抛 :class:`ConnectionUnavailable`(见 resolve_conn_factory)。
""" """
driver = detect_driver()
conf = load_db_conf()
try: try:
writer = configure_runtime_writer( factory = host_conn_factory()
conn_factory=build_connection_factory(conf=conf, driver=driver), except PblRuntimeError as exc:
dialect=resolve_dialect(driver=driver),
close_connection=True)
except RuntimeWriteError as exc:
# 连接参数不全/方言未知:不阻断应用启动(其他模块功能不受影响),
# 但写入契约一旦调用即 fail-closed,且把原因挂在 env 上便于诊断。
env.pbl_runtime_writer_error = exc.as_dict() env.pbl_runtime_writer_error = exc.as_dict()
env.pbl_runtime_writer_ready = False
return None return None
except Exception as exc: # 宿主 API 差异等意外:记录并原样包装,不静默
env.pbl_runtime_writer_error = {
'code': 'PBL_TX_HOST_PROBE_FAILED', 'message': str(exc)}
env.pbl_runtime_writer_ready = False
return None
setattr(env, SERVER_ENV_KEYS[0], factory)
env.pbl_runtime_writer_ready = True env.pbl_runtime_writer_ready = True
env.pbl_runtime_writer_dialect = writer.dialect.name return factory
return writer
def _register_tables_to_env(env): def _register_tables_to_env(env):

View File

@ -1,5 +1,5 @@
# -*- coding: utf-8 -*- # -*- coding: utf-8 -*-
"""M11b-2 运行时单事务写入 —— 异常定义。 """M11b-2 运行时单事务写入 —— 异常定义(唯一真源)。
设计要点(对应需求点 3「错误处理:事务失败时回滚并抛出明确异常,不吞错误」): 设计要点(对应需求点 3「错误处理:事务失败时回滚并抛出明确异常,不吞错误」):
1. 所有异常统一继承 :class:`PblRuntimeError`,携带机器可判定的 ``code`` 与结构化 ``detail``; 1. 所有异常统一继承 :class:`PblRuntimeError`,携带机器可判定的 ``code`` 与结构化 ``detail``;
@ -8,6 +8,12 @@
让上层可以区分「可重试」与「真失败」,而不是笼统的 500。 让上层可以区分「可重试」与「真失败」,而不是笼统的 500。
本模块零依赖(只用标准库),不引入广播 / 轮询相关任何东西。 本模块零依赖(只用标准库),不引入广播 / 轮询相关任何东西。
M11b-2a 迁移补充(响应 QC 退回意见 #1「补齐异常别名」):宿主挂载层
``world_sync/init.py`` 与已迁移的测试用例使用一组**语义别名**
(``RuntimeWriteError`` / ``InvalidWriteRequest`` / ``ConnectionUnavailable`` /
``TransactionAborted``)。别名与真源类是**同一个类对象**(``is`` 相等),
因此 ``except`` 与 ``assertRaises`` 在两套名字下行为完全一致,不存在双真源。
""" """
__all__ = [ __all__ = [
@ -18,6 +24,11 @@ __all__ = [
"StateWriteError", "StateWriteError",
"ConcurrentStateConflict", "ConcurrentStateConflict",
"TxAbortError", "TxAbortError",
# --- M11b-2a 语义别名(与上方真源同一对象,供宿主层/测试层使用)---
"RuntimeWriteError",
"InvalidWriteRequest",
"ConnectionUnavailable",
"TransactionAborted",
] ]
@ -45,6 +56,9 @@ class PblRuntimeError(Exception):
"http_hint": self.http_hint, "http_hint": self.http_hint,
} }
# 兼容旧调用名(init.py / 测试用 as_dict())
as_dict = to_dict
def __str__(self): def __str__(self):
if self.detail: if self.detail:
parts = ", ".join( parts = ", ".join(
@ -98,3 +112,16 @@ class TxAbortError(PblRuntimeError):
code = "PBL_TX_ABORTED" code = "PBL_TX_ABORTED"
http_hint = 500 http_hint = 500
# ===========================================================================
# M11b-2a 语义别名 —— 同一个类对象,不是新异常族(避免双真源)
# RuntimeWriteError → PblRuntimeError(基类,宿主层统一 except 用)
# InvalidWriteRequest → TxConfigError(调用方参数/配置错误)
# ConnectionUnavailable → ConnFactoryNotConfigured(拿不到连接工厂,fail-closed)
# TransactionAborted → TxAbortError(提交阶段失败,已回滚)
# ===========================================================================
RuntimeWriteError = PblRuntimeError
InvalidWriteRequest = TxConfigError
ConnectionUnavailable = ConnFactoryNotConfigured
TransactionAborted = TxAbortError