refactor: unify oper transaction dispatch

This commit is contained in:
jxxghp
2026-08-24 00:04:35 +08:00
parent e9a226af02
commit 25ad924aac
5 changed files with 26 additions and 11 deletions
+2 -7
View File
@@ -14,7 +14,6 @@ from app.db.models.agenttask import (
_list_for_user_statement,
)
from app.db.models.agenttaskrun import AgentTaskRun
from app.db.uow import run_sync_transaction
class AgentTaskOper(DbOper):
@@ -63,9 +62,7 @@ class AgentTaskOper(DbOper):
)
).scalars().first()
if isinstance(self._db, Session):
return query(self._db)
return run_sync_transaction(query)
return self._execute_sync_query(query)
async def async_get(
self,
@@ -104,9 +101,7 @@ class AgentTaskOper(DbOper):
)
).scalars().all())
if isinstance(self._db, Session):
return query(self._db)
return run_sync_transaction(query)
return self._execute_sync_query(query)
def update(
self,
@@ -83,6 +83,9 @@
1. **后台任务的统一所有权已覆盖 API 入口,但仍有更深层任务机制待分级。** `app/runtime/tasks.py` 已建立 lifespan 级 TaskRegistry,启动收尾、插件 Release 刷新、Webhook E0 广播、CookieCloud E1 手工调度、消息入口、Seerr 订阅、整理历史 AI 重做、OpenAI/Anthropic 协议流和 WebAgent 断线后执行/快照保存均不再维护端点模块级任务集合或 Starlette 回调,shutdown 会停止接收、取消并有限等待,且生命周期清单明确登记其顺序。主仓 `app/` 已无裸 FastAPI `BackgroundTasks`;当前仍有约 `50` 个更底层 `create_task`/等价任务创建点,与线程池和 APScheduler 并存,后续需逐项确认 owner、取消、等待、重试、幂等和是否 durable,关键业务副作用优先接入已有 Outbox/恢复表。
2. **动态模块契约仍以 legacy 聚合语义为主。** 当前登记 `212` 个模块方法,其中 `194` 个仍使用 `legacy` aggregation,只有 `14``first_non_empty``4``ordered_list_merge``app/runtime/extensions/module/contracts.py:422-455` 已能登记 family、输入/结果标签和基础签名诊断,但 `193` 个方法没有 required parameters,调度器 `app/runtime/extensions/module/dispatcher.py:109-260` 仍主要依赖运行时反射、返回值形状和短路规则。未知第三方方法保留 legacy fallback 是兼容要求,不应删除;宿主高频能力则应逐族补齐可执行的输入校验、结果校验、超时和错误语义。
3. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。
Oper 内部的执行入口也已统一:最后一处 `AgentTaskOper` 直接 transaction runner 调用已迁入
`DbOper._execute_sync_query`,静态门禁禁止 `app/db/oper` 再绕过统一的 Session 类型分派。
4. **组合根和全局状态仍形成复杂的隐式运行时图。** Singleton 实例、模块级 provider、`configure_*` 注册函数和兼容 Facade 同时存在;它们解决了旧 ABI 和启动顺序问题,但增加测试污染、重复装配、实例身份和初始化顺序风险。`app/startup/lifecycle/__init__.py` 已有声明式生命周期,`app/startup/initializers/modules.py` 也有分阶段关闭,但尚未做到所有进程级资源都只通过 typed HostRuntime 访问。后续应以“新代码禁止新增 Service Locator/Singleton 依赖、旧入口有命中观测”为 ratchet。
Scheduler 已先完成一个可验证切片:API、Agent 与 Command 统一经
+2 -3
View File
@@ -13,8 +13,8 @@
"runtime_to_db": [],
"workflow_to_db": []
},
"edge_count": 6500,
"edge_sha256": "c27b649d471e40e8580426bdee28fdd927be03b7868a1b28bd8d9374c10be280",
"edge_count": 6499,
"edge_sha256": "f21eea0878e60d4b955a87da15cf0653eaecf07d147dc43e9eabc9f9573bb793",
"edges": [
"app -> app.runtime",
"app -> app.runtime.compat",
@@ -3660,7 +3660,6 @@
"app.db.oper.agenttask -> app.db.models",
"app.db.oper.agenttask -> app.db.models.agenttask",
"app.db.oper.agenttask -> app.db.models.agenttaskrun",
"app.db.oper.agenttask -> app.db.uow",
"app.db.oper.downloadfailure -> app.db",
"app.db.oper.downloadfailure -> app.db.base",
"app.db.oper.downloadfailure -> app.db.models",
+17
View File
@@ -432,6 +432,23 @@ def test_database_internals_do_not_import_db_facades():
assert violations == []
def test_database_opers_use_dboper_transaction_dispatchers():
"""Oper 不得绕过 DbOper 的统一 Session 类型分派直接调用事务 runner。"""
runner_names = {"run_sync_transaction", "run_async_transaction"}
violations: list[str] = []
for path in (APP_ROOT / "db" / "oper").rglob("*.py"):
tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path))
for node in ast.walk(tree):
if not isinstance(node, ast.ImportFrom) or node.module != "app.db.uow":
continue
imported = {alias.name for alias in node.names} & runner_names
if imported:
relative = path.relative_to(PROJECT_ROOT).as_posix()
violations.append(f"{relative}:{node.lineno}:{','.join(sorted(imported))}")
assert violations == []
def test_models_and_base_require_explicit_database_sessions():
"""Model/Base 不得装饰事务,且所有 db 参数必须由调用方显式传入。"""
decorator_names = {
@@ -366,7 +366,8 @@ def test_agenttask_oper_reads_with_explicit_session(db, monkeypatch):
task_id = AgentTask.add_task(db.session, **_task("canonical", user_id="alice"))
monkeypatch.setattr(
"app.db.oper.agenttask.run_sync_transaction",
db_base,
"run_sync_transaction",
lambda _query: pytest.fail("显式 Session 查询不应创建兼容事务"),
)