refactor: extract download success settlement

This commit is contained in:
jxxghp
2026-08-22 20:46:48 +08:00
parent 8d0a694816
commit 33ca1f824f
4 changed files with 115 additions and 113 deletions
+3 -4
View File
@@ -279,9 +279,8 @@ class RssHelper:
return False
try:
config = get_chain_runtime_config_snapshot()
ret = RequestUtils(ua=ua,
proxies=config.proxy if proxy else None,
proxies=get_chain_runtime_config_snapshot().proxy if proxy else None,
timeout=timeout or 30, headers=headers).get_res(url)
if not ret:
logger.error(f"获取RSS失败:请求返回空值,URL: {url}")
@@ -307,8 +306,8 @@ class RssHelper:
if raw_data:
ret_xml = RequestUtils.get_decoded_xml_content(
ret,
performance_mode=config.encoding_detection_performance_mode,
confidence_threshold=config.encoding_detection_min_confidence
performance_mode=get_chain_runtime_config_snapshot().encoding_detection_performance_mode,
confidence_threshold=get_chain_runtime_config_snapshot().encoding_detection_min_confidence
)
rust_items = self.__parse_with_rust(ret_xml)
if rust_items is not None:
+109 -108
View File
@@ -1097,6 +1097,95 @@ class DownloadChain(ChainBase):
logger.warn(str(err))
return save_path, str(err)
def _settle_download_success(
self,
*,
context: Context,
media: MediaInfo,
meta: MetaInfo,
torrent: TorrentInfo,
folder_name: str,
file_list: list,
download_dir: Path,
layout: Optional[str],
downloader: Optional[str],
download_hash: str,
download_episodes: Optional[str],
episodes: Optional[Set[int]],
channel: Optional[NotificationChannel],
source: Optional[str],
userid: Union[str, int, None],
username: Optional[str],
torrent_content: Union[str, bytes],
custom_words: Optional[str],
) -> None:
"""提交下载历史、文件明细和 durable 下载事件。"""
if layout == "NoSubfolder" or not folder_name:
download_path = download_dir / file_list[0] if file_list else download_dir
elif folder_name:
download_path = download_dir / folder_name
else:
download_path = download_dir / Path(file_list[0]).stem if file_list else download_dir
save_path = download_dir if layout == "NoSubfolder" or not folder_name else download_path
media_source, media_id = resolve_media_identity(media=media)
history_payload = {
"path": download_path.as_posix(), "type": media.type.value,
"title": media.title, "year": media.year,
"media_source": media_source, "media_id": media_id,
"music_type": getattr(media, "music_type", None), "seasons": meta.season,
"episodes": download_episodes or meta.episode,
"image": media.get_backdrop_image(), "poster": media.get_poster_image(),
"downloader": downloader, "download_hash": download_hash,
"torrent_name": torrent.title, "torrent_description": torrent.description,
"torrent_site": torrent.site_name, "userid": userid, "username": username,
"channel": channel.value if channel else None,
"date": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime()),
"media_category": media.category, "episode_group": media.episode_group,
"note": self._build_download_note(source, media, meta),
"custom_words": custom_words,
}
files_to_add = []
for file in file_list:
if episodes:
file_meta = MetaInfo(Path(file).stem)
if not file_meta.begin_episode or file_meta.begin_episode not in episodes:
continue
if not Path(file).suffix or Path(file).suffix.lower() not in self.runtime_config.media_extensions:
continue
files_to_add.append({
"download_hash": download_hash, "downloader": downloader,
"fullpath": (save_path / file).as_posix(), "savepath": save_path.as_posix(),
"filepath": file, "torrentname": meta.org_string,
})
event_payload = {
"hash": download_hash, "context": context, "username": username,
"downloader": downloader, "episodes": episodes or meta.episode_list, "source": source,
}
def after_commit() -> None:
"""在历史与 intent 提交后保持原有通知和任务编排。"""
self._after_download_history_commit(
context=context, media=media, meta=meta, torrent=torrent,
channel=channel, source=source, userid=userid, username=username,
download_episodes=download_episodes, download_dir=download_dir,
torrent_content=torrent_content,
)
durable_event_writer = getattr(self, "durable_event_writer", None)
if durable_event_writer:
durable_event_writer.download_added(
history_payload=history_payload, file_payloads=files_to_add,
event_payload=event_payload, after_commit=after_commit,
publish=lambda payload: self.eventmanager.send_event(EventType.DownloadAdded, payload),
)
return
downloadhis = DownloadHistoryOper()
downloadhis.add(**history_payload)
if files_to_add:
downloadhis.add_files(files_to_add)
after_commit()
self.eventmanager.send_event(EventType.DownloadAdded, event_payload)
def download_single(self, context: Context,
torrent_file: Path = None,
torrent_content: Optional[Union[str, bytes]] = None,
@@ -1218,114 +1307,26 @@ class DownloadChain(ChainBase):
_downloader, _hash, _layout, error_msg = None, None, None, "未找到下载器"
if _hash:
# `不创建子文件夹` 或 `不存在子文件夹`
if _layout == "NoSubfolder" or not _folder_name:
# 下载路径记录至文件
download_path = download_dir / _file_list[0] if _file_list else download_dir
# 原始布局
elif _folder_name:
download_path = download_dir / _folder_name
# 创建子文件夹
else:
download_path = download_dir / Path(_file_list[0]).stem if _file_list else download_dir
# 文件保存路径
_save_path = download_dir if _layout == "NoSubfolder" or not _folder_name else download_path
media_source, media_id = resolve_media_identity(media=_media)
history_payload = {
"path": download_path.as_posix(),
"type": _media.type.value,
"title": _media.title,
"year": _media.year,
"media_source": media_source,
"media_id": media_id,
"music_type": getattr(_media, "music_type", None),
"seasons": _meta.season,
"episodes": download_episodes or _meta.episode,
"image": _media.get_backdrop_image(),
"poster": _media.get_poster_image(),
"downloader": _downloader,
"download_hash": _hash,
"torrent_name": _torrent.title,
"torrent_description": _torrent.description,
"torrent_site": _torrent.site_name,
"userid": userid,
"username": username,
"channel": channel.value if channel else None,
"date": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime()),
"media_category": _media.category,
"episode_group": _media.episode_group,
"note": self._build_download_note(source, _media, _meta),
"custom_words": custom_words,
}
# 登记下载文件
files_to_add = []
for file in _file_list:
if episodes:
# 识别文件集
file_meta = MetaInfo(Path(file).stem)
if not file_meta.begin_episode \
or file_meta.begin_episode not in episodes:
continue
# 只处理音视频、字幕格式
media_exts = self.runtime_config.media_extensions
if not Path(file).suffix \
or Path(file).suffix.lower() not in media_exts:
continue
files_to_add.append({
"download_hash": _hash,
"downloader": _downloader,
"fullpath": (_save_path / file).as_posix(),
"savepath": _save_path.as_posix(),
"filepath": file,
"torrentname": _meta.org_string,
})
event_payload = {
"hash": _hash,
"context": context,
"username": username,
"downloader": _downloader,
"episodes": episodes or _meta.episode_list,
"source": source,
}
def after_commit() -> None:
"""在历史与 intent 提交后保持原有通知和任务编排。"""
self._after_download_history_commit(
context=context,
media=_media,
meta=_meta,
torrent=_torrent,
channel=channel,
source=source,
userid=userid,
username=username,
download_episodes=download_episodes,
download_dir=download_dir,
torrent_content=torrent_content,
)
durable_event_writer = getattr(self, "durable_event_writer", None)
if durable_event_writer:
durable_event_writer.download_added(
history_payload=history_payload,
file_payloads=files_to_add,
event_payload=event_payload,
after_commit=after_commit,
publish=lambda payload: self.eventmanager.send_event(
EventType.DownloadAdded,
payload,
),
)
else:
# 显式注入旧测试上下文时保持兼容;正式启动上下文总会提供 durable writer。
downloadhis = DownloadHistoryOper()
downloadhis.add(**history_payload)
if files_to_add:
downloadhis.add_files(files_to_add)
after_commit()
self.eventmanager.send_event(EventType.DownloadAdded, event_payload)
self._settle_download_success(
context=context,
media=_media,
meta=_meta,
torrent=_torrent,
folder_name=_folder_name,
file_list=_file_list,
download_dir=download_dir,
layout=_layout,
downloader=_downloader,
download_hash=_hash,
download_episodes=download_episodes,
episodes=episodes,
channel=channel,
source=source,
userid=userid,
username=username,
torrent_content=torrent_content,
custom_words=custom_words,
)
else:
# 下载失败
logger.error(f"{_media.title_year} 添加下载任务失败:"
@@ -983,6 +983,8 @@ Outbox adapter、DB 装饰器、Base 与 UoWstrict 清单扩大到 37 个源
清理服务与 Chain 专项测试通过。
Passkey 的 APP_DOMAIN、NGINX_PORT 和用户验证要求也已接入 API 配置快照,配置债务降至 131 个文件;
MFA/Passkey 专项测试与架构门禁通过,密钥类配置仍保留在安全端口范围内。
`DownloadChain.download_single` 的下载成功结算已提取为独立阶段,入口从 255 行降至 167 行;
历史、文件明细、durable intent、post-commit 通知和旧测试 fallback 语义保持,下载专项测试通过。
#### ARCH-272:异步阻塞检测
+1 -1
View File
@@ -18,7 +18,7 @@
},
"chain_public": {
"app/chain/download.py:DownloadChain.batch_download": 572,
"app/chain/download.py:DownloadChain.download_single": 255,
"app/chain/download.py:DownloadChain.download_single": 167,
"app/chain/mediaserver.py:MediaServerChain.sync": 292,
"app/chain/subscribe.py:SubscribeChain.add": 183,
"app/chain/subscribe.py:SubscribeChain.async_add": 186,