refactor: keep subscription search within architecture budget

This commit is contained in:
jxxghp
2026-09-04 22:15:42 +08:00
parent 5afb1fc03e
commit 54557a991f
2 changed files with 59 additions and 38 deletions
+48 -1
View File
@@ -6,7 +6,12 @@ from dataclasses import dataclass
from typing import Callable, Mapping, Optional, Protocol from typing import Callable, Mapping, Optional, Protocol
from uuid import uuid4 from uuid import uuid4
from app.application.subscription.sitebudget import SiteBudgetClaim from app.application.subscription.sitebudget import (
SiteBudgetClaim,
SubscriptionSearchDeferred,
SubscriptionSiteBudgetDeferral,
)
from app.runtime.log import logger
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
@@ -96,6 +101,48 @@ class SubscriptionExecutionContext:
self.report_phase("submitting") self.report_phase("submitting")
def raise_subscription_site_budget_failures(failures: tuple[str, ...]) -> None:
"""在成功站点结果完成处理后暴露其余站点的聚合失败。"""
if failures:
raise RuntimeError("".join(failures))
def raise_subscription_site_budget_deferral(
deferrals: tuple[SubscriptionSiteBudgetDeferral, ...],
execution_context: Optional[SubscriptionExecutionContext],
) -> None:
"""在没有下载副作用时,将临时站点冲突转换为持久队列延后。"""
if not deferrals or (execution_context and execution_context.download_started):
return
retry_at = min(deferrals, key=lambda item: item.retry_at).retry_at
site_ids = tuple(dict.fromkeys(item.site_id for item in deferrals))
raise SubscriptionSearchDeferred(retry_at=retry_at, site_ids=site_ids)
def handle_subscription_search_deferred(
queue: SubscriptionSearchRepository,
task_id: str,
lease_token: str,
subscribe_name: str,
deferred: SubscriptionSearchDeferred,
record: Callable[[str, Optional[str]], None],
) -> None:
"""把站点预算冲突重新入队,并记录为可恢复而非失败的任务结果。"""
requeued = queue.defer_task(
task_id=task_id,
lease_token=lease_token,
available_at=deferred.retry_at,
)
if requeued:
logger.debug(
f"订阅 {subscribe_name} 站点预算冲突,已排队至 {deferred.retry_at} 后重试,"
f"sites={','.join(str(site_id) for site_id in deferred.site_ids)}"
)
record("requeued", "site_budget_deferred")
else:
logger.debug(f"订阅搜索任务 {task_id} 租约已变化,跳过重复站点预算重排队")
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class SearchBatchSnapshot: class SearchBatchSnapshot:
"""订阅搜索批次的持久业务状态快照。""" """订阅搜索批次的持久业务状态快照。"""
+11 -37
View File
@@ -19,6 +19,9 @@ from app.application.subscription.execution import (
SearchTaskSnapshot, SearchTaskSnapshot,
SubscriptionExecutionContext, SubscriptionExecutionContext,
SubscriptionSearchRepository, SubscriptionSearchRepository,
handle_subscription_search_deferred,
raise_subscription_site_budget_deferral,
raise_subscription_site_budget_failures,
) )
from app.application.subscription.observability import ( from app.application.subscription.observability import (
SearchExecutionSummary, SearchExecutionSummary,
@@ -33,7 +36,6 @@ from app.application.subscription.sitebudget import (
SubscriptionSearchCancelled, SubscriptionSearchCancelled,
SubscriptionSearchDeferred, SubscriptionSearchDeferred,
SubscriptionSiteBudget, SubscriptionSiteBudget,
SubscriptionSiteBudgetDeferral,
) )
from app.chain.media import MediaChain from app.chain.media import MediaChain
from app.chain.search.facade import SearchChain from app.chain.search.facade import SearchChain
@@ -619,19 +621,9 @@ class _SubscribeSearchQueueOwner(_SubscribeSearchQueueCoordinator):
"system_stop" if system_stopped else "cancelled", "system_stop" if system_stopped else "cancelled",
) )
except SubscriptionSearchDeferred as deferred: except SubscriptionSearchDeferred as deferred:
requeued = queue.defer_task( handle_subscription_search_deferred(
task_id=task_id, queue, task_id, lease_token, subscribe.name, deferred, summary.record
lease_token=lease_token,
available_at=deferred.retry_at,
) )
if requeued:
logger.debug(
f"订阅 {subscribe.name} 站点预算冲突,已排队至 {deferred.retry_at} 后重试,"
f"sites={','.join(str(site_id) for site_id in deferred.site_ids)}"
)
summary.record("requeued", "site_budget_deferred")
else:
logger.debug(f"订阅搜索任务 {task_id} 租约已变化,跳过重复站点预算重排队")
except Exception as err: except Exception as err:
logger.error(f"订阅 {subscribe.name} 搜索失败:{str(err)}", exc_info=True) logger.error(f"订阅 {subscribe.name} 搜索失败:{str(err)}", exc_info=True)
queue.finish_task( queue.finish_task(
@@ -848,22 +840,22 @@ class SubscribeSearchOwner(_SubscribeSearchQueueOwner):
if not contexts: if not contexts:
logger.debug(f"订阅 {subscribe.keyword or subscribe.name} 未搜索到资源") logger.debug(f"订阅 {subscribe.keyword or subscribe.name} 未搜索到资源")
if not site_budget_failures: if not site_budget_failures:
self._raise_site_budget_deferral(site_budget_deferrals, execution_context) raise_subscription_site_budget_deferral(site_budget_deferrals, execution_context)
self.finish_subscribe_or_not( self.finish_subscribe_or_not(
subscribe=subscribe, subscribe=subscribe,
meta=meta, meta=meta,
mediainfo=mediainfo, mediainfo=mediainfo,
lefts=no_exists, lefts=no_exists,
) )
self._raise_site_budget_failures(site_budget_failures) raise_subscription_site_budget_failures(site_budget_failures)
return subscribe return subscribe
matched = self._filter_search_contexts(subscribe, contexts) matched = self._filter_search_contexts(subscribe, contexts)
if not matched: if not matched:
logger.debug(f"订阅 {subscribe.name} 没有符合过滤条件的资源") logger.debug(f"订阅 {subscribe.name} 没有符合过滤条件的资源")
if not site_budget_failures: if not site_budget_failures:
self._raise_site_budget_deferral(site_budget_deferrals, execution_context) raise_subscription_site_budget_deferral(site_budget_deferrals, execution_context)
self.finish_subscribe_or_not(subscribe=subscribe, meta=meta, mediainfo=mediainfo, lefts=no_exists) self.finish_subscribe_or_not(subscribe=subscribe, meta=meta, mediainfo=mediainfo, lefts=no_exists)
self._raise_site_budget_failures(site_budget_failures) raise_subscription_site_budget_failures(site_budget_failures)
return subscribe return subscribe
if execution_context: if execution_context:
execution_context.report_phase("preparing") execution_context.report_phase("preparing")
@@ -888,28 +880,10 @@ class SubscribeSearchOwner(_SubscribeSearchQueueOwner):
downloads=downloads, downloads=downloads,
lefts=lefts, lefts=lefts,
) )
self._raise_site_budget_failures(site_budget_failures) raise_subscription_site_budget_failures(site_budget_failures)
self._raise_site_budget_deferral(site_budget_deferrals, execution_context) raise_subscription_site_budget_deferral(site_budget_deferrals, execution_context)
return cast(Optional[SubscriptionSnapshot], current) return cast(Optional[SubscriptionSnapshot], current)
@staticmethod
def _raise_site_budget_failures(failures: tuple[str, ...]) -> None:
"""在成功站点结果完成处理后暴露其余站点的聚合失败。"""
if failures:
raise RuntimeError("".join(failures))
@staticmethod
def _raise_site_budget_deferral(
deferrals: tuple[SubscriptionSiteBudgetDeferral, ...],
execution_context: Optional[SubscriptionExecutionContext],
) -> None:
"""在没有下载副作用时,将临时站点冲突转换为持久队列延后。"""
if not deferrals or (execution_context and execution_context.download_started):
return
retry_at = min(deferrals, key=lambda item: item.retry_at).retry_at
site_ids = tuple(dict.fromkeys(item.site_id for item in deferrals))
raise SubscriptionSearchDeferred(retry_at=retry_at, site_ids=site_ids)
def _filter_search_contexts( def _filter_search_contexts(
self, self,
subscribe: SubscriptionSnapshot, subscribe: SubscriptionSnapshot,