refactor: unify mediaserver chain data port

This commit is contained in:
jxxghp
2026-08-24 04:30:05 +08:00
parent 836cbfbbb2
commit fb08e06277
5 changed files with 47 additions and 15 deletions
+4 -4
View File
@@ -4,7 +4,7 @@ from typing import Callable, Dict, List, Union, Optional, Generator, Any, Tuple
from app.chain import ChainBase from app.chain import ChainBase
from app.runtime.config import global_vars from app.runtime.config import global_vars
from app.application.chain.data import MediaServerPortProxy as MediaServerOper from app.application.chain.data import get_chain_media_server_port
from app.runtime.extensions.service_config import ServiceConfigHelper from app.runtime.extensions.service_config import ServiceConfigHelper
from app.runtime.log import logger from app.runtime.log import logger
from app.schemas.mediaserver import MediaServerLibrary from app.schemas.mediaserver import MediaServerLibrary
@@ -284,7 +284,7 @@ class MediaServerChain(ChainBase):
item.name for item in mediaservers item.name for item in mediaservers
if item and item.enabled and item.name if item and item.enabled and item.name
] ]
dboper = MediaServerOper() dboper = get_chain_media_server_port()
dboper.delete_excluded_servers(enabled_servers) dboper.delete_excluded_servers(enabled_servers)
selected_servers = [ selected_servers = [
item for item in mediaservers item for item in mediaservers
@@ -335,7 +335,7 @@ class MediaServerChain(ChainBase):
server_name: str, server_name: str,
selected_libraries: List[Any], selected_libraries: List[Any],
library_media_counts: Dict[str, Optional[int]], library_media_counts: Dict[str, Optional[int]],
dboper: MediaServerOper, dboper: Any,
sync_time: str, sync_time: str,
progress_callback: Optional[Callable[..., None]], progress_callback: Optional[Callable[..., None]],
server_index: int, server_index: int,
@@ -466,7 +466,7 @@ class MediaServerChain(ChainBase):
with lock: with lock:
# 汇总统计 # 汇总统计
total_count = 0 total_count = 0
dboper = MediaServerOper() dboper = get_chain_media_server_port()
mediaservers, total_servers, server_sync_contexts, global_media_total = ( mediaservers, total_servers, server_sync_contexts, global_media_total = (
self._prepare_sync_contexts(mediaservers, server) self._prepare_sync_contexts(mediaservers, server)
) )
@@ -6,7 +6,7 @@
> 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本
> 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文
> 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` > 相关文档:`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 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口。 > 实施进度:阶段 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 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口
## 当前复核结论(2026-08-24 ## 当前复核结论(2026-08-24
@@ -178,6 +178,15 @@
- 兼容边界不变:数据库 `SiteOper``SitePortProxy``SiteChain/TorrentsChain` 公开方法、站点事件、 - 兼容边界不变:数据库 `SiteOper``SitePortProxy``SiteChain/TorrentsChain` 公开方法、站点事件、
CookieCloud/RSS 语义和插件调用方式均未改动。 CookieCloud/RSS 语义和插件调用方式均未改动。
### 长期整改阶段 16:媒体服务器 Chain 数据端口收口(2026-08-24
- `MediaServerChain` 原先通过 `MediaServerPortProxy as MediaServerOper` 清理已移除服务器、写入媒体项并
清理陈旧记录;现在统一调用 `get_chain_media_server_port()`,同步阶段继续共享同一个端口实例。
- 五个增量同步测试接缝改为替换命名 getter,架构门禁禁止媒体服务器 Chain 重新导入
`MediaServerPortProxy`
- 兼容边界不变:数据库 `MediaServerOper`、Proxy 类、媒体库同步公开方法、进度合同、停止语义和插件调用
方式均未修改。
### 总体判断 ### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
+5 -5
View File
@@ -131,11 +131,11 @@ Session. `app/db/adapters/` is the concrete persistence-adapter layer: it may
depend on Application-owned Protocols, UoW/Session and Oper implementations. depend on Application-owned Protocols, UoW/Session and Oper implementations.
This deliberate dependency inversion is the only `DB implementation -> This deliberate dependency inversion is the only `DB implementation ->
Application contract` direction; Application must remain free of DB imports. Application contract` direction; Application must remain free of DB imports.
Migrated workflow, user, interaction, messaging, music and site Chain consumers use Migrated workflow, user, interaction, messaging, music, site and media-server
the named `get_chain_*_port()` functions from `app/application/chain/data.py`; Chain consumers use the named `get_chain_*_port()` functions from
they must not alias migration-time `*PortProxy` classes back to database Oper `app/application/chain/data.py`; they must not alias migration-time `*PortProxy`
names. Those proxy classes remain compatibility boundaries while the other classes back to database Oper names. Those proxy classes remain compatibility
established Chain domains migrate independently. boundaries while the other established Chain domains migrate independently.
### Adapter boundaries ### Adapter boundaries
+15
View File
@@ -486,6 +486,21 @@ def test_site_chains_use_explicit_site_data_port_getter():
assert violations == [] assert violations == []
def test_mediaserver_chain_uses_explicit_data_port_getter():
"""媒体服务器链不得把 MediaServerPortProxy 伪装成 MediaServerOper。"""
path = APP_ROOT / "chain" / "mediaserver.py"
tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path))
violations = [
f"{path.relative_to(PROJECT_ROOT).as_posix()}:{node.lineno}"
for node in ast.walk(tree)
if isinstance(node, ast.ImportFrom)
and node.module == "app.application.chain.data"
and any(alias.name == "MediaServerPortProxy" for alias in node.names)
]
assert violations == []
def test_plugin_components_do_not_reexport_legacy_abi_names(): def test_plugin_components_do_not_reexport_legacy_abi_names():
"""新插件组件只提供 canonical 能力,不得复制旧 Helper、Manager 或 Oper 导出。""" """新插件组件只提供 canonical 能力,不得复制旧 Helper、Manager 或 Oper 导出。"""
violations: list[str] = [] violations: list[str] = []
+13 -5
View File
@@ -103,7 +103,7 @@ def test_sync_persists_music_without_querying_tv_episodes(database):
with database() as session: with database() as session:
with patch.object( with patch.object(
MEDIA_SERVER_CHAIN_MODULE, MEDIA_SERVER_CHAIN_MODULE,
"MediaServerOper", "get_chain_media_server_port",
lambda: MediaServerOper(session), lambda: MediaServerOper(session),
), patch.object( ), patch.object(
MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper, MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper,
@@ -196,7 +196,7 @@ def test_sync_updates_rows_and_removes_stale_entries(database):
with database() as session: with database() as session:
with patch.object( with patch.object(
MEDIA_SERVER_CHAIN_MODULE, MEDIA_SERVER_CHAIN_MODULE,
"MediaServerOper", "get_chain_media_server_port",
lambda: MediaServerOper(session), lambda: MediaServerOper(session),
), patch.object( ), patch.object(
MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper, MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper,
@@ -280,7 +280,7 @@ def test_sync_queries_counts_before_items_and_reports_media_progress(database):
with database() as session: with database() as session:
with patch.object( with patch.object(
MEDIA_SERVER_CHAIN_MODULE, MEDIA_SERVER_CHAIN_MODULE,
"MediaServerOper", "get_chain_media_server_port",
lambda: MediaServerOper(session), lambda: MediaServerOper(session),
), patch.object( ), patch.object(
MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper, MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper,
@@ -334,7 +334,11 @@ def test_sync_targets_one_server_without_excluding_other_enabled_servers(monkeyp
excluded_server_calls.append(servers) excluded_server_calls.append(servers)
chain.librarys = lambda server: library_calls.append(server) or [] chain.librarys = lambda server: library_calls.append(server) or []
monkeypatch.setattr(MEDIA_SERVER_CHAIN_MODULE, "MediaServerOper", FakeMediaServerOper) monkeypatch.setattr(
MEDIA_SERVER_CHAIN_MODULE,
"get_chain_media_server_port",
FakeMediaServerOper,
)
monkeypatch.setattr( monkeypatch.setattr(
MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper, MEDIA_SERVER_CHAIN_MODULE.ServiceConfigHelper,
"get_mediaserver_configs", "get_mediaserver_configs",
@@ -367,7 +371,11 @@ def test_sync_stops_without_emitting_completion_after_stop_signal(monkeypatch):
global_vars.stop_system() global_vars.stop_system()
return 0, 0 return 0, 0
monkeypatch.setattr(MEDIA_SERVER_CHAIN_MODULE, "MediaServerOper", FakeMediaServerOper) monkeypatch.setattr(
MEDIA_SERVER_CHAIN_MODULE,
"get_chain_media_server_port",
FakeMediaServerOper,
)
monkeypatch.setattr( monkeypatch.setattr(
chain, chain,
"_prepare_sync_contexts", "_prepare_sync_contexts",