mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-09-06 07:56:52 +08:00
refactor: split api dependencies by domain
This commit is contained in:
@@ -0,0 +1 @@
|
||||
"""按业务领域拆分的 FastAPI 依赖工厂。"""
|
||||
@@ -0,0 +1,29 @@
|
||||
"""Agent 与消息查询依赖。"""
|
||||
|
||||
from fastapi import Depends
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.api.context import get_agent_chat_repository, get_agent_chat_transaction
|
||||
from app.api.data import get_async_db
|
||||
from app.api.dependencies.data import repository
|
||||
from app.application.messaging.chat import (
|
||||
AgentChatService,
|
||||
AsyncAgentChatRepository,
|
||||
AsyncUnitOfWork,
|
||||
)
|
||||
from app.application.messaging.message import MessageQueryService
|
||||
|
||||
|
||||
def get_agent_chat_service(
|
||||
chat_repository: AsyncAgentChatRepository = Depends(get_agent_chat_repository),
|
||||
unit_of_work: AsyncUnitOfWork = Depends(get_agent_chat_transaction),
|
||||
) -> AgentChatService:
|
||||
"""组装类型化 Agent 会话历史查询和删除服务。"""
|
||||
return AgentChatService(chat_repository, unit_of_work)
|
||||
|
||||
|
||||
def get_message_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> MessageQueryService:
|
||||
"""组装消息历史异步查询服务。"""
|
||||
return MessageQueryService(repository=repository("message", db))
|
||||
@@ -0,0 +1,116 @@
|
||||
"""用户身份、授权与认证服务依赖。"""
|
||||
|
||||
from typing import Any
|
||||
|
||||
from fastapi import Depends, HTTPException
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.adapters.web.security.access import verify_token
|
||||
from app.api.data import get_async_db, get_db
|
||||
from app.api.dependencies.data import repository, standalone_repository
|
||||
from app.application.security.auth import AuthService
|
||||
from app.application.security.passkeys import PasskeyService
|
||||
from app.application.security.user import UserService
|
||||
from app.schemas.token import TokenPayload as _SchemaTokenPayload
|
||||
|
||||
|
||||
def get_user_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> UserService:
|
||||
"""组装用户管理应用服务。"""
|
||||
return UserService(repository=repository("user", db))
|
||||
|
||||
|
||||
def get_auth_service() -> AuthService:
|
||||
"""组装同步认证应用服务。"""
|
||||
return AuthService(
|
||||
users=standalone_repository("user"),
|
||||
config=standalone_repository("system_config"),
|
||||
passkeys=standalone_repository("passkey"),
|
||||
)
|
||||
|
||||
|
||||
def get_passkey_service() -> PasskeyService:
|
||||
"""组装 PassKey 应用服务。"""
|
||||
return PasskeyService(repository=standalone_repository("passkey"))
|
||||
|
||||
|
||||
def get_current_user(
|
||||
db: Session = Depends(get_db),
|
||||
token_data: _SchemaTokenPayload = Depends(verify_token),
|
||||
) -> Any:
|
||||
"""读取令牌对应用户,不存在时返回 403。"""
|
||||
user = repository("user", db).get_by_id(token_data.sub)
|
||||
if not user:
|
||||
raise HTTPException(status_code=403, detail="用户不存在")
|
||||
return user
|
||||
|
||||
|
||||
async def get_current_user_async(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
token_data: _SchemaTokenPayload = Depends(verify_token),
|
||||
) -> Any:
|
||||
"""异步读取令牌对应用户,不存在时返回 403。"""
|
||||
user = await repository("user", db).async_get_by_id(token_data.sub)
|
||||
if not user:
|
||||
raise HTTPException(status_code=403, detail="用户不存在")
|
||||
return user
|
||||
|
||||
|
||||
def get_current_active_user(
|
||||
current_user: Any = Depends(get_current_user),
|
||||
) -> Any:
|
||||
"""校验并返回当前激活用户。"""
|
||||
if not current_user.is_active:
|
||||
raise HTTPException(status_code=403, detail="用户未激活")
|
||||
return current_user
|
||||
|
||||
|
||||
async def get_current_active_user_async(
|
||||
current_user: Any = Depends(get_current_user_async),
|
||||
) -> Any:
|
||||
"""异步校验并返回当前激活用户。"""
|
||||
if not current_user.is_active:
|
||||
raise HTTPException(status_code=403, detail="用户未激活")
|
||||
return current_user
|
||||
|
||||
|
||||
def _ensure_manage_user(current_user: Any) -> Any:
|
||||
"""校验用户具备全局管理权限。"""
|
||||
permissions = current_user.permissions or {}
|
||||
if not current_user.is_superuser and not bool(permissions.get("manage")):
|
||||
raise HTTPException(status_code=400, detail="用户权限不足")
|
||||
return current_user
|
||||
|
||||
|
||||
def get_current_active_manage_user(
|
||||
current_user: Any = Depends(get_current_active_user),
|
||||
) -> Any:
|
||||
"""返回当前拥有管理权限的激活用户。"""
|
||||
return _ensure_manage_user(current_user)
|
||||
|
||||
|
||||
async def get_current_active_manage_user_async(
|
||||
current_user: Any = Depends(get_current_active_user_async),
|
||||
) -> Any:
|
||||
"""异步返回当前拥有管理权限的激活用户。"""
|
||||
return _ensure_manage_user(current_user)
|
||||
|
||||
|
||||
def get_current_active_superuser(
|
||||
current_user: Any = Depends(get_current_user),
|
||||
) -> Any:
|
||||
"""校验并返回当前激活超级管理员。"""
|
||||
if not current_user.is_superuser:
|
||||
raise HTTPException(status_code=400, detail="用户权限不足")
|
||||
return current_user
|
||||
|
||||
|
||||
async def get_current_active_superuser_async(
|
||||
current_user: Any = Depends(get_current_user_async),
|
||||
) -> Any:
|
||||
"""异步校验并返回当前激活超级管理员。"""
|
||||
if not current_user.is_superuser:
|
||||
raise HTTPException(status_code=400, detail="用户权限不足")
|
||||
return current_user
|
||||
@@ -0,0 +1,20 @@
|
||||
"""未迁移领域共用的 API 数据兼容 Facade。"""
|
||||
|
||||
from typing import Any
|
||||
|
||||
from app.api.data import get_api_data_ports
|
||||
|
||||
|
||||
def repository(name: str, session: Any) -> Any:
|
||||
"""按旧能力名构造绑定当前请求会话的仓储。"""
|
||||
return get_api_data_ports().repository(name, session)
|
||||
|
||||
|
||||
def standalone_repository(name: str) -> Any:
|
||||
"""按旧能力名构造无需绑定请求会话的仓储。"""
|
||||
return get_api_data_ports().standalone_repository(name)
|
||||
|
||||
|
||||
def transaction(name: str, session: Any) -> Any:
|
||||
"""按旧能力名构造绑定当前请求会话的事务端口。"""
|
||||
return get_api_data_ports().transaction(name, session)
|
||||
@@ -0,0 +1,86 @@
|
||||
"""历史、媒体服务器与 Dashboard 查询依赖。"""
|
||||
|
||||
from fastapi import Depends
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.api.data import get_async_db, get_db
|
||||
from app.api.dependencies.data import repository, transaction
|
||||
from app.application.dashboard import DashboardQueryService
|
||||
from app.application.history import (
|
||||
DownloadHistoryMutationCommand,
|
||||
HistoryQueryService,
|
||||
TransferHistoryLookupService,
|
||||
TransferHistoryMutationCommand,
|
||||
clear_transfer_failures,
|
||||
)
|
||||
from app.application.mediaserver import MediaServerQueryService
|
||||
from app.chain.storage import StorageChain
|
||||
from app.runtime.events import eventmanager
|
||||
from app.schemas.types import EventType
|
||||
from app.schemas.workflow import FileItem as _SchemaFileItem
|
||||
|
||||
|
||||
def get_mediaserver_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> MediaServerQueryService:
|
||||
"""组装媒体服务器本地条目异步查询服务。"""
|
||||
return MediaServerQueryService(repository=repository("media_server", db))
|
||||
|
||||
|
||||
def get_dashboard_query_service(
|
||||
db: Session = Depends(get_db),
|
||||
) -> DashboardQueryService:
|
||||
"""组装 Dashboard 媒体与整理历史统计查询服务。"""
|
||||
from app.chain.dashboard import DashboardChain
|
||||
|
||||
return DashboardQueryService(
|
||||
repository=repository("transfer_history", db),
|
||||
media_statistics=DashboardChain().media_statistic,
|
||||
)
|
||||
|
||||
|
||||
def get_download_history_mutation_command(
|
||||
db: Session = Depends(get_db),
|
||||
) -> DownloadHistoryMutationCommand:
|
||||
"""组装下载历史删除用例及其请求级事务。"""
|
||||
return DownloadHistoryMutationCommand(
|
||||
repository=repository("download_history", db),
|
||||
unit_of_work=transaction("sync", db),
|
||||
)
|
||||
|
||||
|
||||
def get_history_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> HistoryQueryService:
|
||||
"""组装历史列表和详情异步查询服务。"""
|
||||
return HistoryQueryService(
|
||||
download_repository=repository("download_history", db),
|
||||
transfer_repository=repository("transfer_history", db),
|
||||
)
|
||||
|
||||
|
||||
def get_transfer_history_lookup_service(
|
||||
db: Session = Depends(get_db),
|
||||
) -> TransferHistoryLookupService:
|
||||
"""组装手动整理使用的同步历史投影服务。"""
|
||||
return TransferHistoryLookupService(repository("transfer_history", db))
|
||||
|
||||
|
||||
def get_transfer_history_mutation_command(
|
||||
db: Session = Depends(get_db),
|
||||
) -> TransferHistoryMutationCommand:
|
||||
"""组装整理历史删除、文件处理和事件发布用例。"""
|
||||
storage_chain = StorageChain()
|
||||
return TransferHistoryMutationCommand(
|
||||
repository=repository("transfer_history", db),
|
||||
download_repository=repository("download_history", db),
|
||||
unit_of_work=transaction("sync", db),
|
||||
file_item_factory=lambda payload: _SchemaFileItem(**payload),
|
||||
delete_media_file=storage_chain.delete_media_file,
|
||||
publish_download_file_deleted=lambda payload: eventmanager.send_event(
|
||||
EventType.DownloadFileDeleted,
|
||||
payload,
|
||||
),
|
||||
clear_failures=clear_transfer_failures,
|
||||
)
|
||||
@@ -0,0 +1,43 @@
|
||||
"""插件配置与运行态刷新依赖。"""
|
||||
|
||||
from app.application.commands import init_commands
|
||||
from app.application.plugin.config import PluginConfigCommand
|
||||
from app.application.plugin.routes import register_plugin_api
|
||||
from app.application.plugin.runtime import get_plugin_manager
|
||||
from app.application.scheduling import update_plugin_job
|
||||
from app.runtime.events import eventmanager
|
||||
from app.schemas.event import PluginDataResetEventData
|
||||
from app.schemas.types import ChainEventType
|
||||
|
||||
|
||||
def get_plugin_config_command() -> PluginConfigCommand:
|
||||
"""组装插件配置更新与重置用例,隔离 API 对运行时写操作的编排。"""
|
||||
manager = get_plugin_manager()
|
||||
|
||||
def publish_reset(plugin_id: str) -> None:
|
||||
"""在清理持久化数据前通知目标插件执行补偿。"""
|
||||
eventmanager.send_event(
|
||||
ChainEventType.PluginDataReset,
|
||||
PluginDataResetEventData(
|
||||
plugin_id=plugin_id,
|
||||
reset_config=True,
|
||||
reset_data=True,
|
||||
),
|
||||
)
|
||||
|
||||
def refresh_registrations(plugin_id: str) -> None:
|
||||
"""按服务、命令、动态路由顺序刷新插件宿主注册。"""
|
||||
update_plugin_job(plugin_id)
|
||||
init_commands(plugin_id)
|
||||
register_plugin_api(plugin_id)
|
||||
|
||||
return PluginConfigCommand(
|
||||
save_config=manager.save_plugin_config,
|
||||
initialize=manager.init_plugin,
|
||||
stop=manager.stop,
|
||||
delete_config=manager.delete_plugin_config,
|
||||
delete_data=manager.delete_plugin_data,
|
||||
reload_runtime=manager.reload_plugin,
|
||||
publish_reset=publish_reset,
|
||||
refresh_registrations=refresh_registrations,
|
||||
)
|
||||
@@ -0,0 +1,62 @@
|
||||
"""站点领域的请求级 command/query 依赖。"""
|
||||
|
||||
from fastapi import Depends
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.api.data import get_async_db, get_db
|
||||
from app.api.dependencies.data import repository, transaction
|
||||
from app.application.site.mutation import SiteMutationCommand
|
||||
from app.application.site.query import SiteQueryService
|
||||
from app.application.site.sites import SitesHelper # pylint: disable=no-name-in-module
|
||||
from app.domain import site as site_rules
|
||||
from app.foundation import url as url_tools
|
||||
from app.runtime.events import eventmanager
|
||||
from app.schemas.types import EventType
|
||||
|
||||
|
||||
async def _publish_site_updated(payload: dict) -> None:
|
||||
"""发布已提交的站点更新事件。"""
|
||||
await eventmanager.async_send_event(EventType.SiteUpdated, payload)
|
||||
|
||||
|
||||
async def _publish_site_deleted(payload: dict) -> None:
|
||||
"""发布已提交的站点删除事件。"""
|
||||
await eventmanager.async_send_event(EventType.SiteDeleted, payload)
|
||||
|
||||
|
||||
def get_site_mutation_command(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> SiteMutationCommand:
|
||||
"""组装请求级站点写用例及其事务和外部目录依赖。"""
|
||||
sites_helper = SitesHelper()
|
||||
|
||||
def normalize_url(value: str) -> str:
|
||||
"""沿用站点接口的 scheme/netloc 规范化格式。"""
|
||||
scheme, netloc = url_tools.split_netloc(value)
|
||||
return f"{scheme}://{netloc}/"
|
||||
|
||||
return SiteMutationCommand(
|
||||
repository=repository("site", db),
|
||||
unit_of_work=transaction("async", db),
|
||||
auth_level_provider=lambda: sites_helper.auth_level,
|
||||
indexer_loader=sites_helper.async_get_indexer,
|
||||
domain_extractor=site_rules.extract_domain,
|
||||
url_normalizer=normalize_url,
|
||||
publish_updated=_publish_site_updated,
|
||||
publish_deleted=_publish_site_deleted,
|
||||
)
|
||||
|
||||
|
||||
def get_site_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> SiteQueryService:
|
||||
"""组装站点异步查询服务。"""
|
||||
return SiteQueryService(repository=repository("site", db))
|
||||
|
||||
|
||||
def get_site_sync_query_service(
|
||||
db: Session = Depends(get_db),
|
||||
) -> SiteQueryService:
|
||||
"""组装站点同步查询服务,用于同步 Chain 路由。"""
|
||||
return SiteQueryService(repository=repository("site", db))
|
||||
@@ -0,0 +1,125 @@
|
||||
"""订阅领域的请求级 command/query 依赖。"""
|
||||
|
||||
from fastapi import BackgroundTasks, Depends
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.adapters.external.server import MoviePilotServerHelper
|
||||
from app.api.data import get_async_db, get_db
|
||||
from app.api.dependencies.data import repository, transaction
|
||||
from app.application.scheduling import Scheduler
|
||||
from app.application.servarr import ServarrSubscriptionService
|
||||
from app.application.subscription.delete import DeleteSubscribeCommand
|
||||
from app.application.subscription.identity import DeleteSubscriptionsByIdentityCommand
|
||||
from app.application.subscription.mutation import SubscriptionMutationService
|
||||
from app.application.subscription.query import SubscriptionQueryService
|
||||
from app.application.subscription.search import SearchSubscriptionsCommand
|
||||
from app.runtime.events import eventmanager
|
||||
from app.runtime.log import logger
|
||||
from app.schemas.types import EventType
|
||||
|
||||
|
||||
async def _publish_subscribe_deleted(
|
||||
subscribe_id: int,
|
||||
subscribe_info: dict,
|
||||
) -> None:
|
||||
"""通过宿主事件总线发布已提交的订阅删除事件。"""
|
||||
await eventmanager.async_send_event(
|
||||
EventType.SubscribeDeleted,
|
||||
{"subscribe_id": subscribe_id, "subscribe_info": subscribe_info},
|
||||
)
|
||||
|
||||
|
||||
def get_delete_subscribe_command(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> DeleteSubscribeCommand:
|
||||
"""组装请求级订阅删除用例及其具体适配器。"""
|
||||
return DeleteSubscribeCommand(
|
||||
repository=repository("subscribe", db),
|
||||
unit_of_work=transaction("async", db),
|
||||
publish_deleted=_publish_subscribe_deleted,
|
||||
report_deleted=MoviePilotServerHelper.sub_done_async,
|
||||
)
|
||||
|
||||
|
||||
def _log_subscribe_deleted_event_error(
|
||||
subscribe_id: int,
|
||||
error: Exception,
|
||||
) -> None:
|
||||
"""记录按媒体身份删除时的单条事件失败并允许后续事件继续。"""
|
||||
logger.error(
|
||||
f"发送订阅删除事件失败:{subscribe_id} - {error}",
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
|
||||
def get_delete_subscriptions_by_identity_command(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> DeleteSubscriptionsByIdentityCommand:
|
||||
"""组装请求级按媒体身份删除订阅用例。"""
|
||||
return DeleteSubscriptionsByIdentityCommand(
|
||||
repository=repository("subscribe", db),
|
||||
unit_of_work=transaction("async", db),
|
||||
publish_deleted=_publish_subscribe_deleted,
|
||||
handle_event_error=_log_subscribe_deleted_event_error,
|
||||
)
|
||||
|
||||
|
||||
def get_search_subscriptions_command(
|
||||
background_tasks: BackgroundTasks,
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> SearchSubscriptionsCommand:
|
||||
"""组装手工订阅搜索用例,并把调度延迟到响应后的后台任务。"""
|
||||
def schedule_search(subscribe_id: int | None, state: str | None) -> None:
|
||||
"""按历史参数提交订阅搜索调度任务。"""
|
||||
background_tasks.add_task(
|
||||
Scheduler().start,
|
||||
job_id="subscribe_search",
|
||||
sid=subscribe_id,
|
||||
state=state,
|
||||
manual=True,
|
||||
)
|
||||
|
||||
return SearchSubscriptionsCommand(
|
||||
repository=repository("subscribe", db),
|
||||
schedule_search=schedule_search,
|
||||
)
|
||||
|
||||
|
||||
def get_subscription_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> SubscriptionQueryService:
|
||||
"""组装订阅和订阅历史异步查询服务。"""
|
||||
return SubscriptionQueryService(
|
||||
repository=repository("subscribe", db),
|
||||
async_repository=repository("subscribe", db),
|
||||
history_repository=repository("subscribe_history", db),
|
||||
)
|
||||
|
||||
|
||||
def get_subscription_mutation_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> SubscriptionMutationService:
|
||||
"""组装异步订阅写服务。"""
|
||||
return SubscriptionMutationService(
|
||||
repository=repository("subscribe", db),
|
||||
history_repository=repository("subscribe_history", db),
|
||||
)
|
||||
|
||||
|
||||
def get_subscription_sync_mutation_service(
|
||||
db: Session = Depends(get_db),
|
||||
) -> SubscriptionMutationService:
|
||||
"""组装同步订阅查询服务,供文件信息接口使用。"""
|
||||
return SubscriptionMutationService(repository=repository("subscribe", db))
|
||||
|
||||
|
||||
def get_servarr_subscription_service(
|
||||
async_db: AsyncSession = Depends(get_async_db),
|
||||
db: Session = Depends(get_db),
|
||||
) -> ServarrSubscriptionService:
|
||||
"""组装 Servarr 兼容路由的请求级订阅数据用例。"""
|
||||
return ServarrSubscriptionService(
|
||||
async_repository=repository("subscribe", async_db),
|
||||
sync_repository=repository("subscribe", db),
|
||||
)
|
||||
@@ -0,0 +1,60 @@
|
||||
"""工作流领域的请求级 command/query 依赖。"""
|
||||
|
||||
from fastapi import Depends
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.adapters.external.server import MoviePilotServerHelper
|
||||
from app.api.data import get_async_db, get_db
|
||||
from app.api.dependencies.data import repository, standalone_repository, transaction
|
||||
from app.application.scheduling import Scheduler
|
||||
from app.application.workflow import (
|
||||
WorkflowDefinitionCommand,
|
||||
WorkflowMutationCommand,
|
||||
WorkflowQueryService,
|
||||
)
|
||||
from app.runtime.config import global_vars
|
||||
from app.workflow import WorkFlowManager
|
||||
|
||||
|
||||
def get_workflow_mutation_command(
|
||||
db: Session = Depends(get_db),
|
||||
) -> WorkflowMutationCommand:
|
||||
"""组装请求级工作流写用例和提交后的调度副作用。"""
|
||||
scheduler = Scheduler()
|
||||
workflow_manager = WorkFlowManager()
|
||||
return WorkflowMutationCommand(
|
||||
repository=repository("workflow", db),
|
||||
unit_of_work=transaction("sync", db),
|
||||
add_timer=scheduler.update_workflow_job,
|
||||
remove_timer=scheduler.remove_workflow_job,
|
||||
load_event=workflow_manager.load_workflow_events,
|
||||
remove_event=workflow_manager.remove_workflow_event,
|
||||
refresh_event=workflow_manager.update_workflow_event,
|
||||
stop_running=global_vars.stop_workflow,
|
||||
delete_cache=lambda workflow_id: standalone_repository(
|
||||
"system_config"
|
||||
).delete(f"WorkflowCache-{workflow_id}"),
|
||||
)
|
||||
|
||||
|
||||
def get_workflow_definition_command(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> WorkflowDefinitionCommand:
|
||||
"""组装工作流创建、复用和重置的异步写用例。"""
|
||||
return WorkflowDefinitionCommand(
|
||||
repository=repository("workflow", db),
|
||||
unit_of_work=transaction("async", db),
|
||||
stop_running=global_vars.stop_workflow,
|
||||
delete_cache=lambda workflow_id: standalone_repository(
|
||||
"system_config"
|
||||
).delete(f"WorkflowCache-{workflow_id}"),
|
||||
report_fork=MoviePilotServerHelper.async_workflow_fork_by_id,
|
||||
)
|
||||
|
||||
|
||||
def get_workflow_query_service(
|
||||
db: AsyncSession = Depends(get_async_db),
|
||||
) -> WorkflowQueryService:
|
||||
"""组装工作流只读查询用例,避免端点直接持有数据库操作器。"""
|
||||
return WorkflowQueryService(repository=repository("workflow", db))
|
||||
Reference in New Issue
Block a user