refactor: unify plugin file operation cancellation

This commit is contained in:
jxxghp
2026-08-24 03:17:07 +08:00
parent e768eb46fe
commit 10582277e2
7 changed files with 38 additions and 35 deletions
+3 -14
View File
@@ -44,6 +44,9 @@ from app.adapters.system.plugin.manifest import (
)
from app.runtime.log import logger
from app.runtime.observability import observe_compat_facade
from app.runtime.execution import (
run_in_threadpool_to_completion as _await_thread_operation,
)
from app.runtime.tasks import get_task_registry
from app.adapters.network.http import RequestUtils, AsyncRequestUtils
from app.foundation.singleton import WeakSingleton
@@ -82,20 +85,6 @@ def _empty_installed_plugins() -> List[str]:
_installed_plugins_provider: InstalledPluginsProvider = _empty_installed_plugins
async def _await_thread_operation(func, *args, **kwargs):
"""取消请求到达时先等待同步插件操作收口,避免目录写入继续进行。"""
task = asyncio.create_task(asyncio.to_thread(func, *args, **kwargs))
try:
return await asyncio.shield(task)
except asyncio.CancelledError:
try:
await asyncio.shield(task)
except BaseException:
pass
raise
def configure_installed_plugins_provider(
provider: InstalledPluginsProvider,
) -> None:
+3 -15
View File
@@ -2,7 +2,6 @@
from __future__ import annotations
import asyncio
import re
import shutil
import uuid
@@ -11,6 +10,9 @@ from pathlib import Path
from typing import Optional
from app.adapters.external.market import PluginHelper as _PluginHelper
from app.runtime.execution import (
run_in_threadpool_to_completion as _await_thread_operation,
)
from app.runtime.log import logger
from app.runtime.settings import RuntimeSettingsCompat
@@ -18,20 +20,6 @@ from app.runtime.settings import RuntimeSettingsCompat
# 保留旧模块级入口,插件本地同步测试和旧扩展仍可能覆盖这些设置。
settings = RuntimeSettingsCompat()
async def _await_thread_operation(func, *args, **kwargs):
"""取消请求到达时先等待文件操作收口,避免后台线程继续写运行目录。"""
task = asyncio.create_task(asyncio.to_thread(func, *args, **kwargs))
try:
return await asyncio.shield(task)
except asyncio.CancelledError:
try:
await asyncio.shield(task)
except BaseException:
pass
raise
@dataclass(frozen=True, slots=True)
class PluginPackageCheckpoint:
"""记录一次插件包变更前可用于补偿恢复的文件快照。"""
@@ -95,6 +95,14 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同
- Scheduler 的协程作业与异步进度收尾由 Scheduler 自有句柄表持有;同步 `start()` / `stop()` ABI 保持,
生命周期关闭入口等待目标事件循环确认真实收尾,跨线程取消代理不作为任务完成凭据。
### Plugin package mutations
- 插件快照、安装、回滚和持久备份刷新中的同步文件操作属于 E3 步骤;取消只能延迟传播到同步 worker
到达终态后,不能让文件仍在写入时释放插件 mutation owner。
- 市场适配器和插件包适配器统一复用 `runtime.execution.run_in_threadpool_to_completion`;两个模块内的
`_await_thread_operation` 私有接缝保留为同一函数别名,不改变 `PluginHelper``PluginPackageManager`
的同步/异步调用合同。
### Transfer pending / 文件整理
- transfer pending、队列任务和实际文件移动:E3。完成点是文件步骤、历史状态和必要清理均达到一致终态。
@@ -6,7 +6,7 @@
> 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本
> 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文
> 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md`
> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权。
> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义
## 当前复核结论(2026-08-24
@@ -15,7 +15,7 @@
### 长期整改阶段 0:治理门禁恢复(2026-08-23
- 宿主依赖基线已审查 TaskRegistry、有界后台 owner 与插件变更准入接入后的语义差异:当前为 `806` 个模块、`6529` 条内部导入边,12 组重点禁止边继续全部为 `0`,唯一非平凡 SCC 仍是隔离的 TMDB 移植包。
- 宿主依赖基线已审查 TaskRegistry、有界后台 owner 与插件变更准入接入后的语义差异:当前为 `806` 个模块、`6531` 条内部导入边,12 组重点禁止边继续全部为 `0`,唯一非平凡 SCC 仍是隔离的 TMDB 移植包。
- 启动性能探针会在隔离生命周期中真实创建并释放 TaskRegistrynormal/safe 组件数分别为 `23`/`11`,CI 只读检查使用稳定的宿主模块集合和生命周期组件顺序,不再把 Python/平台模块数量当作硬合同。
- 官方插件快照覆盖 `plugins.v3``plugins.v2` 以及 V3 实际会从 `package.json` 回退加载的 31 个默认实现;`app/plugins/**` 仍只是宿主运行副本,不进入扫描。
- SDK 快照以各模块显式 `__all__` 为公开合同,能够记录赋值别名;`typing``__future__` 等实现期导入不再被误冻结,既有数据库备份门面已补精确导出清单。
@@ -112,7 +112,7 @@
- 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。
- `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。
- 依赖图当前为 `806` 个 Python 模块、`6529` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
- 依赖图当前为 `806` 个 Python 模块、`6531` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
- 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。
综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。
@@ -136,6 +136,8 @@
渠道交给生命周期持有的共享线程池;这仍是 E0 投递,不宣称跨进程恢复。
图片代理安全日志的窗口聚合也已补齐内部任务所有权:timer 到期创建的 flush task 由 coalescer 持有,
`stop_modules()` 会刷新未到期摘要并等待已启动回调;它属于 E1 观测,不扩大 TaskRegistry 或 Outbox 范围。
插件市场与插件包适配器原有两套线程取消 wrapper 也已统一到 `runtime.execution`:连续取消必须等同步
文件 worker 到达终态后再传播,避免提前释放 mutation owner;旧模块内私有名称保留为 canonical 别名。
2. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。
Oper 内部的执行入口也已统一:最后一处 `AgentTaskOper` 直接 transaction runner 调用已迁入
+4 -2
View File
@@ -13,8 +13,8 @@
"runtime_to_db": [],
"workflow_to_db": []
},
"edge_count": 6529,
"edge_sha256": "daf8459c6194e28766d62c5f16a0e1bb64caf3f3430f51736b57da505e4401e1",
"edge_count": 6531,
"edge_sha256": "86b2ad0b585deba78c5e0e5e7e4b55cd66e05a2d5e988ed28ac7f6a4ae15ff2f",
"edges": [
"app -> app.runtime",
"app -> app.runtime.compat",
@@ -62,6 +62,7 @@
"app.adapters.external.market -> app.foundation.version",
"app.adapters.external.market -> app.runtime",
"app.adapters.external.market -> app.runtime.cache",
"app.adapters.external.market -> app.runtime.execution",
"app.adapters.external.market -> app.runtime.log",
"app.adapters.external.market -> app.runtime.observability",
"app.adapters.external.market -> app.runtime.settings",
@@ -145,6 +146,7 @@
"app.adapters.system.plugin.package -> app.adapters.external",
"app.adapters.system.plugin.package -> app.adapters.external.market",
"app.adapters.system.plugin.package -> app.runtime",
"app.adapters.system.plugin.package -> app.runtime.execution",
"app.adapters.system.plugin.package -> app.runtime.log",
"app.adapters.system.plugin.package -> app.runtime.settings",
"app.adapters.system.resource -> app.adapters",
+4 -1
View File
@@ -196,7 +196,10 @@ def _patch_async_remote_install(helper, monkeypatch, meta: dict,
monkeypatch.setattr(helper, "_PluginHelper__async_install_dependencies_if_required", fake_dependencies)
monkeypatch.setattr(helper, "_PluginHelper__async_install_from_release", fake_release)
monkeypatch.setattr(helper, "_PluginHelper__prepare_content_via_filelist_async", fake_filelist)
monkeypatch.setattr("app.adapters.external.market.asyncio.to_thread", fake_to_thread)
monkeypatch.setattr(
"app.adapters.external.market._await_thread_operation",
fake_to_thread,
)
return calls
+11
View File
@@ -6,9 +6,20 @@ import threading
import pytest
from anyio.to_thread import current_default_thread_limiter
from app.adapters.external import market as market_adapter
from app.adapters.system.plugin import package as plugin_package_adapter
from app.runtime.execution import run_in_threadpool_to_completion
def test_plugin_file_adapters_share_runtime_completion_contract() -> None:
"""市场与插件包适配器不得各自维护另一套线程取消实现。"""
assert market_adapter._await_thread_operation is run_in_threadpool_to_completion
assert (
plugin_package_adapter._await_thread_operation
is run_in_threadpool_to_completion
)
@pytest.mark.asyncio
async def test_threadpool_capacity_is_held_until_cancelled_call_finishes() -> None:
"""调用方取消后,执行令牌必须由真实同步调用持有到终态。"""