Skip to content

Commit aac184a

Browse files
committed
refactor(transfer): make queue admission durable
1 parent 1133557 commit aac184a

23 files changed

Lines changed: 1548 additions & 195 deletions

app/application/chain/data.py

Lines changed: 6 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,11 @@
1010
from dataclasses import dataclass
1111
from typing import Any, Optional
1212

13+
from app.application.transfer import TransferAdmissionRepository
14+
1315

1416
OperFactory = Callable[[], Any]
17+
TransferAdmissionRepositoryFactory = Callable[[], TransferAdmissionRepository]
1518

1619

1720
@dataclass(frozen=True, slots=True)
@@ -23,7 +26,7 @@ class ChainDataPorts:
2326
workflow: OperFactory
2427
download_history: OperFactory
2528
transfer_history: OperFactory
26-
transfer_pending: OperFactory
29+
transfer_pending: TransferAdmissionRepositoryFactory
2730
media_server: OperFactory
2831
download_failure: OperFactory
2932
user: OperFactory
@@ -77,12 +80,6 @@ class TransferHistoryPortProxy(_ChainDataPortProxy):
7780
port_name = "transfer_history"
7881

7982

80-
class TransferPendingPortProxy(_ChainDataPortProxy):
81-
"""待整理数据端口代理。"""
82-
83-
port_name = "transfer_pending"
84-
85-
8683
class MediaServerPortProxy(_ChainDataPortProxy):
8784
"""媒体服务器数据端口代理。"""
8885

@@ -156,8 +153,8 @@ def get_chain_transfer_history_port() -> Any:
156153
return get_chain_data_ports().transfer_history()
157154

158155

159-
def get_chain_transfer_pending_port() -> Any:
160-
"""创建待整理数据端口实例。"""
156+
def get_chain_transfer_pending_port() -> TransferAdmissionRepository:
157+
"""创建类型化的整理任务 durable admission 仓储。"""
161158
return get_chain_data_ports().transfer_pending()
162159

163160

app/application/transfer.py

Lines changed: 72 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -19,14 +19,10 @@
1919
from dataclasses import dataclass
2020
from pathlib import Path
2121
from time import monotonic
22-
from typing import Callable, Dict, List, Optional, Tuple, Union
22+
from typing import Callable, Dict, List, Optional, Protocol, Tuple, Union
2323

24-
from pydantic import BaseModel, ConfigDict
24+
from pydantic import BaseModel, ConfigDict, PrivateAttr
2525

26-
from app.schemas.transfer import MetaInfo as _SchemaMetaInfo
27-
from app.schemas.transfer import MusicInfo as _SchemaMusicInfo
28-
from app.schemas.transfer import MusicMeta as _SchemaMusicMeta
29-
from app.schemas.workflow import MediaInfo as _SchemaMediaInfo
3026
from app.adapters.system.host import SystemUtils
3127
from app.application.agent import get_prompt_manager, get_running_agent_manager
3228
from app.domain.context import MediaInfo, MusicInfo
@@ -40,6 +36,9 @@
4036
from app.schemas.media import OptionalMediaIdentityMixin, resolve_media_identity
4137
from app.schemas.system import TransferDirectoryConf
4238
from app.schemas.tmdb import TmdbEpisode
39+
from app.schemas.transfer import MetaInfo as _SchemaMetaInfo
40+
from app.schemas.transfer import MusicInfo as _SchemaMusicInfo
41+
from app.schemas.transfer import MusicMeta as _SchemaMusicMeta
4342
from app.schemas.transfer import TransferInfo, TransferJob, TransferJobTask
4443
from app.schemas.types import (
4544
MUSIC_ENTITY_ALBUM,
@@ -48,7 +47,7 @@
4847
MediaType,
4948
ReplyMode,
5049
)
51-
50+
from app.schemas.workflow import MediaInfo as _SchemaMediaInfo
5251

5352

