refactor: isolate workflow queries

This commit is contained in:
jxxghp
2026-08-23 14:06:14 +08:00
parent 4074fa4e42
commit be1071b6cd
8 changed files with 90 additions and 71 deletions
+9 -9
View File
@@ -7,7 +7,7 @@ from sqlalchemy.orm import Mapped, mapped_column
from sqlalchemy.ext.asyncio import AsyncSession
from app.db.base import Base, get_id_column
from app.db.decorators import db_query, async_db_query
from app.db.decorators import legacy_async_db_query, legacy_db_query
class Workflow(Base):
@@ -56,18 +56,18 @@ class Workflow(Base):
)
@classmethod
@db_query
@legacy_db_query
def get_enabled_workflows(cls, db):
return list(db.execute(select(cls).where(cls.state != 'P')).scalars().all())
@classmethod
@async_db_query
@legacy_async_db_query
async def async_get_enabled_workflows(cls, db: AsyncSession):
result = await db.execute(select(cls).where(cls.state != 'P'))
return list(result.scalars().all())
@classmethod
@db_query
@legacy_db_query
def get_timer_triggered_workflows(cls, db):
"""获取定时触发的工作流"""
return list(db.execute(select(cls).where(
@@ -81,7 +81,7 @@ class Workflow(Base):
)).scalars().all())
@classmethod
@async_db_query
@legacy_async_db_query
async def async_get_timer_triggered_workflows(cls, db: AsyncSession):
"""异步获取定时触发的工作流"""
result = await db.execute(select(cls).where(
@@ -96,7 +96,7 @@ class Workflow(Base):
return list(result.scalars().all())
@classmethod
@db_query
@legacy_db_query
def get_event_triggered_workflows(cls, db):
"""获取事件触发的工作流"""
return list(db.execute(select(cls).where(
@@ -107,7 +107,7 @@ class Workflow(Base):
)).scalars().all())
@classmethod
@async_db_query
@legacy_async_db_query
async def async_get_event_triggered_workflows(cls, db: AsyncSession):
"""异步获取事件触发的工作流"""
result = await db.execute(select(cls).where(
@@ -119,12 +119,12 @@ class Workflow(Base):
return list(result.scalars().all())
@classmethod
@db_query
@legacy_db_query
def get_by_name(cls, db, name: str):
return db.execute(select(cls).where(cls.name == name)).scalars().first()
@classmethod
@async_db_query
@legacy_async_db_query
async def async_get_by_name(cls, db: AsyncSession, name: str):
result = await db.execute(select(cls).where(cls.name == name))
return result.scalars().first()
+24 -10
View File
@@ -66,7 +66,7 @@ class WorkflowOper(DbOper):
新增工作流
"""
wf = Workflow(**kwargs)
if not wf.get_by_name(self._db, kwargs.get("name")):
if not self.get_by_name(kwargs.get("name")):
self._stage_create(wf)
return True, "新增工作流成功"
return False, "工作流已存在"
@@ -75,7 +75,7 @@ class WorkflowOper(DbOper):
"""
查询单个工作流
"""
return Workflow.get(self._db, wid)
return self._execute_sync_query(lambda session: Workflow.get(session, wid))
def stage_state(self, workflow_id: int, state: str) -> bool:
"""暂存工作流状态变更,不由模型方法自行提交。"""
@@ -109,49 +109,63 @@ class WorkflowOper(DbOper):
"""
异步查询单个工作流
"""
return await Workflow.async_get(self._db, wid)
return await self._execute_async_query(
lambda session: Workflow.async_get(session, wid)
)
def list(self) -> List[Workflow]:
"""
获取所有工作流列表
"""
return Workflow.list(self._db)
return self._execute_sync_query(lambda session: Workflow.list(session))
async def async_list(self) -> List[Workflow]:
"""
异步获取所有工作流列表
"""
return await Workflow.async_list(self._db)
return await self._execute_async_query(
lambda session: Workflow.async_list(session)
)
def list_enabled(self) -> List[Workflow]:
"""
获取启用的工作流列表
"""
return Workflow.get_enabled_workflows(self._db)
return self._execute_sync_query(
lambda session: Workflow.get_enabled_workflows(session)
)
def get_timer_triggered_workflows(self) -> List[Workflow]:
"""
获取定时触发的工作流列表
"""
return Workflow.get_timer_triggered_workflows(self._db)
return self._execute_sync_query(
lambda session: Workflow.get_timer_triggered_workflows(session)
)
def get_event_triggered_workflows(self) -> List[Workflow]:
"""
获取事件触发的工作流列表
"""
return Workflow.get_event_triggered_workflows(self._db)
return self._execute_sync_query(
lambda session: Workflow.get_event_triggered_workflows(session)
)
def get_by_name(self, name: str) -> Workflow:
"""
按名称获取工作流
"""
return Workflow.get_by_name(self._db, name)
return self._execute_sync_query(
lambda session: Workflow.get_by_name(session, name)
)
async def async_get_by_name(self, name: str) -> Optional[Workflow]:
"""
异步按名称获取工作流
"""
return await Workflow.async_get_by_name(self._db, name)
return await self._execute_async_query(
lambda session: Workflow.async_get_by_name(session, name)
)
async def stage_create(self, payload: Mapping[str, Any]) -> Workflow:
"""暂存新工作流,不在操作器内提交事务。"""
+2 -2
View File
@@ -378,9 +378,9 @@ flowchart LR
成功后执行。订阅新增样板由 `startup/subscription.py` 创建独占 Session
`application/subscription/write.py` 决定事务与 post-commit 边界,`SubscribeOper.stage_add()`
只查重、`add``flush`。旧 SDK 显式构造的无会话 Oper 暂留兼容自动短会话,不得被新代码复用。
`transaction-debt-baseline.json` 当前冻结 38 个正式只读查询装饰器;原有同步/异步写装饰器
`transaction-debt-baseline.json` 当前冻结 30 个正式只读查询装饰器;原有同步/异步写装饰器
已全部移除,`db_update``async_db_update` 必须持续保持为 0。下载/整理历史的旧插件 Model
调用由 `legacy_*` 兼容外壳承接,宿主 Oper 必须显式传递 Session。宿主 Oper 也不得调用 Base 保留的
与工作流旧插件 Model 调用由 `legacy_*` 兼容外壳承接,宿主 Oper 必须显式传递 Session。宿主 Oper 也不得调用 Base 保留的
`create/update/delete/truncate` 兼容包装器;AST 门禁保证显式 Session 的提交权不会被底层抢走。
- 站点、历史、工作流、Agent 会话删除和插件数据重置已经形成同构事务切片;对应 Application
Command/Service 持有 UoWOper 的 `stage_*` 方法只修改当前会话。插件数据重置从
@@ -28,7 +28,7 @@
1. **后台任务的统一所有权已覆盖 API 入口,但仍有更深层任务机制待分级。** `app/runtime/tasks.py` 已建立 lifespan 级 TaskRegistry,启动收尾、插件 Release 刷新、Webhook E0 广播、CookieCloud E1 手工调度、消息入口、Seerr 订阅和 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. **查询侧数据库兼容 ABI 仍未完全收口。** 写事务装饰器已降为 `0`,正式 `db_query/async_db_query` 已降至 `38` 个(`20` 个同步、`18` 个异步)。站点、消息、用户、订阅以及下载/整理历史的宿主 Oper 已迁到显式 Session 路径;下载/整理历史的旧插件 Model 调用由独立 `legacy_*` 外壳保留,可同时接受显式 Session 与无 Session 的位置/关键字参数。剩余正式装饰器仍会隐式创建会话,查询返回的 ORM 对象也可能跨层流转,后续继续按 Workflow、MediaServer 等风险切片迁移。
3. **查询侧数据库兼容 ABI 仍未完全收口。** 写事务装饰器已降为 `0`,正式 `db_query/async_db_query` 已降至 `30` 个(`16` 个同步、`14` 个异步)。站点、消息、用户、订阅下载/整理历史和工作流的宿主 Oper 已迁到显式 Session 路径;对应旧插件 Model 调用由独立 `legacy_*` 外壳保留,可同时接受显式 Session 与无 Session 的位置/关键字参数。剩余正式装饰器仍会隐式创建会话,查询返回的 ORM 对象也可能跨层流转,后续继续按 MediaServer、SiteUserData 等风险切片迁移。
4. **组合根和全局状态仍形成复杂的隐式运行时图。** Singleton 实例、模块级 provider、`configure_*` 注册函数和兼容 Facade 同时存在;它们解决了旧 ABI 和启动顺序问题,但增加测试污染、重复装配、实例身份和初始化顺序风险。`app/startup/lifecycle/__init__.py:161-376` 已有声明式生命周期,`app/startup/modules_initializer.py:505-530` 也有分阶段关闭,但尚未做到所有进程级资源都只通过 typed HostRuntime 访问。后续应以“新代码禁止新增 Service Locator/Singleton 依赖、旧入口有命中观测”为 ratchet。
### P2:中长期可演进性债务
@@ -904,7 +904,7 @@ ADR 必须逐个映射当前 Event、BackgroundTasks、Scheduler job、Agent tas
事务低水位从 174 降到 168Oper 仍不创建 Session、也不直接 commit/rollback。
- 剩余同步/异步 Model 写装饰器已全部迁移:AgentTask、PassKey、User、消息、历史清理、
站点快照、媒体服务器、插件数据、TransferPending 等写入由调用方 Session 和 UoW 收口;无 Session
的旧 Oper ABI 委托 Startup 注入的短事务执行器。当前 Model 正式查询装饰器仅剩 38 个(同步 20、异步 18),
的旧 Oper ABI 委托 Startup 注入的短事务执行器。当前 Model 正式查询装饰器仅剩 30 个(同步 16、异步 14),
`db_update``async_db_update` 均为 0Oper 自建 Session/直接提交仍为 0。
- 数据清理按批次显式提交 UoW,单表失败先回滚会话再继续汇总后续表;不再依赖删除 Model 的隐式提交。
- 收尾批次进一步移除宿主 Oper 对 `Base.create/update/delete/truncate` 八个兼容包装器的调用:显式
@@ -1183,6 +1183,11 @@ host/plugin 架构基线均通过。
按签名插入一次性会话,兼容无 Session 的位置参数和关键字参数,同时显式 Session 不创建额外会话。
历史查询、删除工具、类型门禁和插件架构专项共 `101 passed`host/plugin 架构基线通过。
2026-08-23 完成 Workflow 查询切片:`WorkflowOper` 的同步/异步查询入口统一通过
`_execute_sync_query` / `_execute_async_query` 复用调用方 Session,正式查询装饰器由 38 降至
30 个且写装饰器保持 0。旧插件仍可直接调用 Workflow Model 方法,显式 Session 与无 Session 的
关键字调用均有回归覆盖;Workflow、架构基线专项共 `76 passed`host/plugin 架构基线通过。
#### ARCH-272:异步阻塞检测
**目标**:对新 API/Agent/Application async 路径检测 `open`、文件遍历、同步 HTTP、阻塞 sleep 和重 CPU 解析。
@@ -1360,7 +1365,7 @@ rollback:
| 基线写入行为 | 默认命令可能覆盖 fixture | 所有默认/check 命令保证工作树不变;write 必须显式 scope |
| 全功能 worker | 配置允许 >1,控制面会复制 | 启动期明确拒绝 >1;文档与配置一致 |
| 健康接口 | 认证 `/system/ping` 为主 | 分离公开 live 与受限/安全 ready;失败原因可诊断 |
| Model 事务装饰器 | 当前 38 个且全部只读;写装饰器 0 | 查询债务只降不增;写事务不回退到 Model/Base 隐式提交 |
| Model 事务装饰器 | 当前 30 个且全部只读;写装饰器 0 | 查询债务只降不增;写事务不回退到 Model/Base 隐式提交 |
| 新写用例事务 | 宿主写 Oper 已脱离 Base 隐式提交 | 100% 由入口/Application 边界拥有 Session/UoW |
| 高频 Module 契约 | 212 个宿主能力显式登记 | 新观察到的宿主方法必须同步登记完整契约 |
| Event payload | 53 类型全部登记 typed payload 与可靠性 | 新事件必须同步登记,不回退裸 dict |
+1 -1
View File
@@ -84,7 +84,7 @@ Oper classes accept and return persistence values. Turning a `MediaInfo` or
### Transaction ownership ratchet
- `tests/fixtures/architecture/transaction-debt-baseline.json` records the
existing Model transaction decorators. The current 38 decorators are query-only
existing Model transaction decorators. The current 30 decorators are query-only
migration debt: they may decrease but must never increase or move to a new
Model method. Both `db_update` and `async_db_update` must remain at zero.
- `legacy_db_query` / `legacy_async_db_query` are compatibility-only shells for
+3 -43
View File
@@ -1,12 +1,12 @@
{
"model_decorators": {
"by_kind": {
"async_db_query": 18,
"async_db_query": 14,
"async_db_update": 0,
"db_query": 20,
"db_query": 16,
"db_update": 0
},
"count": 38,
"count": 30,
"methods": [
{
"decorator": "async_db_query",
@@ -157,46 +157,6 @@
"decorator": "db_query",
"file": "app/db/models/transferpending.py",
"method": "TransferPending.list_all"
},
{
"decorator": "async_db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.async_get_by_name"
},
{
"decorator": "async_db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.async_get_enabled_workflows"
},
{
"decorator": "async_db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.async_get_event_triggered_workflows"
},
{
"decorator": "async_db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.async_get_timer_triggered_workflows"
},
{
"decorator": "db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.get_by_name"
},
{
"decorator": "db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.get_enabled_workflows"
},
{
"decorator": "db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.get_event_triggered_workflows"
},
{
"decorator": "db_query",
"file": "app/db/models/workflow.py",
"method": "Workflow.get_timer_triggered_workflows"
}
]
},
+2 -2
View File
@@ -126,8 +126,8 @@ def test_transaction_debt_baseline_is_a_model_and_oper_ratchet() -> None:
baseline = json.loads(baseline_path.read_text(encoding="utf-8"))
assert baseline["schema_version"] == 1
assert baseline["model_decorators"]["count"] == 38
assert sum(baseline["model_decorators"]["by_kind"].values()) == 38
assert baseline["model_decorators"]["count"] == 30
assert sum(baseline["model_decorators"]["by_kind"].values()) == 30
assert baseline["model_decorators"]["by_kind"]["db_update"] == 0
assert baseline["model_decorators"]["by_kind"]["async_db_update"] == 0
assert baseline["model_transaction_calls"] == {"count": 0, "calls": []}
+41 -1
View File
@@ -9,8 +9,10 @@ import asyncio
import pytest
from app.db import decorators
from app.db.models.workflow import Workflow
from app.db.session import async_session_scope
from app.db.oper.workflow import WorkflowOper
from app.db.session import SessionFactory, async_session_scope
@pytest.fixture(autouse=True)
@@ -58,6 +60,44 @@ def test_list_and_get_by_name_match_async_twins(db):
assert sync_ids == async_ids
def test_workflow_oper_reuses_explicit_query_sessions(db, monkeypatch):
"""WorkflowOper 绑定显式会话后不得再创建兼容查询会话。"""
created = db.add(_flow("wf-explicit-session"))
monkeypatch.setattr(
decorators,
"ScopedSession",
lambda: (_ for _ in ()).throw(AssertionError("不应创建额外同步会话")),
)
assert WorkflowOper(db.session).get_by_name(created.name).id == created.id
async def check() -> None:
"""验证异步 Oper 同样复用调用方会话。"""
async with async_session_scope() as session:
monkeypatch.setattr(
decorators,
"async_session_scope",
lambda: (_ for _ in ()).throw(AssertionError("不应创建额外异步会话")),
)
assert (await WorkflowOper(session).async_get_by_name(created.name)).id == created.id
asyncio.run(check())
def test_workflow_model_legacy_queries_keep_no_session_abi(db, monkeypatch):
"""旧插件直接调用 Workflow Model 时仍应按签名自动补入短会话。"""
created = db.add(_flow("wf-legacy-query"))
opened = []
monkeypatch.setattr(
decorators,
"ScopedSession",
lambda: (opened.append(True) or SessionFactory()),
)
assert Workflow.get_by_name(name=created.name).id == created.id
assert opened == [True]
def test_enabled_workflows_exclude_paused(db):
"""
启用列表排除暂停状态。