mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-09-04 23:17:20 +08:00
refactor(transfer): make queue admission durable
This commit is contained in:
@@ -10,8 +10,11 @@ from collections.abc import Callable
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Optional
|
||||
|
||||
from app.application.transfer import TransferAdmissionRepository
|
||||
|
||||
|
||||
OperFactory = Callable[[], Any]
|
||||
TransferAdmissionRepositoryFactory = Callable[[], TransferAdmissionRepository]
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
@@ -23,7 +26,7 @@ class ChainDataPorts:
|
||||
workflow: OperFactory
|
||||
download_history: OperFactory
|
||||
transfer_history: OperFactory
|
||||
transfer_pending: OperFactory
|
||||
transfer_pending: TransferAdmissionRepositoryFactory
|
||||
media_server: OperFactory
|
||||
download_failure: OperFactory
|
||||
user: OperFactory
|
||||
@@ -77,12 +80,6 @@ class TransferHistoryPortProxy(_ChainDataPortProxy):
|
||||
port_name = "transfer_history"
|
||||
|
||||
|
||||
class TransferPendingPortProxy(_ChainDataPortProxy):
|
||||
"""待整理数据端口代理。"""
|
||||
|
||||
port_name = "transfer_pending"
|
||||
|
||||
|
||||
class MediaServerPortProxy(_ChainDataPortProxy):
|
||||
"""媒体服务器数据端口代理。"""
|
||||
|
||||
@@ -156,8 +153,8 @@ def get_chain_transfer_history_port() -> Any:
|
||||
return get_chain_data_ports().transfer_history()
|
||||
|
||||
|
||||
def get_chain_transfer_pending_port() -> Any:
|
||||
"""创建待整理数据端口实例。"""
|
||||
def get_chain_transfer_pending_port() -> TransferAdmissionRepository:
|
||||
"""创建类型化的整理任务 durable admission 仓储。"""
|
||||
return get_chain_data_ports().transfer_pending()
|
||||
|
||||
|
||||
|
||||
+72
-13
@@ -19,14 +19,10 @@ from copy import deepcopy
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from time import monotonic
|
||||
from typing import Callable, Dict, List, Optional, Tuple, Union
|
||||
from typing import Callable, Dict, List, Optional, Protocol, Tuple, Union
|
||||
|
||||
from pydantic import BaseModel, ConfigDict
|
||||
from pydantic import BaseModel, ConfigDict, PrivateAttr
|
||||
|
||||
from app.schemas.transfer import MetaInfo as _SchemaMetaInfo
|
||||
from app.schemas.transfer import MusicInfo as _SchemaMusicInfo
|
||||
from app.schemas.transfer import MusicMeta as _SchemaMusicMeta
|
||||
from app.schemas.workflow import MediaInfo as _SchemaMediaInfo
|
||||
from app.adapters.system.host import SystemUtils
|
||||
from app.application.agent import get_prompt_manager, get_running_agent_manager
|
||||
from app.domain.context import MediaInfo, MusicInfo
|
||||
@@ -40,6 +36,9 @@ from app.schemas.history import DownloadHistory
|
||||
from app.schemas.media import OptionalMediaIdentityMixin, resolve_media_identity
|
||||
from app.schemas.system import TransferDirectoryConf
|
||||
from app.schemas.tmdb import TmdbEpisode
|
||||
from app.schemas.transfer import MetaInfo as _SchemaMetaInfo
|
||||
from app.schemas.transfer import MusicInfo as _SchemaMusicInfo
|
||||
from app.schemas.transfer import MusicMeta as _SchemaMusicMeta
|
||||
from app.schemas.transfer import TransferInfo, TransferJob, TransferJobTask
|
||||
from app.schemas.types import (
|
||||
MUSIC_ENTITY_ALBUM,
|
||||
@@ -48,7 +47,7 @@ from app.schemas.types import (
|
||||
MediaType,
|
||||
ReplyMode,
|
||||
)
|
||||
|
||||
from app.schemas.workflow import MediaInfo as _SchemaMediaInfo
|
||||
|
||||
|
||||
class TransferTask(OptionalMediaIdentityMixin, BaseModel):
|
||||
@@ -81,6 +80,16 @@ class TransferTask(OptionalMediaIdentityMixin, BaseModel):
|
||||
manual: Optional[bool] = False
|
||||
background: Optional[bool] = True
|
||||
preview: Optional[bool] = False
|
||||
_admission_task_id: Optional[str] = PrivateAttr(default=None)
|
||||
|
||||
@property
|
||||
def admission_task_id(self) -> Optional[str]:
|
||||
"""返回仅供宿主持久准入和终态结算使用的内部任务标识。"""
|
||||
return self._admission_task_id
|
||||
|
||||
def bind_admission_task_id(self, task_id: str) -> None:
|
||||
"""绑定持久准入生成的稳定身份,不改变插件可见序列化字段。"""
|
||||
self._admission_task_id = task_id
|
||||
|
||||
def to_dict(self):
|
||||
"""
|
||||
@@ -113,6 +122,42 @@ class TransferQueue(BaseModel):
|
||||
result: Optional[TransferInfo] = None
|
||||
|
||||
|
||||
TRANSFER_ADMISSION_ACCEPTED = "accepted"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TransferAdmission:
|
||||
"""描述已经持久化、可在进程退出后恢复的整理任务准入事实。"""
|
||||
|
||||
task_id: str
|
||||
storage: str
|
||||
src_path: str
|
||||
state: str
|
||||
created_at: str
|
||||
updated_at: str
|
||||
last_error: Optional[str] = None
|
||||
|
||||
|
||||
class TransferAdmissionRepository(Protocol):
|
||||
"""整理任务 durable admission 所需的类型化持久化端口。"""
|
||||
|
||||
def admit(self, *, storage: str, src_path: str) -> TransferAdmission:
|
||||
"""幂等登记源文件并返回稳定任务身份。"""
|
||||
...
|
||||
|
||||
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
|
||||
"""按登记顺序返回等待恢复或执行的任务。"""
|
||||
...
|
||||
|
||||
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
|
||||
"""记录内存队列接收失败,保留任务供后续恢复。"""
|
||||
...
|
||||
|
||||
def discard_task(self, *, task_id: str) -> int:
|
||||
"""按稳定任务身份删除已经到达终态的登记。"""
|
||||
...
|
||||
|
||||
|
||||
class TransferQueueService:
|
||||
"""协调整理任务登记、入队、移除和队列视图查询。"""
|
||||
|
||||
@@ -120,29 +165,43 @@ class TransferQueueService:
|
||||
self,
|
||||
*,
|
||||
register_task: Callable[[TransferTask], bool],
|
||||
admit_task: Callable[[TransferTask], TransferAdmission],
|
||||
enqueue: Callable[[TransferQueue], None],
|
||||
before_enqueue: Callable[[TransferTask], None],
|
||||
after_enqueue: Callable[[TransferTask], None],
|
||||
enqueue_failed: Callable[[TransferTask, Exception], None],
|
||||
remove_task: Callable[[FileItem], None],
|
||||
list_tasks: Callable[[], List[TransferJob]],
|
||||
expire_tasks: Callable[[], None],
|
||||
) -> None:
|
||||
"""保存队列用例依赖,避免 Application 服务绑定具体线程队列实现。"""
|
||||
self._register_task = register_task
|
||||
self._admit_task = admit_task
|
||||
self._enqueue = enqueue
|
||||
self._before_enqueue = before_enqueue
|
||||
self._after_enqueue = after_enqueue
|
||||
self._enqueue_failed = enqueue_failed
|
||||
self._remove_task = remove_task
|
||||
self._list_tasks = list_tasks
|
||||
self._expire_tasks = expire_tasks
|
||||
|
||||
def put(self, task: TransferTask, callback: Callable) -> bool:
|
||||
"""登记并入队一个整理任务,保持原有副作用顺序。"""
|
||||
"""先持久化准入事实再入队;任何前置失败都撤销内存作业视图。"""
|
||||
if not task or not self._register_task(task):
|
||||
return False
|
||||
self._before_enqueue(task)
|
||||
self._enqueue(TransferQueue(task=task, callback=callback))
|
||||
self._after_enqueue(task)
|
||||
try:
|
||||
admission = self._admit_task(task)
|
||||
task.bind_admission_task_id(admission.task_id)
|
||||
except Exception:
|
||||
self._remove_task(task.fileitem)
|
||||
raise
|
||||
try:
|
||||
self._before_enqueue(task)
|
||||
self._enqueue(TransferQueue(task=task, callback=callback))
|
||||
except Exception as err:
|
||||
try:
|
||||
self._enqueue_failed(task, err)
|
||||
finally:
|
||||
self._remove_task(task.fileitem)
|
||||
raise
|
||||
return True
|
||||
|
||||
def remove(self, fileitem: FileItem) -> None:
|
||||
|
||||
+58
-23
@@ -52,6 +52,7 @@ from app.application.outbox import (
|
||||
from app.application.transfer import (
|
||||
FailedRetryScheduler,
|
||||
JobManager,
|
||||
TransferAdmission,
|
||||
TransferFailureNotification,
|
||||
TransferFailureNotificationAggregator,
|
||||
TransferQueue,
|
||||
@@ -225,8 +226,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
self.retry_scheduler = FailedRetryScheduler()
|
||||
# 整理失败通知聚合器
|
||||
self.failure_notification_aggregator = TransferFailureNotificationAggregator()
|
||||
# 待整理文件落盘登记,用于进程重启后回放内存队列里未完成的任务
|
||||
self._pendingoper = get_chain_transfer_pending_port()
|
||||
# durable admission 仓储先于内存入队保存任务,进程退出后仍可恢复。
|
||||
self._transfer_admissions = get_chain_transfer_pending_port()
|
||||
# 转移成功的文件清单
|
||||
self._success_target_files: Dict[Tuple, List[str]] = {}
|
||||
# 批次级刮削缓冲,避免同一批多文件入库重复触发目录刮削
|
||||
@@ -867,7 +868,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
"""
|
||||
添加到待整理队列
|
||||
:param task: 任务信息
|
||||
:return: True表示任务已添加到队列,False表示任务无效或已存在(重复)
|
||||
:return: True表示任务已添加,False表示链已关闭或任务无效/重复
|
||||
:raises Exception: 持久准入、批次登记或内存入队失败
|
||||
"""
|
||||
with self._worker_state_lock:
|
||||
if self._closing:
|
||||
@@ -879,9 +881,10 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
"""构建保持旧队列对象和私有兼容接缝的应用服务。"""
|
||||
return TransferQueueService(
|
||||
register_task=self.__put_to_jobview,
|
||||
admit_task=self.__admit_transfer,
|
||||
enqueue=self._queue.put,
|
||||
before_enqueue=self._register_scrape_batch_task,
|
||||
after_enqueue=self.__register_pending,
|
||||
enqueue_failed=self.__record_enqueue_failure,
|
||||
remove_task=self.jobview.remove_task,
|
||||
list_tasks=self.jobview.list_jobs,
|
||||
expire_tasks=self.__expire_stale_transfer_tasks,
|
||||
@@ -936,7 +939,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
if stop_event.is_set():
|
||||
return
|
||||
try:
|
||||
pendings = self._pendingoper.list_all()
|
||||
pendings = self._transfer_admissions.list_accepted()
|
||||
except Exception as err:
|
||||
logger.error(f"读取待整理文件登记失败:{err}")
|
||||
return
|
||||
@@ -944,9 +947,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
return
|
||||
logger.info(f"发现 {len(pendings)} 个上次未整理完的文件,正在重新送入整理链 ...")
|
||||
replayed = 0
|
||||
for storage, src_path in pendings:
|
||||
for admission in pendings:
|
||||
if stop_event.is_set():
|
||||
break
|
||||
storage = admission.storage
|
||||
src_path = admission.src_path
|
||||
try:
|
||||
fileitem, should_discard = self.__build_replay_fileitem(storage, src_path)
|
||||
# stat 等同步 I/O 返回后重新检查,关闭期间不得注销尚未完成的登记。
|
||||
@@ -955,7 +960,9 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
if not fileitem:
|
||||
if should_discard:
|
||||
# 源文件确认已消失,注销登记避免每次启动重复回放
|
||||
self._pendingoper.discard(storage=storage, src_path=src_path)
|
||||
self._transfer_admissions.discard_task(
|
||||
task_id=admission.task_id
|
||||
)
|
||||
continue
|
||||
self.do_transfer(fileitem=fileitem)
|
||||
replayed += 1
|
||||
@@ -1009,19 +1016,35 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
modify_time=modify_time,
|
||||
), False
|
||||
|
||||
def __register_pending(self, task: TransferTask):
|
||||
"""
|
||||
落盘登记一个待整理文件,登记失败不影响正常入队。
|
||||
:param task: 任务信息
|
||||
"""
|
||||
def __admit_transfer(self, task: TransferTask) -> TransferAdmission:
|
||||
"""在内存入队前持久化源文件并返回稳定任务身份。"""
|
||||
fileitem = task.fileitem if task else None
|
||||
if not fileitem or not fileitem.path:
|
||||
return
|
||||
if not fileitem or not fileitem.storage or not fileitem.path:
|
||||
raise ValueError("整理任务缺少源文件身份")
|
||||
return self._transfer_admissions.admit(
|
||||
storage=fileitem.storage,
|
||||
src_path=fileitem.path,
|
||||
)
|
||||
|
||||
def __record_enqueue_failure(
|
||||
self,
|
||||
task: TransferTask,
|
||||
error: Exception,
|
||||
) -> None:
|
||||
"""记录内存入队失败并撤销该任务的批次占位。"""
|
||||
try:
|
||||
self._pendingoper.register(storage=fileitem.storage, src_path=fileitem.path)
|
||||
except Exception as err:
|
||||
# 登记只是重启后的补救手段,失败不能阻断正常整理
|
||||
logger.debug(f"登记待整理文件失败: {fileitem.path} - {err}")
|
||||
if task.admission_task_id:
|
||||
self._transfer_admissions.record_enqueue_failure(
|
||||
task_id=task.admission_task_id,
|
||||
error=str(error),
|
||||
)
|
||||
except Exception as record_error:
|
||||
logger.error(
|
||||
"记录整理任务入队失败原因异常:"
|
||||
f"{task.admission_task_id} - {record_error}"
|
||||
)
|
||||
finally:
|
||||
self._finish_scrape_batch_task(task)
|
||||
|
||||
def __discard_pending(self, task: TransferTask):
|
||||
"""
|
||||
@@ -1031,13 +1054,17 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
重试预算重新送入整理链;留在本表反而会每次重启都重复回放。
|
||||
:param task: 任务信息
|
||||
"""
|
||||
fileitem = task.fileitem if task else None
|
||||
if not fileitem or not fileitem.path:
|
||||
if not task or not task.admission_task_id:
|
||||
return
|
||||
try:
|
||||
self._pendingoper.discard(storage=fileitem.storage, src_path=fileitem.path)
|
||||
self._transfer_admissions.discard_task(
|
||||
task_id=task.admission_task_id
|
||||
)
|
||||
except Exception as err:
|
||||
logger.debug(f"注销待整理文件登记失败: {fileitem.path} - {err}")
|
||||
logger.error(
|
||||
"注销整理任务 durable admission 失败: "
|
||||
f"{task.admission_task_id} - {err}"
|
||||
)
|
||||
|
||||
def __put_to_jobview(self, task: TransferTask) -> bool:
|
||||
"""
|
||||
@@ -2751,7 +2778,15 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
|
||||
preview=preview,
|
||||
)
|
||||
if background:
|
||||
if self.put_to_queue(task=transfer_task):
|
||||
try:
|
||||
queued = self.put_to_queue(task=transfer_task)
|
||||
except Exception as err:
|
||||
all_success = False
|
||||
message = f"{file_path.name} 加入整理队列失败:{err}"
|
||||
err_msgs.append(message)
|
||||
logger.error(message)
|
||||
continue
|
||||
if queued:
|
||||
logger.info(f"{file_path.name} 已添加到整理队列")
|
||||
else:
|
||||
logger.debug(f"{file_path.name} 已在整理队列中,跳过")
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
"""整理任务持久准入端口的 SQLAlchemy 适配器。"""
|
||||
|
||||
from collections.abc import Callable
|
||||
from datetime import datetime
|
||||
from uuid import uuid4
|
||||
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.application.transfer import (
|
||||
TRANSFER_ADMISSION_ACCEPTED,
|
||||
TransferAdmission,
|
||||
)
|
||||
from app.db.models.transferpending import TransferPending
|
||||
from app.db.oper.transferpending import TransferPendingOper
|
||||
from app.db.uow import SqlAlchemyUnitOfWork
|
||||
|
||||
|
||||
class TransactionalTransferAdmissionRepository:
|
||||
"""以短生命周期 Session 实现整理任务持久准入端口。"""
|
||||
|
||||
def __init__(self, session_factory: Callable[[], Session]) -> None:
|
||||
"""保存由组合根提供的同步会话工厂。"""
|
||||
self._session_factory = session_factory
|
||||
|
||||
@staticmethod
|
||||
def _now() -> str:
|
||||
"""生成与历史登记时间可按字典序比较的当前时间。"""
|
||||
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
@staticmethod
|
||||
def _project(pending: TransferPending) -> TransferAdmission:
|
||||
"""在 Session 有效期内把 ORM 行冻结为应用层 DTO。"""
|
||||
created_at = pending.created_at or pending.updated_at
|
||||
return TransferAdmission(
|
||||
task_id=pending.task_id,
|
||||
storage=pending.storage,
|
||||
src_path=pending.src_path,
|
||||
state=pending.state,
|
||||
created_at=created_at,
|
||||
updated_at=pending.updated_at,
|
||||
last_error=pending.last_error,
|
||||
)
|
||||
|
||||
def admit(self, *, storage: str, src_path: str) -> TransferAdmission:
|
||||
"""幂等持久化准入事实,并返回跨重启稳定的任务标识。"""
|
||||
now_time = self._now()
|
||||
try:
|
||||
with self._session_factory() as session:
|
||||
transaction = SqlAlchemyUnitOfWork(session)
|
||||
try:
|
||||
pending = TransferPendingOper(db=session).stage_admit(
|
||||
task_id=uuid4().hex,
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
state=TRANSFER_ADMISSION_ACCEPTED,
|
||||
now_time=now_time,
|
||||
)
|
||||
if pending is None:
|
||||
raise ValueError("整理任务的存储与源路径不能为空")
|
||||
session.flush()
|
||||
admission = self._project(pending)
|
||||
transaction.commit()
|
||||
return admission
|
||||
except Exception:
|
||||
transaction.rollback()
|
||||
raise
|
||||
except IntegrityError as error:
|
||||
# 并发准入可能同时通过查询;唯一约束决定赢家,输家回读稳定身份。
|
||||
with self._session_factory() as session:
|
||||
pending = TransferPendingOper(db=session).get_by_identity(
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
)
|
||||
if pending is None:
|
||||
raise RuntimeError("并发准入冲突后未找到已提交记录") from error
|
||||
return self._project(pending)
|
||||
|
||||
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
|
||||
"""在独立只读会话中投影等待恢复或执行的准入记录。"""
|
||||
with self._session_factory() as session:
|
||||
pending_items = TransferPendingOper(db=session).list_by_state(
|
||||
state=TRANSFER_ADMISSION_ACCEPTED,
|
||||
limit=limit,
|
||||
)
|
||||
return [self._project(pending) for pending in pending_items]
|
||||
|
||||
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
|
||||
"""独立提交最近一次入队失败,保留准入记录供后续恢复。"""
|
||||
with self._session_factory() as session:
|
||||
transaction = SqlAlchemyUnitOfWork(session)
|
||||
try:
|
||||
TransferPendingOper(db=session).stage_record_enqueue_failure(
|
||||
task_id=task_id,
|
||||
error=error,
|
||||
now_time=self._now(),
|
||||
)
|
||||
transaction.commit()
|
||||
except Exception:
|
||||
transaction.rollback()
|
||||
raise
|
||||
|
||||
def discard_task(self, *, task_id: str) -> int:
|
||||
"""在独立事务中按稳定任务标识删除已到终态的准入记录。"""
|
||||
with self._session_factory() as session:
|
||||
transaction = SqlAlchemyUnitOfWork(session)
|
||||
try:
|
||||
deleted = TransferPendingOper(db=session).stage_discard_task(
|
||||
task_id=task_id,
|
||||
)
|
||||
transaction.commit()
|
||||
return deleted
|
||||
except Exception:
|
||||
transaction.rollback()
|
||||
raise
|
||||
@@ -1,6 +1,8 @@
|
||||
from typing import List, Optional
|
||||
from datetime import datetime
|
||||
from typing import List, Optional, cast
|
||||
from uuid import uuid4
|
||||
|
||||
from sqlalchemy import Index, String, delete, select
|
||||
from sqlalchemy import Index, String, Text, UniqueConstraint, delete, select, update
|
||||
from sqlalchemy.orm import Mapped, Session, mapped_column
|
||||
|
||||
from app.db.base import Base, execute_dml, get_id_column
|
||||
@@ -15,22 +17,43 @@ class TransferPending(Base):
|
||||
蒸发。而已经稳定落地的文件不会再产生任何监控事件,也不会有新的补偿扫描起点
|
||||
——结果就是永久漏件,只能靠人工比对补整理。
|
||||
|
||||
这里只落盘最小事实:存储与源文件路径。重启后重新走一遍整理入口,由整理历史
|
||||
查重挡掉已经完成的,因此不需要序列化 meta/mediainfo 这些重对象,也不存在
|
||||
识别结果陈旧的问题。
|
||||
这里只落盘恢复所需的最小事实:稳定任务身份、存储、源文件路径、准入状态和
|
||||
最近入队错误。重启后重新走一遍整理入口,由整理历史查重挡掉已经完成的,
|
||||
因此不需要序列化 meta/mediainfo 这些重对象,也不存在识别结果陈旧的问题。
|
||||
"""
|
||||
|
||||
id = get_id_column()
|
||||
# 稳定任务标识
|
||||
task_id: Mapped[str] = mapped_column(
|
||||
String(64), nullable=False, default=lambda: uuid4().hex
|
||||
)
|
||||
# 存储
|
||||
storage: Mapped[str] = mapped_column(String, nullable=False)
|
||||
# 源文件路径
|
||||
src_path: Mapped[str] = mapped_column(String, nullable=False)
|
||||
# 登记时间
|
||||
created_at: Mapped[Optional[str]] = mapped_column(String)
|
||||
# 持久状态
|
||||
state: Mapped[str] = mapped_column(String(32), nullable=False, default="accepted")
|
||||
# 最后更新时间
|
||||
updated_at: Mapped[str] = mapped_column(
|
||||
String(40), nullable=False,
|
||||
default=lambda: datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||||
)
|
||||
# 最近一次入队失败原因
|
||||
last_error: Mapped[Optional[str]] = mapped_column(Text)
|
||||
|
||||
__table_args__ = (
|
||||
# 同一个文件重复入队只保留一条,回放时不会重复送入整理链
|
||||
Index("ux_transferpending_storage_path", "storage", "src_path", unique=True),
|
||||
# 恢复主查询按状态过滤、登记时间与主键稳定排序
|
||||
Index(
|
||||
"ix_transferpending_state_created",
|
||||
"state",
|
||||
"created_at",
|
||||
"id",
|
||||
),
|
||||
UniqueConstraint("task_id", name="uq_transferpending_task_id"),
|
||||
)
|
||||
|
||||
@classmethod
|
||||
@@ -50,11 +73,128 @@ class TransferPending(Base):
|
||||
select(cls).where(cls.storage == storage, cls.src_path == src_path)
|
||||
).scalars().first()
|
||||
if pending:
|
||||
return pending
|
||||
pending = cls(storage=storage, src_path=src_path, created_at=now_time)
|
||||
return cast("TransferPending", pending)
|
||||
pending = cls(
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
state="accepted",
|
||||
created_at=now_time,
|
||||
updated_at=now_time,
|
||||
)
|
||||
db.add(pending)
|
||||
return pending
|
||||
|
||||
@classmethod
|
||||
def stage_admit(cls, db: Session, *, task_id: str, storage: str,
|
||||
src_path: str, state: str,
|
||||
now_time: str) -> Optional["TransferPending"]:
|
||||
"""
|
||||
在调用方会话中暂存一条持久接纳记录。
|
||||
|
||||
相同存储与路径已存在时返回原记录,确保重复监控事件复用同一个任务标识。
|
||||
:param db: 数据库会话
|
||||
:param task_id: 任务标识
|
||||
:param storage: 存储
|
||||
:param src_path: 源文件路径
|
||||
:param state: 持久状态
|
||||
:param now_time: 当前时间
|
||||
:return: 接纳记录
|
||||
"""
|
||||
if not task_id or not storage or not src_path or not state:
|
||||
return None
|
||||
pending = db.execute(
|
||||
select(cls).where(cls.storage == storage, cls.src_path == src_path)
|
||||
).scalars().first()
|
||||
if pending:
|
||||
return cast("TransferPending", pending)
|
||||
pending = cls(
|
||||
task_id=task_id,
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
state=state,
|
||||
created_at=now_time,
|
||||
updated_at=now_time,
|
||||
)
|
||||
db.add(pending)
|
||||
return pending
|
||||
|
||||
@classmethod
|
||||
def list_by_state(cls, db: Session, *, state: str,
|
||||
limit: Optional[int] = 5000) -> List["TransferPending"]:
|
||||
"""
|
||||
按登记顺序列出指定持久状态的接纳记录。
|
||||
:param db: 数据库会话
|
||||
:param state: 持久状态
|
||||
:param limit: 单次读取上限
|
||||
:return: 接纳记录列表
|
||||
"""
|
||||
if not state:
|
||||
return []
|
||||
return list(db.execute(
|
||||
select(cls)
|
||||
.where(cls.state == state)
|
||||
.order_by(cls.created_at.asc(), cls.id.asc())
|
||||
.limit(limit)
|
||||
).scalars().all())
|
||||
|
||||
@classmethod
|
||||
def get_by_identity(cls, db: Session, *, storage: str,
|
||||
src_path: str) -> Optional["TransferPending"]:
|
||||
"""
|
||||
按存储与源路径查询一条持久接纳记录。
|
||||
:param db: 数据库会话
|
||||
:param storage: 存储
|
||||
:param src_path: 源文件路径
|
||||
:return: 接纳记录
|
||||
"""
|
||||
if not storage or not src_path:
|
||||
return None
|
||||
return cast(
|
||||
Optional["TransferPending"],
|
||||
db.execute(
|
||||
select(cls).where(
|
||||
cls.storage == storage,
|
||||
cls.src_path == src_path,
|
||||
)
|
||||
).scalars().first(),
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def record_enqueue_failure(cls, db: Session, *, task_id: str,
|
||||
error: str, now_time: str) -> int:
|
||||
"""
|
||||
在调用方会话中记录任务最近一次入队失败。
|
||||
:param db: 数据库会话
|
||||
:param task_id: 任务标识
|
||||
:param error: 失败原因
|
||||
:param now_time: 当前时间
|
||||
:return: 更新的记录数
|
||||
"""
|
||||
if not task_id:
|
||||
return 0
|
||||
return execute_dml(
|
||||
db,
|
||||
update(cls)
|
||||
.where(cls.task_id == task_id)
|
||||
.values(last_error=error, updated_at=now_time),
|
||||
execution_options={"synchronize_session": False},
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def discard_task(cls, db: Session, *, task_id: str) -> int:
|
||||
"""
|
||||
在调用方会话中按任务标识删除接纳记录。
|
||||
:param db: 数据库会话
|
||||
:param task_id: 任务标识
|
||||
:return: 删除的记录数
|
||||
"""
|
||||
if not task_id:
|
||||
return 0
|
||||
return execute_dml(
|
||||
db, delete(cls).where(cls.task_id == task_id),
|
||||
execution_options={"synchronize_session": False},
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def discard(cls, db: Session, storage: str, src_path: str) -> int:
|
||||
"""
|
||||
|
||||
@@ -9,8 +9,9 @@ class TransferPendingOper(DbOper):
|
||||
"""
|
||||
待整理文件登记管理。
|
||||
|
||||
只保存「存储 + 源文件路径」这一最小事实,用于在进程重启后把没走完整理链的
|
||||
文件重新送回去,避免挂载故障重启后永久漏件。
|
||||
保存稳定任务身份、存储、源文件路径和准入状态,用于在进程重启后把没走完
|
||||
整理链的文件重新送回去,避免挂载故障重启后永久漏件。旧版路径登记接口继续
|
||||
保留,供插件和兼容调用方使用。
|
||||
"""
|
||||
|
||||
def register(self, storage: str, src_path: str) -> Optional[TransferPending]:
|
||||
@@ -30,6 +31,93 @@ class TransferPendingOper(DbOper):
|
||||
)
|
||||
)
|
||||
|
||||
def stage_admit(self, *, task_id: str, storage: str, src_path: str,
|
||||
state: str, now_time: str) -> Optional[TransferPending]:
|
||||
"""
|
||||
在当前会话中暂存一条持久接纳记录。
|
||||
|
||||
适配器应传入显式 Session,使提交与回滚仍由应用用例对应的 UoW 管理。
|
||||
:param task_id: 任务标识
|
||||
:param storage: 存储
|
||||
:param src_path: 源文件路径
|
||||
:param state: 持久状态
|
||||
:param now_time: 当前时间
|
||||
:return: 接纳记录
|
||||
"""
|
||||
return self._execute_sync_write(
|
||||
lambda session: TransferPending.stage_admit(
|
||||
session,
|
||||
task_id=task_id,
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
state=state,
|
||||
now_time=now_time,
|
||||
)
|
||||
)
|
||||
|
||||
def list_by_state(self, *, state: str,
|
||||
limit: Optional[int] = 5000) -> List[TransferPending]:
|
||||
"""
|
||||
使用当前会话列出指定状态记录。
|
||||
:param state: 持久状态
|
||||
:param limit: 单次读取上限
|
||||
:return: ORM 接纳记录列表
|
||||
"""
|
||||
return self._execute_sync_query(
|
||||
lambda session: TransferPending.list_by_state(
|
||||
session,
|
||||
state=state,
|
||||
limit=limit,
|
||||
)
|
||||
) or []
|
||||
|
||||
def get_by_identity(self, *, storage: str,
|
||||
src_path: str) -> Optional[TransferPending]:
|
||||
"""
|
||||
使用当前会话按存储与源路径查询接纳记录。
|
||||
:param storage: 存储
|
||||
:param src_path: 源文件路径
|
||||
:return: 接纳记录
|
||||
"""
|
||||
return self._execute_sync_query(
|
||||
lambda session: TransferPending.get_by_identity(
|
||||
session,
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
)
|
||||
)
|
||||
|
||||
def stage_record_enqueue_failure(self, *, task_id: str, error: str,
|
||||
now_time: str) -> int:
|
||||
"""
|
||||
在当前会话中暂存最近一次入队失败。
|
||||
:param task_id: 任务标识
|
||||
:param error: 失败原因
|
||||
:param now_time: 当前时间
|
||||
:return: 更新的记录数
|
||||
"""
|
||||
return self._execute_sync_write(
|
||||
lambda session: TransferPending.record_enqueue_failure(
|
||||
session,
|
||||
task_id=task_id,
|
||||
error=error,
|
||||
now_time=now_time,
|
||||
)
|
||||
)
|
||||
|
||||
def stage_discard_task(self, *, task_id: str) -> int:
|
||||
"""
|
||||
在当前会话中暂存按任务标识删除接纳记录。
|
||||
:param task_id: 任务标识
|
||||
:return: 删除的记录数
|
||||
"""
|
||||
return self._execute_sync_write(
|
||||
lambda session: TransferPending.discard_task(
|
||||
session,
|
||||
task_id=task_id,
|
||||
)
|
||||
)
|
||||
|
||||
def discard(self, storage: str, src_path: str) -> int:
|
||||
"""
|
||||
注销一个待整理文件登记。
|
||||
|
||||
@@ -110,6 +110,7 @@ from app.db.adapters.outbox import SqlAlchemyAsyncOutboxStager, SqlAlchemyOutbox
|
||||
from app.db.adapters.site import TransactionalSiteRepository
|
||||
from app.db.adapters.subscription import TransactionalSubscribeWriter
|
||||
from app.db.adapters.transaction import TransactionalWriteRunner
|
||||
from app.db.adapters.transfer import TransactionalTransferAdmissionRepository
|
||||
from app.db.adapters.workflow import TransactionalWorkflowExecutionService
|
||||
from app.db.oper.agentchat import AgentChatOper
|
||||
from app.db.oper.agenttask import AgentTaskOper
|
||||
@@ -123,7 +124,6 @@ from app.db.oper.subscribe import SubscribeOper
|
||||
from app.db.oper.subscribehistory import SubscribeHistoryOper
|
||||
from app.db.oper.systemconfig import SystemConfigOper
|
||||
from app.db.oper.transferhistory import TransferHistoryOper
|
||||
from app.db.oper.transferpending import TransferPendingOper
|
||||
from app.db.oper.user import UserOper
|
||||
from app.db.oper.userconfig import UserConfigOper
|
||||
from app.db.oper.workflow import WorkflowOper, configure_workflow_legacy_writer
|
||||
@@ -854,7 +854,9 @@ async def init_modules() -> HostRuntime:
|
||||
workflow=lambda: WorkflowOper(),
|
||||
download_history=lambda: DownloadHistoryOper(),
|
||||
transfer_history=lambda: TransferHistoryOper(),
|
||||
transfer_pending=lambda: TransferPendingOper(),
|
||||
transfer_pending=lambda: TransactionalTransferAdmissionRepository(
|
||||
SessionFactory
|
||||
),
|
||||
media_server=lambda: MediaServerOper(),
|
||||
download_failure=lambda: TransactionalDownloadFailureRepository(
|
||||
SessionFactory
|
||||
|
||||
Reference in New Issue
Block a user