diff --git a/app/chain/search.py b/app/chain/search.py index 93bbcb1ca..2a61c0c2c 100644 --- a/app/chain/search.py +++ b/app/chain/search.py @@ -5,8 +5,9 @@ import random import re import time from concurrent.futures import FIRST_COMPLETED, ThreadPoolExecutor, as_completed, wait +from contextlib import aclosing from datetime import datetime -from typing import AsyncIterator, Any, Dict, Iterable, Tuple +from typing import AsyncIterator, Any, Awaitable, Callable, Dict, Iterable, Tuple from typing import List, Optional from unicodedata import normalize @@ -2420,6 +2421,68 @@ class SearchChain(ChainBase): # 返回 return results + async def _iter_site_page_results( + self, + *, + indexer_sites: List[dict], + search_pages: List[int], + search_page: Callable[[dict, int], Awaitable[Optional[List[Any]]]], + should_continue: Callable[[dict, List[Any]], bool], + task_owner: str, + ) -> AsyncIterator[Tuple[dict, int, List[Any], bool]]: + """统一调度站点逐页请求,并在调用方退出时取消、等待全部请求。""" + total_num = len(indexer_sites) * len(search_pages) + semaphore = asyncio.Semaphore( + self.runtime_config.search_threadpool_size or max(1, total_num) + ) + pending_tasks: dict[ + asyncio.Task[List[Any]], + Tuple[dict, int, int], + ] = {} + + async def run_site_page(site: dict, page_number: int) -> List[Any]: + """在共享并发预算内执行一页站点请求并规范化空结果。""" + async with semaphore: + return await search_page(site, page_number) or [] + + def submit_site_page(site: dict, page_index: int) -> None: + """登记一页请求及其续页位置,供统一终态收口。""" + page_number = search_pages[page_index] + task = asyncio.create_task( + run_site_page(site, page_number), + name=task_owner, + ) + pending_tasks[task] = (site, page_index, page_number) + + for site in indexer_sites: + submit_site_page(site, 0) + + try: + while pending_tasks: + if global_vars.is_system_stopped: + break + done_tasks, _ = await asyncio.wait( + pending_tasks, + return_when=asyncio.FIRST_COMPLETED, + ) + for task in done_tasks: + site, page_index, page_number = pending_tasks.pop(task) + page_results = await task + continued = ( + should_continue(site, page_results) + and page_index + 1 < len(search_pages) + ) + if continued: + submit_site_page(site, page_index + 1) + yield site, page_number, page_results, continued + finally: + tasks = tuple(pending_tasks) + for task in tasks: + if not task.done(): + task.cancel() + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + async def __async_search_all_sites(self, keyword: str, mediainfo: Optional[MediaInfo] = None, sites: List[int] = None, @@ -2472,75 +2535,57 @@ class SearchChain(ChainBase): text=f"开始搜索,共 {len(indexer_sites)} 个站点,{len(search_pages)} 页 ...") # 结果集 results = list(plugin_results) - semaphore = asyncio.Semaphore( - self.runtime_config.search_threadpool_size or total_num - ) async def search_site_page(site: dict, search_page: int) -> List[TorrentInfo]: - """ - 控制单次站点页请求的并发量,并返回该页的资源列表。 - """ - async with semaphore: - if area == "imdbid": - # 搜索IMDBID - return await self.async_search_site_torrents(site=site, - keyword=mediainfo.imdb_id if mediainfo else None, - mtype=mediainfo.type if mediainfo else mtype, - page=search_page) - # 搜索标题 - return await self.async_search_site_torrents(site=site, - keyword=keyword, - mtype=mediainfo.type if mediainfo else mtype, - page=search_page) + """调用既有站点资源接口,具体并发由统一逐页编排器控制。""" + search_keyword = ( + mediainfo.imdb_id + if area == "imdbid" and mediainfo + else keyword + ) + return await self.async_search_site_torrents( + site=site, + keyword=search_keyword, + mtype=mediainfo.type if mediainfo else mtype, + page=search_page, + ) - pending_tasks = {} + def should_continue(site: dict, page_results: List[Any]) -> bool: + """按既有站点分页规则判断是否提交下一页资源请求。""" + search_keyword = ( + mediainfo.imdb_id + if area == "imdbid" and mediainfo + else keyword + ) + return self._should_continue_search_pages( + site=site, + page_results=page_results, + keyword=search_keyword, + ) - def submit_site_page(site: dict, page_index: int): - """ - 提交异步站点页搜索任务,并记录该任务对应的站点和页码位置。 - """ - search_page = search_pages[page_index] - search_keyword = mediainfo.imdb_id if area == "imdbid" and mediainfo else keyword - task = asyncio.create_task(search_site_page(site=site, search_page=search_page)) - pending_tasks[task] = (site, page_index, search_page, search_keyword) - - for site in indexer_sites: - submit_site_page(site=site, page_index=0) - - try: - while pending_tasks: - if global_vars.is_system_stopped: - break - done_tasks, _ = await asyncio.wait( - pending_tasks.keys(), - return_when=asyncio.FIRST_COMPLETED, + page_iterator = self._iter_site_page_results( + indexer_sites=indexer_sites, + search_pages=search_pages, + search_page=search_site_page, + should_continue=should_continue, + task_owner="chain.search.media.site_page", + ) + async with aclosing(page_iterator): + async for site, search_page, result, continued in page_iterator: + finish_count += 1 + results.extend(result) + if not continued: + logger.debug( + f"{site.get('name')} 第 {search_page} 页返回 {len(result)} 条,停止继续翻页" + ) + logger.info(f"站点搜索进度:{finish_count} / {total_num}") + await progress.update( + value=finish_count / total_num * 100, + text=( + f"正在搜索{keyword or ''},已完成 " + f"{finish_count} / {total_num} 个请求 ..." + ), ) - for future in done_tasks: - site, page_index, search_page, search_keyword = pending_tasks.pop(future) - finish_count += 1 - result = await future - if result: - results.extend(result) - if ( - self._should_continue_search_pages( - site=site, page_results=result, keyword=search_keyword - ) - and page_index + 1 < len(search_pages) - ): - submit_site_page(site=site, page_index=page_index + 1) - else: - logger.debug( - f"{site.get('name')} 第 {search_page} 页返回 {len(result or [])} 条,停止继续翻页" - ) - logger.info(f"站点搜索进度:{finish_count} / {total_num}") - await progress.update(value=finish_count / total_num * 100, - text=f"正在搜索{keyword or ''},已完成 {finish_count} / {total_num} 个请求 ...") - finally: - for task in pending_tasks: - if not task.done(): - task.cancel() - if pending_tasks: - await asyncio.gather(*pending_tasks.keys(), return_exceptions=True) # 计算耗时 end_time = datetime.now() @@ -2630,89 +2675,69 @@ class SearchChain(ChainBase): "total": total_num } - semaphore = asyncio.Semaphore( - self.runtime_config.search_threadpool_size or total_num - ) - async def search_site(site: dict, search_page: int) -> List[TorrentInfo]: - """ - 搜索单个站点页,用于渐进式返回入口。 - """ - async with semaphore: - if area == "imdbid": - site_result = await self.async_search_site_torrents(site=site, - keyword=mediainfo.imdb_id if mediainfo else None, - mtype=mediainfo.type if mediainfo else mtype, - page=search_page) - else: - site_result = await self.async_search_site_torrents(site=site, - keyword=keyword, - mtype=mediainfo.type if mediainfo else mtype, - page=search_page) - return site_result or [] + """调用既有站点资源接口,具体并发由统一逐页编排器控制。""" + search_keyword = ( + mediainfo.imdb_id + if area == "imdbid" and mediainfo + else keyword + ) + return await self.async_search_site_torrents( + site=site, + keyword=search_keyword, + mtype=mediainfo.type if mediainfo else mtype, + page=search_page, + ) - tasks = {} - - def submit_site_page(site: dict, page_index: int): - """ - 提交渐进式站点页搜索任务,并保留站点和页码上下文。 - """ - search_page = search_pages[page_index] - search_keyword = mediainfo.imdb_id if area == "imdbid" and mediainfo else keyword - task = asyncio.create_task(search_site(site=site, search_page=search_page)) - tasks[task] = (site, page_index, search_page, search_keyword) - - for site in indexer_sites: - submit_site_page(site=site, page_index=0) + def should_continue(site: dict, page_results: List[Any]) -> bool: + """按既有站点分页规则判断是否提交下一页资源请求。""" + search_keyword = ( + mediainfo.imdb_id + if area == "imdbid" and mediainfo + else keyword + ) + return self._should_continue_search_pages( + site=site, + page_results=page_results, + keyword=search_keyword, + ) results_count = len(plugin_results) - try: - while tasks: - if global_vars.is_system_stopped: - break - done_tasks, _ = await asyncio.wait( - tasks.keys(), - return_when=asyncio.FIRST_COMPLETED, + page_iterator = self._iter_site_page_results( + indexer_sites=indexer_sites, + search_pages=search_pages, + search_page=search_site, + should_continue=should_continue, + task_owner="chain.search.media.site_page", + ) + async with aclosing(page_iterator): + async for site, search_page, result, continued in page_iterator: + finish_count += 1 + results_count += len(result) + if not continued: + logger.debug( + f"{site.get('name')} 第 {search_page} 页返回 {len(result)} 条,停止继续翻页" + ) + logger.info(f"站点搜索进度:{finish_count} / {total_num}") + progress_value = finish_count / total_num * 100 + progress_text = ( + f"正在搜索{keyword or ''},已完成 " + f"{finish_count} / {total_num} 个请求 ..." ) - for future in done_tasks: - site, page_index, search_page, search_keyword = tasks.pop(future) - finish_count += 1 - result = await future - results_count += len(result) - if ( - self._should_continue_search_pages( - site=site, page_results=result, keyword=search_keyword - ) - and page_index + 1 < len(search_pages) - ): - submit_site_page(site=site, page_index=page_index + 1) - else: - logger.debug( - f"{site.get('name')} 第 {search_page} 页返回 {len(result)} 条,停止继续翻页" - ) - logger.info(f"站点搜索进度:{finish_count} / {total_num}") - progress_value = finish_count / total_num * 100 - progress_text = f"正在搜索{keyword or ''},已完成 {finish_count} / {total_num} 个请求 ..." - await progress.update(value=progress_value, text=progress_text) - yield { - "type": "append", - "stage": "searching", - "value": progress_value, - "text": progress_text, - "items": result, - "site": site.get("name"), - "site_id": site.get("id"), - "page": search_page, - "finished": finish_count, - "total": total_num, - "total_items": results_count - } - finally: - for task in tasks: - if not task.done(): - task.cancel() - if tasks: - await asyncio.gather(*tasks.keys(), return_exceptions=True) + await progress.update(value=progress_value, text=progress_text) + yield { + "type": "append", + "stage": "searching", + "value": progress_value, + "text": progress_text, + "items": result, + "site": site.get("name"), + "site_id": site.get("id"), + "page": search_page, + "finished": finish_count, + "total": total_num, + "total_items": results_count + } end_time = datetime.now() await progress.update(value=100, @@ -2754,65 +2779,46 @@ class SearchChain(ChainBase): await progress.update(value=0, text=f"开始搜索字幕,共 {len(indexer_sites)} 个站点,{len(search_pages)} 页 ...") results = [] - semaphore = asyncio.Semaphore( - self.runtime_config.search_threadpool_size or total_num - ) async def search_site_page(site: dict, search_page: int) -> List[SubtitleInfo]: - """ - 控制单次字幕站点页请求的并发量,并返回该页的字幕列表。 - """ - async with semaphore: - return await self.async_search_subtitles( - site=site, keyword=keyword, page=search_page + """调用既有字幕接口,具体并发由统一逐页编排器控制。""" + return await self.async_search_subtitles( + site=site, + keyword=keyword, + page=search_page, + ) + + def should_continue(site: dict, page_results: List[Any]) -> bool: + """按既有字幕分页规则判断是否提交下一页请求。""" + return self._should_continue_subtitle_search_pages( + site=site, + page_results=page_results, + ) + + page_iterator = self._iter_site_page_results( + indexer_sites=indexer_sites, + search_pages=search_pages, + search_page=search_site_page, + should_continue=should_continue, + task_owner="chain.search.subtitle.site_page", + ) + async with aclosing(page_iterator): + async for site, search_page, result, continued in page_iterator: + finish_count += 1 + results.extend(result) + if not continued: + logger.debug( + f"{site.get('name')} 字幕第 {search_page} 页返回 {len(result)} 条,停止继续翻页" + ) + logger.info(f"站点字幕搜索进度:{finish_count} / {total_num}") + await progress.update( + value=finish_count / total_num * 100, + text=( + f"正在搜索字幕{keyword or ''},已完成 " + f"{finish_count} / {total_num} 个请求 ..." + ), ) - pending_tasks = {} - - def submit_site_page(site: dict, page_index: int): - """ - 提交异步字幕站点页搜索任务,并记录站点和页码位置。 - """ - search_page = search_pages[page_index] - task = asyncio.create_task(search_site_page(site=site, search_page=search_page)) - pending_tasks[task] = (site, page_index, search_page) - - for site in indexer_sites: - submit_site_page(site=site, page_index=0) - - try: - while pending_tasks: - if global_vars.is_system_stopped: - break - done_tasks, _ = await asyncio.wait( - pending_tasks.keys(), - return_when=asyncio.FIRST_COMPLETED, - ) - for future in done_tasks: - site, page_index, search_page = pending_tasks.pop(future) - finish_count += 1 - result = await future - if result: - results.extend(result) - if ( - self._should_continue_subtitle_search_pages(site=site, page_results=result) - and page_index + 1 < len(search_pages) - ): - submit_site_page(site=site, page_index=page_index + 1) - else: - logger.debug( - f"{site.get('name')} 字幕第 {search_page} 页返回 {len(result or [])} 条,停止继续翻页" - ) - logger.info(f"站点字幕搜索进度:{finish_count} / {total_num}") - await progress.update(value=finish_count / total_num * 100, - text=f"正在搜索字幕{keyword or ''},已完成 {finish_count} / {total_num} 个请求 ...") - finally: - for task in pending_tasks: - if not task.done(): - task.cancel() - if pending_tasks: - await asyncio.gather(*pending_tasks.keys(), return_exceptions=True) - end_time = datetime.now() await progress.update(value=100, text=f"站点字幕搜索完成,有效字幕数:{len(results)},总耗时 {(end_time - start_time).seconds} 秒") @@ -2871,79 +2877,57 @@ class SearchChain(ChainBase): "total": total_num } - semaphore = asyncio.Semaphore( - self.runtime_config.search_threadpool_size or total_num - ) - async def search_site(site: dict, search_page: int) -> List[SubtitleInfo]: - """ - 搜索单个站点字幕页,用于渐进式返回入口。 - """ - async with semaphore: - site_result = await self.async_search_subtitles( - site=site, keyword=keyword, page=search_page - ) - return site_result or [] + """调用既有字幕接口,具体并发由统一逐页编排器控制。""" + return await self.async_search_subtitles( + site=site, + keyword=keyword, + page=search_page, + ) - tasks = {} - - def submit_site_page(site: dict, page_index: int): - """ - 提交渐进式字幕站点页搜索任务,并保留站点和页码上下文。 - """ - search_page = search_pages[page_index] - task = asyncio.create_task(search_site(site=site, search_page=search_page)) - tasks[task] = (site, page_index, search_page) - - for site in indexer_sites: - submit_site_page(site=site, page_index=0) + def should_continue(site: dict, page_results: List[Any]) -> bool: + """按既有字幕分页规则判断是否提交下一页请求。""" + return self._should_continue_subtitle_search_pages( + site=site, + page_results=page_results, + ) results_count = 0 - try: - while tasks: - if global_vars.is_system_stopped: - break - done_tasks, _ = await asyncio.wait( - tasks.keys(), - return_when=asyncio.FIRST_COMPLETED, + page_iterator = self._iter_site_page_results( + indexer_sites=indexer_sites, + search_pages=search_pages, + search_page=search_site, + should_continue=should_continue, + task_owner="chain.search.subtitle.site_page", + ) + async with aclosing(page_iterator): + async for site, search_page, result, continued in page_iterator: + finish_count += 1 + results_count += len(result) + if not continued: + logger.debug( + f"{site.get('name')} 字幕第 {search_page} 页返回 {len(result)} 条,停止继续翻页" + ) + logger.info(f"站点字幕搜索进度:{finish_count} / {total_num}") + progress_value = finish_count / total_num * 100 + progress_text = ( + f"正在搜索字幕{keyword or ''},已完成 " + f"{finish_count} / {total_num} 个请求 ..." ) - for future in done_tasks: - site, page_index, search_page = tasks.pop(future) - finish_count += 1 - result = await future - results_count += len(result) - if ( - self._should_continue_subtitle_search_pages(site=site, page_results=result) - and page_index + 1 < len(search_pages) - ): - submit_site_page(site=site, page_index=page_index + 1) - else: - logger.debug( - f"{site.get('name')} 字幕第 {search_page} 页返回 {len(result)} 条,停止继续翻页" - ) - logger.info(f"站点字幕搜索进度:{finish_count} / {total_num}") - progress_value = finish_count / total_num * 100 - progress_text = f"正在搜索字幕{keyword or ''},已完成 {finish_count} / {total_num} 个请求 ..." - await progress.update(value=progress_value, text=progress_text) - yield { - "type": "append", - "stage": "searching", - "value": progress_value, - "text": progress_text, - "items": result, - "site": site.get("name"), - "site_id": site.get("id"), - "page": search_page, - "finished": finish_count, - "total": total_num, - "total_items": results_count - } - finally: - for task in tasks: - if not task.done(): - task.cancel() - if tasks: - await asyncio.gather(*tasks.keys(), return_exceptions=True) + await progress.update(value=progress_value, text=progress_text) + yield { + "type": "append", + "stage": "searching", + "value": progress_value, + "text": progress_text, + "items": result, + "site": site.get("name"), + "site_id": site.get("id"), + "page": search_page, + "finished": finish_count, + "total": total_num, + "total_items": results_count + } end_time = datetime.now() await progress.update(value=100, diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 95df2d1fb..3944320a3 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -6,7 +6,7 @@ > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`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 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解。 +> 实施进度:阶段 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 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排。 ## 当前复核结论(2026-08-24) @@ -472,6 +472,17 @@ - `pytest.ini` 将未处理线程异常提升为测试失败,后续不会再以绿色结果掩盖后台线程越过事件循环生命周期。 本阶段不改变 WebAgent API/SSE、Agent 模块、SDK/Compat 或 V1/V2/V3 插件合同,也未修改插件仓。 +### 长期整改阶段 45:影视与字幕搜索逐页任务编排统一(2026-08-24) + +- 影视普通搜索、影视流式搜索、字幕普通搜索和字幕流式搜索原先各自维护一套站点 task 字典、并发信号量、 + 续页判断以及取消 `gather`;四套近似实现使断线和异常收尾修复容易只覆盖部分入口。 +- `SearchChain._iter_site_page_results()` 现在统一持有请求级站点 task,使用稳定 owner 名称,集中执行并发 + 限制、逐站点续页、系统停止检查和 `cancel + gather` 终态等待;四个既有入口只保留各自的结果累计、进度 + 文案和 SSE 包装。外层使用 `aclosing`,消费者提前关闭流时会立即触发统一收尾。 +- 回归测试覆盖迭代器提前关闭时取消并等待其他站点请求,以及字幕普通/流式入口的页序、续页条件和结果 + 一致性。公开搜索方法、SSE 字段、插件资源源、站点模块方法、SDK/Compat 与 V1/V2/V3 插件合同均未改变; + 本阶段未修改插件仓。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/tests/test_search_page_orchestration.py b/tests/test_search_page_orchestration.py new file mode 100644 index 000000000..4288c0e9c --- /dev/null +++ b/tests/test_search_page_orchestration.py @@ -0,0 +1,124 @@ +"""站点资源与字幕搜索的统一逐页任务编排测试。""" + +import asyncio +from types import SimpleNamespace + +import pytest + +import app.chain.search as search_module +from app.chain.search import SearchChain + + +class _Progress: + """提供不访问 Redis 的异步搜索进度替身。""" + + def __init__(self, *_args, **_kwargs) -> None: + """接受生产构造参数但不创建外部资源。""" + + async def start(self) -> None: + """模拟开始进度。""" + + async def update(self, **_kwargs) -> None: + """模拟更新进度。""" + + async def end(self) -> None: + """模拟结束进度。""" + + +def _make_chain() -> SearchChain: + """构造只包含逐页搜索所需运行时配置的 SearchChain。""" + chain = object.__new__(SearchChain) + chain.runtime_config = SimpleNamespace(search_threadpool_size=2) + return chain + + +@pytest.mark.asyncio +async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None: + """调用方提前关闭迭代器时,统一编排器必须收口其他站点请求。""" + chain = _make_chain() + blocked_started = asyncio.Event() + blocked_cancelled = asyncio.Event() + task_names: list[str] = [] + + async def search_page(site: dict, _page: int) -> list[str]: + """让一个站点立即完成,另一个停留到被编排器取消。""" + task = asyncio.current_task() + task_names.append(task.get_name() if task else "") + if site["id"] == 1: + return ["ready"] + blocked_started.set() + try: + await asyncio.Event().wait() + except asyncio.CancelledError: + blocked_cancelled.set() + raise + + iterator = chain._iter_site_page_results( + indexer_sites=[{"id": 1}, {"id": 2}], + search_pages=[0], + search_page=search_page, + should_continue=lambda _site, _results: False, + task_owner="test.search.site_page", + ) + + assert await anext(iterator) == ({"id": 1}, 0, ["ready"], False) + await blocked_started.wait() + await iterator.aclose() + + assert blocked_cancelled.is_set() + assert task_names == ["test.search.site_page", "test.search.site_page"] + + +@pytest.mark.asyncio +async def test_subtitle_list_and_stream_share_page_continuation(monkeypatch) -> None: + """字幕列表与流式入口应复用同一续页规则、页序和停止条件。""" + chain = _make_chain() + site = {"id": 7, "name": "SubtitleSite", "subtitles": {"page_size": 2}} + calls: list[int] = [] + + class _Sites: + """返回固定字幕站点的异步索引替身。""" + + async def async_get_indexers(self) -> list[dict]: + """返回唯一启用的字幕站点。""" + return [site] + + async def search_subtitles(*, page: int, **_kwargs) -> list[str]: + """第一页满页触发续页,第二页不足一页后停止。""" + calls.append(page) + return ["one", "two"] if page == 0 else ["three"] + + monkeypatch.setattr(search_module, "SitesHelper", _Sites) + monkeypatch.setattr(search_module, "AsyncProgressHelper", _Progress) + monkeypatch.setattr( + search_module, + "get_configured_system_config", + lambda: SimpleNamespace(get=lambda _key: [7]), + ) + chain._build_search_pages = lambda _page: [0, 1, 2] + chain._should_continue_subtitle_search_pages = ( + lambda *, site, page_results: len(page_results) == 2 + ) + chain.async_search_subtitles = search_subtitles + + listed = await chain._SearchChain__async_search_subtitles_all_sites( + keyword="demo" + ) + assert listed == ["one", "two", "three"] + assert calls == [0, 1] + + calls.clear() + events = [ + event + async for event in chain._SearchChain__async_search_subtitles_all_sites_stream( + keyword="demo" + ) + ] + appended = [event for event in events if event["type"] == "append"] + + assert [event["page"] for event in appended] == [0, 1] + assert [event["items"] for event in appended] == [ + ["one", "two"], + ["three"], + ] + assert calls == [0, 1]