Merge pull request #6591 from jxxghp/fix/issue-6588

feat(search): log site durations
This commit is contained in:
jxxghp
2026-09-06 16:30:08 +08:00
committed by GitHub
3 changed files with 89 additions and 13 deletions
+8 -6
View File
@@ -1,6 +1,7 @@
"""搜索分页计划与异步逐页调度 owner。"""
import asyncio
import time
from typing import Any, AsyncIterator, Awaitable, Callable, Dict, List, Optional, Tuple
from app.application.configuration import (
@@ -10,7 +11,7 @@ from app.chain.search.contract import _SearchOwnerBase
from app.runtime.stop import runtime_stop_state
from app.runtime.tasks import get_task_registry
PageResult = Tuple[List[Any], Optional[Exception]]
PageResult = Tuple[List[Any], Optional[Exception], float]
PageTask = asyncio.Task[PageResult]
SiteIndexer = Dict[str, Any]
PendingPages = dict[PageTask, Tuple[SiteIndexer, int, int]]
@@ -24,14 +25,15 @@ async def _run_site_page(
) -> PageResult:
"""执行一页请求,并把业务异常交回编排层统一传播。"""
async with semaphore:
started_at = time.perf_counter()
try:
return await search_page(site, page_number) or [], None
return await search_page(site, page_number) or [], None, time.perf_counter() - started_at
except asyncio.CancelledError:
raise
except Exception as error: # noqa: BLE001
# TaskRegistry 会报告未收口的后台异常;这里由等待方同步传播,
# 避免同一 provider 故障同时被登记器和调用链重复记录。
return [], error
return [], error, time.perf_counter() - started_at
def _submit_site_page(
@@ -126,7 +128,7 @@ class SearchPaginationOwner(_SearchOwnerBase):
search_page: Callable[[SiteIndexer, int], Awaitable[Optional[List[Any]]]],
should_continue: Callable[[SiteIndexer, List[Any]], bool],
task_owner: str,
) -> AsyncIterator[Tuple[SiteIndexer, int, List[Any], bool]]:
) -> AsyncIterator[Tuple[SiteIndexer, int, List[Any], bool, float]]:
"""统一调度站点逐页请求,并在调用方退出时取消、等待全部请求。"""
total_num = len(indexer_sites) * len(search_pages)
semaphore = asyncio.Semaphore(self.runtime_config.search_threadpool_size or max(1, total_num))
@@ -153,7 +155,7 @@ class SearchPaginationOwner(_SearchOwnerBase):
)
for task in done_tasks:
site, page_index, page_number = pending_tasks.pop(task)
page_results, error = await task
page_results, error, elapsed = await task
if error is not None:
raise error
continued = should_continue(site, page_results) and page_index + 1 < len(search_pages)
@@ -167,6 +169,6 @@ class SearchPaginationOwner(_SearchOwnerBase):
search_page=search_page,
task_owner=task_owner,
)
yield site, page_number, page_results, continued
yield site, page_number, page_results, continued, elapsed
finally:
await _cancel_pending_pages(pending_tasks)
+68 -6
View File
@@ -35,6 +35,27 @@ _site_request_schedule_lock = threading.Lock()
_site_next_request_at: Dict[str, float] = {}
def _site_duration_key(site: SiteIndexer) -> str:
"""返回站点耗时统计使用的稳定键。"""
site_id = site.get("id")
return str(site_id if site_id is not None else site.get("domain") or site.get("name") or "未知站点")
def _site_duration_label(site: SiteIndexer) -> str:
"""返回站点耗时日志中的可读名称。"""
return str(site.get("name") or site.get("domain") or site.get("id") or "未知站点")
def _format_site_durations(durations: Dict[str, tuple[str, float]]) -> str:
"""格式化各站点耗时,并标出本轮最慢站点。"""
if not durations:
return ""
ordered = sorted(durations.values(), key=lambda item: item[1], reverse=True)
details = "".join(f"{name}{seconds:.2f}" for name, seconds in ordered)
slowest_name, slowest_seconds = ordered[0]
return f";各站点耗时:{details};最慢站点:{slowest_name}{slowest_seconds:.2f} 秒)"
def _site_request_interval(site: SiteIndexer) -> float:
"""读取站点管理中已有的单次访问间隔配置。"""
try:
@@ -90,6 +111,25 @@ def _search_site_page(
)
def _timed_sync_site_page(
owner: "_SearchProviderSyncOwner",
*,
site: SiteIndexer,
keyword: str,
mtype: Optional[MediaType],
page: int,
) -> tuple[List[TorrentInfo], float]:
"""执行同步站点搜索并返回本页耗时,不包含线程池排队时间。"""
started_at = time.perf_counter()
result = owner._search_site_torrents_with_budget(
site=site,
keyword=keyword,
mtype=mtype,
page=page,
)
return result, time.perf_counter() - started_at
@dataclass(frozen=True)
class ProviderBatch:
"""一次 provider 页完成后发布的统一事实。"""
@@ -101,6 +141,7 @@ class ProviderBatch:
total: int
total_items: int
continued: bool
elapsed: float
class _SearchProviderSyncOwner(_SearchOwnerBase):
@@ -186,7 +227,8 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
"""向进程共享线程 owner 提交一页,并登记该站点的续页位置。"""
page_number = search_pages[page_index]
future = ThreadHelper().submit(
self._search_site_torrents_with_budget,
_timed_sync_site_page,
self,
site=site,
keyword=search_keyword,
mtype=media_type,
@@ -312,6 +354,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
self.runtime_config.search_threadpool_size or len(indexer_sites),
)
finish_count = 0
site_durations: Dict[str, tuple[str, float]] = {}
for _ in range(max_workers):
SearchProviderOwner._submit_next_sync_site(
@@ -332,7 +375,14 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
)
for future in done:
site, page_index, page_number = pending.pop(future)
page_results = future.result() or []
page_results, elapsed = future.result()
page_results = page_results or []
duration_key = _site_duration_key(site)
previous = site_durations.get(duration_key)
site_durations[duration_key] = (
_site_duration_label(site),
(previous[1] if previous else 0.0) + elapsed,
)
finish_count += 1
results.extend(page_results)
continued = self._should_continue_search_pages(
@@ -368,6 +418,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
finally:
for future in pending:
future.cancel()
return site_durations
def _search_all_sites(
self,
@@ -402,7 +453,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
value=0,
text=(f"开始搜索,共 {len(indexer_sites)} 个站点,{len(search_pages)} 页 ..."),
)
SearchProviderOwner._collect_sync_site_results(
site_durations = SearchProviderOwner._collect_sync_site_results(
self,
keyword=keyword,
indexer_sites=indexer_sites,
@@ -423,7 +474,10 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
if isinstance(getattr(self, "_subscription_site_budget", None), SubscriptionSiteBudget)
else logger.info
)
log(f"站点搜索完成,有效资源数:{len(results)},总耗时 {elapsed}")
log(
f"站点搜索完成,有效资源数:{len(results)},总耗时 {elapsed}"
f"{_format_site_durations(site_durations)}"
)
return results
finally:
progress.end()
@@ -457,7 +511,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
task_owner=task_owner,
)
async with aclosing(page_iterator):
async for site, page_number, page_results, continued in page_iterator:
async for site, page_number, page_results, continued, elapsed in page_iterator:
finish_count += 1
total_items += len(page_results)
yield ProviderBatch(
@@ -468,6 +522,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
total=total,
total_items=total_items,
continued=continued,
elapsed=elapsed,
)
@staticmethod
@@ -615,6 +670,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
started_at = datetime.now()
last_total_items = len(initial_items)
site_durations: Dict[str, tuple[str, float]] = {}
# 既有流式协议约定插件资源先于站点进度发布;列表入口也消费同一事件流。
if initial_items:
yield SearchProviderOwner._plugin_event(
@@ -646,6 +702,12 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
async with aclosing(batches):
async for batch in batches:
last_total_items = batch.total_items
duration_key = _site_duration_key(batch.site)
previous = site_durations.get(duration_key)
site_durations[duration_key] = (
_site_duration_label(batch.site),
(previous[1] if previous else 0.0) + batch.elapsed,
)
yield await SearchProviderOwner._batch_event(
batch=batch,
keyword=keyword,
@@ -658,7 +720,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
f"站点{label}搜索完成,有效{'字幕' if subtitle else '资源'}数:{last_total_items},总耗时 {elapsed}"
)
await progress.update(value=100, text=done_text)
logger.info(done_text)
logger.info(f"{done_text}{_format_site_durations(site_durations)}")
finally:
await progress.end()
+13 -1
View File
@@ -32,6 +32,16 @@ def _make_chain() -> SearchChain:
return chain
def test_format_site_durations_orders_sites_and_marks_slowest() -> None:
"""站点耗时日志应按耗时倒序,并明确标出最慢站点。"""
formatted = search_module._format_site_durations({
"1": ("站点一", 1.2),
"2": ("站点二", 3.4),
})
assert formatted == ";各站点耗时:站点二:3.40 秒,站点一:1.20 秒;最慢站点:站点二(3.40 秒)"
@pytest.mark.asyncio
async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None:
"""调用方提前关闭迭代器时,统一编排器必须收口其他站点请求。"""
@@ -61,7 +71,9 @@ async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None:
task_owner="test.search.site_page",
)
assert await anext(iterator) == ({"id": 1}, 0, ["ready"], False)
result = await anext(iterator)
assert result[:4] == ({"id": 1}, 0, ["ready"], False)
assert result[4] >= 0
await blocked_started.wait()
await iterator.aclose()