diff --git a/app/db/oper/subscriptionsearch.py b/app/db/oper/subscriptionsearch.py index 2f3a37d91..1e780cd9e 100644 --- a/app/db/oper/subscriptionsearch.py +++ b/app/db/oper/subscriptionsearch.py @@ -105,6 +105,9 @@ class SubscriptionSearchOper(DbOper): if not isinstance(self._db, Session): raise RuntimeError("订阅搜索认领需要调用方提供同步 Session") now = utc_now_text() + fairness_before = ( + datetime.now(timezone.utc) - timedelta(minutes=15) + ).isoformat(timespec="seconds") lease_expires_at = ( datetime.now(timezone.utc) + timedelta(seconds=max(1, lease_seconds)) ).isoformat(timespec="seconds") @@ -129,6 +132,10 @@ class SubscriptionSearchOper(DbOper): ), ) .order_by( + case( + (SubscriptionSearchTask.created_at <= fairness_before, 1), + else_=0, + ).desc(), SubscriptionSearchTask.priority.desc(), SubscriptionSearchTask.available_at.asc(), SubscriptionSearchTask.created_at.asc(), diff --git a/docs/refactor/moviepilot-subscription-governance-roadmap.md b/docs/refactor/moviepilot-subscription-governance-roadmap.md index 71f475d71..0a9f9a001 100644 --- a/docs/refactor/moviepilot-subscription-governance-roadmap.md +++ b/docs/refactor/moviepilot-subscription-governance-roadmap.md @@ -1,7 +1,7 @@ # MoviePilot 订阅执行治理 > 状态:`active(2026-09-01 已由 MoviePilot v3 接管)` -> 当前叶:`SUB-GOV-002C` +> 当前叶:`SUB-GOV-003A` > 迁移来源:`V3-RDY-009A`、`V3-RDY-009A1`、`V3-RDY-009A1B`、 > `V3-RDY-009A1C`、`V3-RDY-009A1D` > 适用范围:V3 订阅刷新、匹配、兜底搜索、下载提交与用户状态闭环 @@ -234,12 +234,12 @@ RSS/Spider 保持现有逐站点串行刷新,不通过并发化换取几秒收 | `SUB-GOV-001E` | 若 001B–D 后仍有可重复等待热点,再通过隔离压测评估 2/4 worker 的有界准备;提交保持串行 | 001D | `not_activated(2026-09-01)` | | `SUB-GOV-002A` | 将 24 小时兜底搜索移出 Match/提交长锁,建立批次、订阅 single-flight、可取消等待和恢复游标 | 001C | `completed(2026-09-01)` | | `SUB-GOV-002B` | 为兜底搜索建立每站点并发、间隔和错误冷却预算;RSS/Spider 保持基线压力 | 002A | `completed(2026-09-01)` | -| `SUB-GOV-002C` | 联合运行日常 Match、手工搜索和兜底搜索,验证公平性、锁等待、站点压力和失败恢复 | 001D, 002B | `in_progress` | -| `SUB-GOV-003A` | 建立跨入口的订阅级下载幂等、下载器不确定终态和取消补偿 | 001C, 002A | `pending` | +| `SUB-GOV-002C` | 联合运行日常 Match、手工搜索和兜底搜索,验证公平性、锁等待、站点压力和失败恢复 | 001D, 002B | `completed(2026-09-01)` | +| `SUB-GOV-003A` | 建立跨入口的订阅级下载幂等、下载器不确定终态和取消补偿 | 001C, 002A | `in_progress` | | `SUB-GOV-003B` | 后端提供订阅及批次业务状态,前端展示排队、匹配、搜索、提交、完成、失败和取消 | 002A, 003A | `pending` | | `SUB-GOV-004` | 多条订阅记录指向同一媒体时的跨记录季集去重和产品规则 | 003A | `pending(最低优先级)` | -当前只激活 `SUB-GOV-002C`。001B–D 是日常 Match 主线;002A–C 是 24 小时兜底搜索主线;003A/B +当前只激活 `SUB-GOV-003A`。001B–D 是日常 Match 主线;002A–C 是 24 小时兜底搜索主线;003A/B 只有在前两条链路的身份和终态稳定后实施。每个叶子完成验收并更新本表后,才激活下一个满足依赖的叶子。 ### 6.1 SUB-GOV-001A 验收证据 @@ -330,6 +330,19 @@ RSS/Spider 保持现有逐站点串行刷新,不通过并发化换取几秒收 `37 passed, 1 skipped`,架构合同组合 `149 passed`,错误级 Pylint 为 0,Alembic 唯一 head 为 `d2a7c5e9f1b4`,Host 架构基线与 `git diff --check` 通过。 +### 6.8 SUB-GOV-002C 验收证据 + +- 队列认领增加 15 分钟持久老化窗口:新手工任务继续优先,但等待超过窗口的 fallback 任务先于后续高优先级 + 新任务执行,防止持续手工流量造成后台饥饿; +- 正式搜索队列继续完全不访问 Match 类级锁;联合测试在同一 provider fan-out 中证明冷却站点返回可聚合失败, + 健康站点仍独立完成并释放预算,慢站点不拖住整轮; +- 同站点第二个 worker 无法认领,另一站点可同时认领;任务租约与站点租约均验证过期后以新 token 恢复, + 单订阅失败继续后续任务且批次保持聚合 failed; +- 手工入口仅改变队列优先级和启动时机,不改变站点预算、失败冷却或 single-flight;RSS/Spider 仍未接入该预算, + 因此其请求模型保持基线; +- 验证:联合队列/站点预算/搜索状态 `35 passed`,订阅/搜索/调度架构 `58 passed`,错误级 Pylint 为 0, + Host 架构基线与 `git diff --check` 通过。 + ## 7. 上线前验证与验收 ### 7.1 场景 diff --git a/tests/test_subscription_governance_joint.py b/tests/test_subscription_governance_joint.py new file mode 100644 index 000000000..1d8db8e06 --- /dev/null +++ b/tests/test_subscription_governance_joint.py @@ -0,0 +1,99 @@ +"""日常 Match、手工/兜底队列和站点预算的联合治理测试。""" + +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace + +from app.application.site.search_observation import report_site_search_outcome +from app.application.subscription.sitebudget import SiteBudgetClaim, SubscriptionSiteBudget +from app.chain.search.facade import SearchChain +from app.chain.search.provider import SearchProviderOwner +from app.runtime.stop import ProcessStopState + + +class _Progress: + """收集 provider 进度而不依赖全局进度缓存。""" + + def __init__(self) -> None: + self.values: list[float] = [] + + def update(self, *, value: float, **_kwargs) -> None: + """记录一次进度值。""" + self.values.append(value) + + +class _MixedBudgetRepository: + """模拟一个冷却站点和一个立即可用站点。""" + + def __init__(self) -> None: + self.finished_sites: list[int] = [] + + def claim_site(self, *, site_id: int, owner: str, lease_seconds: int) -> SiteBudgetClaim: + """站点 1 保持冷却,站点 2 返回独占租约。""" + del owner, lease_seconds + now = datetime.now(timezone.utc) + if site_id == 1: + return SiteBudgetClaim( + site_id=site_id, + acquired=False, + retry_at=(now + timedelta(minutes=10)).isoformat(timespec="seconds"), + consecutive_failures=1, + ) + return SiteBudgetClaim( + site_id=site_id, + acquired=True, + retry_at=now.isoformat(timespec="seconds"), + consecutive_failures=0, + lease_token=f"lease-{site_id}", + ) + + def finish_site(self, *, site_id: int, **_kwargs) -> bool: + """记录成功执行并释放的站点。""" + self.finished_sites.append(site_id) + return True + + +def test_cooled_site_does_not_block_independent_site_or_hide_batch_failure(): + """慢站点跳过后其它站点仍完成,调用方同时收到可聚合失败。""" + repository = _MixedBudgetRepository() + budget = SubscriptionSiteBudget( + repository=repository, + owner="fallback-task", + cancelled=lambda: False, + stop_state=ProcessStopState(), + max_wait_seconds=0, + random_uniform=lambda _low, _high: 60.0, + ) + chain = object.__new__(SearchChain) + chain._runtime_config = SimpleNamespace(search_threadpool_size=2) + chain.configure_subscription_site_budget(budget) + chain._should_continue_search_pages = lambda **_kwargs: False + + def search_site_torrents(*, site, **_kwargs): + """模拟可用站点返回一条结果并发布成功事实。""" + report_site_search_outcome(attempted=True, outcome="success") + return [site["id"]] + + chain.search_site_torrents = search_site_torrents + results = [] + progress = _Progress() + + SearchProviderOwner._collect_sync_site_results( # pylint: disable=protected-access + chain, + keyword="movie", + indexer_sites=[ + {"id": 1, "name": "Cooling"}, + {"id": 2, "name": "Healthy"}, + ], + search_pages=[0], + search_keyword="movie", + media_type=None, + results=results, + progress=progress, + ) + + assert results == [2] + assert repository.finished_sites == [2] + failures = chain.consume_subscription_site_budget_failures() + assert len(failures) == 1 + assert "站点 1" in failures[0] + assert progress.values[-1] == 100 diff --git a/tests/test_subscription_search_queue.py b/tests/test_subscription_search_queue.py index 25757da61..615f65fd7 100644 --- a/tests/test_subscription_search_queue.py +++ b/tests/test_subscription_search_queue.py @@ -130,3 +130,23 @@ def test_search_queue_finishes_batch_with_aggregated_failure(tmp_path): assert batch.finished_count == 1 assert batch.failed_count == 1 assert batch.last_error == "site timeout" + + +def test_search_queue_ages_old_fallback_ahead_of_new_manual_work(tmp_path): + """手工任务可优先,但等待超过公平窗口的兜底任务不得持续饥饿。""" + repository, engine = _repository(tmp_path) + repository.enqueue(subscription_ids=(8,), source="fallback", priority=10) + aged_at = (datetime.now(timezone.utc) - timedelta(minutes=16)).isoformat(timespec="seconds") + with Session(engine) as session: + session.execute( + update(SubscriptionSearchTask) + .where(SubscriptionSearchTask.subscription_id == 8) + .values(created_at=aged_at) + ) + session.commit() + repository.enqueue(subscription_ids=(9,), source="manual", priority=100) + + claimed = repository.claim_next(owner="worker-a") + + assert claimed.subscription_id == 8 + assert claimed.source == "fallback" diff --git a/tests/test_subscription_site_budget.py b/tests/test_subscription_site_budget.py index 14510a048..3487d07aa 100644 --- a/tests/test_subscription_site_budget.py +++ b/tests/test_subscription_site_budget.py @@ -44,6 +44,25 @@ def test_site_budget_allows_one_inflight_per_site_and_independent_sites(tmp_path assert other_site.acquired is True +def test_site_budget_recovers_expired_inflight_lease(tmp_path): + """进程遗留的过期站点租约可被新 worker 以新 token 恢复。""" + repository, engine = _repository(tmp_path) + first = repository.claim_site(site_id=20, owner="worker-a", lease_seconds=900) + expired_at = (datetime.now(timezone.utc) - timedelta(seconds=1)).isoformat(timespec="seconds") + with Session(engine) as session: + session.execute( + SiteBudgetRecord.__table__.update() + .where(SiteBudgetRecord.site_id == 20) + .values(lease_expires_at=expired_at) + ) + session.commit() + + recovered = repository.claim_site(site_id=20, owner="worker-b", lease_seconds=900) + + assert recovered.acquired is True + assert recovered.lease_token != first.lease_token + + def test_site_budget_applies_error_cooldown_and_gradual_success_recovery(tmp_path): """失败增加冷却计数,后续成功每次只恢复一级而非直接清零。""" repository, engine = _repository(tmp_path)