5453
class TransferTask(OptionalMediaIdentityMixin, BaseModel):
@@ -81,6 +80,16 @@ class TransferTask(OptionalMediaIdentityMixin, BaseModel):
8180
manual: Optional[bool] = False
8281
background: Optional[bool] = True
8382
preview: Optional[bool] = False
83+
_admission_task_id: Optional[str] = PrivateAttr(default=None)
84+
85+
@property
86+
def admission_task_id(self) -> Optional[str]:
87+
"""返回仅供宿主持久准入和终态结算使用的内部任务标识。"""
88+
return self._admission_task_id
89+
90+
def bind_admission_task_id(self, task_id: str) -> None:
91+
"""绑定持久准入生成的稳定身份,不改变插件可见序列化字段。"""
92+
self._admission_task_id = task_id
8493

8594
def to_dict(self):
8695
"""
@@ -113,36 +122,86 @@ class TransferQueue(BaseModel):
113122
result: Optional[TransferInfo] = None
114123

115124

125+
TRANSFER_ADMISSION_ACCEPTED = "accepted"
126+
127+
128+
@dataclass(frozen=True, slots=True)
129+
class TransferAdmission:
130+
"""描述已经持久化、可在进程退出后恢复的整理任务准入事实。"""
131+
132+
task_id: str
133+
storage: str
134+
src_path: str
135+
state: str
136+
created_at: str
137+
updated_at: str
138+
last_error: Optional[str] = None
139+
140+
141+
class TransferAdmissionRepository(Protocol):
142+
"""整理任务 durable admission 所需的类型化持久化端口。"""
143+
144+
def admit(self, *, storage: str, src_path: str) -> TransferAdmission:
145+
"""幂等登记源文件并返回稳定任务身份。"""
146+
...
147+
148+
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
149+
"""按登记顺序返回等待恢复或执行的任务。"""
150+
...
151+
152+
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
153+
"""记录内存队列接收失败,保留任务供后续恢复。"""
154+
...
155+
156+
def discard_task(self, *, task_id: str) -> int:
157+
"""按稳定任务身份删除已经到达终态的登记。"""
158+
...
159+
160+
116161
class TransferQueueService:
117162
"""协调整理任务登记、入队、移除和队列视图查询。"""
118163

119164
def __init__(
120165
self,
121166
*,
122167
register_task: Callable[[TransferTask], bool],
168+
admit_task: Callable[[TransferTask], TransferAdmission],
123169
enqueue: Callable[[TransferQueue], None],
124170
before_enqueue: Callable[[TransferTask], None],
125-
after_enqueue: Callable[[TransferTask], None],
171+
enqueue_failed: Callable[[TransferTask, Exception], None],
126172
remove_task: Callable[[FileItem], None],
127173
list_tasks: Callable[[], List[TransferJob]],
128174
expire_tasks: Callable[[], None],
129175
) -> None:
130176
"""保存队列用例依赖,避免 Application 服务绑定具体线程队列实现。"""
131177
self._register_task = register_task
178+
self._admit_task = admit_task
132179
self._enqueue = enqueue
133180
self._before_enqueue = before_enqueue
134-
self._after_enqueue = after_enqueue
181+
self._enqueue_failed = enqueue_failed
135182
self._remove_task = remove_task
136183
self._list_tasks = list_tasks
137184
self._expire_tasks = expire_tasks
138185

139186
def put(self, task: TransferTask, callback: Callable) -> bool:
140-
"""登记并入队一个整理任务,保持原有副作用顺序。"""
187+
"""先持久化准入事实再入队;任何前置失败都撤销内存作业视图。"""
141188
if not task or not self._register_task(task):
142189
return False
143-
self._before_enqueue(task)
144-
self._enqueue(TransferQueue(task=task, callback=callback))
145-
self._after_enqueue(task)
190+
try:
191+
admission = self._admit_task(task)
192+
task.bind_admission_task_id(admission.task_id)
193+
except Exception:
194+
self._remove_task(task.fileitem)
195+
raise
196+
try:
197+
self._before_enqueue(task)
198+
self._enqueue(TransferQueue(task=task, callback=callback))
199+
except Exception as err:
200+
try:
201+
self._enqueue_failed(task, err)
202+
finally:
203+
self._remove_task(task.fileitem)
204+
raise
146205
return True
147206

148207
def remove(self, fileitem: FileItem) -> None:

app/chain/transfer.py

Lines changed: 58 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@
5252
from app.application.transfer import (
5353
FailedRetryScheduler,
5454
JobManager,
55+
TransferAdmission,
5556
TransferFailureNotification,
5657
TransferFailureNotificationAggregator,
5758
TransferQueue,
@@ -225,8 +226,8 @@ def __init__(self) -> None:
225226
self.retry_scheduler = FailedRetryScheduler()
226227
# 整理失败通知聚合器
227228
self.failure_notification_aggregator = TransferFailureNotificationAggregator()
228-
# 待整理文件落盘登记,用于进程重启后回放内存队列里未完成的任务
229-
self._pendingoper = get_chain_transfer_pending_port()
229+
# durable admission 仓储先于内存入队保存任务,进程退出后仍可恢复。
230+
self._transfer_admissions = get_chain_transfer_pending_port()
230231
# 转移成功的文件清单
231232
self._success_target_files: Dict[Tuple, List[str]] = {}
232233
# 批次级刮削缓冲,避免同一批多文件入库重复触发目录刮削
@@ -867,7 +868,8 @@ def put_to_queue(self, task: TransferTask) -> bool:
867868
"""
868869
添加到待整理队列
869870
:param task: 任务信息
870-
:return: True表示任务已添加到队列,False表示任务无效或已存在(重复)
871+
:return: True表示任务已添加,False表示链已关闭或任务无效/重复
872+
:raises Exception: 持久准入、批次登记或内存入队失败
871873
"""
872874
with self._worker_state_lock:
873875
if self._closing:
@@ -879,9 +881,10 @@ def _transfer_queue_service(self) -> TransferQueueService:
879881
"""构建保持旧队列对象和私有兼容接缝的应用服务。"""
880882
return TransferQueueService(
881883
register_task=self.__put_to_jobview,
884+
admit_task=self.__admit_transfer,
882885
enqueue=self._queue.put,
883886
before_enqueue=self._register_scrape_batch_task,
884-
after_enqueue=self.__register_pending,
887+
enqueue_failed=self.__record_enqueue_failure,
885888
remove_task=self.jobview.remove_task,
886889
list_tasks=self.jobview.list_jobs,
887890
expire_tasks=self.__expire_stale_transfer_tasks,
@@ -936,17 +939,19 @@ def __replay_pending(
936939
if stop_event.is_set():
937940
return
938941
try:
939-
pendings = self._pendingoper.list_all()
942+
pendings = self._transfer_admissions.list_accepted()
940943
except Exception as err:
941944
logger.error(f"读取待整理文件登记失败:{err}")
942945
return
943946
if not pendings:
944947
return
945948
logger.info(f"发现 {len(pendings)} 个上次未整理完的文件,正在重新送入整理链 ...")
946949
replayed = 0
947-
for storage, src_path in pendings:
950+
for admission in pendings:
948951
if stop_event.is_set():
949952
break
953+
storage = admission.storage
954+
src_path = admission.src_path
950955
try:
951956
fileitem, should_discard = self.__build_replay_fileitem(storage, src_path)
952957
# stat 等同步 I/O 返回后重新检查,关闭期间不得注销尚未完成的登记。
@@ -955,7 +960,9 @@ def __replay_pending(
955960
if not fileitem:
956961
if should_discard:
957962
# 源文件确认已消失,注销登记避免每次启动重复回放
958-
self._pendingoper.discard(storage=storage, src_path=src_path)
963+
self._transfer_admissions.discard_task(
964+
task_id=admission.task_id
965+
)
959966
continue
960967
self.do_transfer(fileitem=fileitem)
961968
replayed += 1
@@ -1009,19 +1016,35 @@ def __build_replay_fileitem(storage: str, src_path: str) -> Tuple[Optional[FileI
10091016
modify_time=modify_time,
10101017
), False
10111018

1012-
def __register_pending(self, task: TransferTask):
1013-
"""
1014-
落盘登记一个待整理文件,登记失败不影响正常入队。
1015-
:param task: 任务信息
1016-
"""
1019+
def __admit_transfer(self, task: TransferTask) -> TransferAdmission:
1020+
"""在内存入队前持久化源文件并返回稳定任务身份。"""
10171021
fileitem = task.fileitem if task else None
1018-
if not fileitem or not fileitem.path:
1019-
return
1022+
if not fileitem or not fileitem.storage or not fileitem.path:
1023+
raise ValueError("整理任务缺少源文件身份")
1024+
return self._transfer_admissions.admit(
1025+
storage=fileitem.storage,
1026+
src_path=fileitem.path,
1027+
)
1028+
1029+
def __record_enqueue_failure(
1030+
self,
1031+
task: TransferTask,
1032+
error: Exception,
1033+
) -> None:
1034+
"""记录内存入队失败并撤销该任务的批次占位。"""
10201035
try:
1021-
self._pendingoper.register(storage=fileitem.storage, src_path=fileitem.path)
1022-
except Exception as err:
1023-
# 登记只是重启后的补救手段,失败不能阻断正常整理
1024-
logger.debug(f"登记待整理文件失败: {fileitem.path} - {err}")
1036+
if task.admission_task_id:
1037+
self._transfer_admissions.record_enqueue_failure(
1038+
task_id=task.admission_task_id,
1039+
error=str(error),
1040+
)
1041+
except Exception as record_error:
1042+
logger.error(
1043+
"记录整理任务入队失败原因异常:"
1044+
f"{task.admission_task_id} - {record_error}"
1045+
)
1046+
finally:
1047+
self._finish_scrape_batch_task(task)
10251048

10261049
def __discard_pending(self, task: TransferTask):
10271050
"""
@@ -1031,13 +1054,17 @@ def __discard_pending(self, task: TransferTask):
10311054
重试预算重新送入整理链;留在本表反而会每次重启都重复回放。
10321055
:param task: 任务信息
10331056
"""
1034-
fileitem = task.fileitem if task else None
1035-
if not fileitem or not fileitem.path:
1057+
if not task or not task.admission_task_id:
10361058
return
10371059
try:
1038-
self._pendingoper.discard(storage=fileitem.storage, src_path=fileitem.path)
1060+
self._transfer_admissions.discard_task(
1061+
task_id=task.admission_task_id
1062+
)
10391063
except Exception as err:
1040-
logger.debug(f"注销待整理文件登记失败: {fileitem.path} - {err}")
1064+
logger.error(
1065+
"注销整理任务 durable admission 失败: "
1066+
f"{task.admission_task_id} - {err}"
1067+
)
10411068

10421069
def __put_to_jobview(self, task: TransferTask) -> bool:
10431070
"""
@@ -2751,7 +2778,15 @@ def _get_cached_extra_meta(
27512778
preview=preview,
27522779
)
27532780
if background:
2754-
if self.put_to_queue(task=transfer_task):
2781+
try:
2782+
queued = self.put_to_queue(task=transfer_task)
2783+
except Exception as err:
2784+
all_success = False
2785+
message = f"{file_path.name} 加入整理队列失败:{err}"
2786+
err_msgs.append(message)
2787+
logger.error(message)
2788+
continue
2789+
if queued:
27552790
logger.info(f"{file_path.name} 已添加到整理队列")
27562791
else:
27572792
logger.debug(f"{file_path.name} 已在整理队列中,跳过")

0 commit comments

Comments
 (0)