test(subscribe): validate joint search governance

This commit is contained in:
jxxghp
2026-09-01 17:09:54 +08:00
parent bab1d2862e
commit 9c4598c86b
5 changed files with 162 additions and 4 deletions
+7
View File
@@ -105,6 +105,9 @@ class SubscriptionSearchOper(DbOper):
if not isinstance(self._db, Session): if not isinstance(self._db, Session):
raise RuntimeError("订阅搜索认领需要调用方提供同步 Session") raise RuntimeError("订阅搜索认领需要调用方提供同步 Session")
now = utc_now_text() now = utc_now_text()
fairness_before = (
datetime.now(timezone.utc) - timedelta(minutes=15)
).isoformat(timespec="seconds")
lease_expires_at = ( lease_expires_at = (
datetime.now(timezone.utc) + timedelta(seconds=max(1, lease_seconds)) datetime.now(timezone.utc) + timedelta(seconds=max(1, lease_seconds))
).isoformat(timespec="seconds") ).isoformat(timespec="seconds")
@@ -129,6 +132,10 @@ class SubscriptionSearchOper(DbOper):
), ),
) )
.order_by( .order_by(
case(
(SubscriptionSearchTask.created_at <= fairness_before, 1),
else_=0,
).desc(),
SubscriptionSearchTask.priority.desc(), SubscriptionSearchTask.priority.desc(),
SubscriptionSearchTask.available_at.asc(), SubscriptionSearchTask.available_at.asc(),
SubscriptionSearchTask.created_at.asc(), SubscriptionSearchTask.created_at.asc(),
@@ -1,7 +1,7 @@
# MoviePilot 订阅执行治理 # MoviePilot 订阅执行治理
> 状态:`active2026-09-01 已由 MoviePilot v3 接管)` > 状态:`active2026-09-01 已由 MoviePilot v3 接管)`
> 当前叶:`SUB-GOV-002C` > 当前叶:`SUB-GOV-003A`
> 迁移来源:`V3-RDY-009A`、`V3-RDY-009A1`、`V3-RDY-009A1B`、 > 迁移来源:`V3-RDY-009A`、`V3-RDY-009A1`、`V3-RDY-009A1B`、
> `V3-RDY-009A1C`、`V3-RDY-009A1D` > `V3-RDY-009A1C`、`V3-RDY-009A1D`
> 适用范围:V3 订阅刷新、匹配、兜底搜索、下载提交与用户状态闭环 > 适用范围:V3 订阅刷新、匹配、兜底搜索、下载提交与用户状态闭环
@@ -234,12 +234,12 @@ RSS/Spider 保持现有逐站点串行刷新,不通过并发化换取几秒收
| `SUB-GOV-001E` | 若 001B–D 后仍有可重复等待热点,再通过隔离压测评估 2/4 worker 的有界准备;提交保持串行 | 001D | `not_activated2026-09-01` | | `SUB-GOV-001E` | 若 001B–D 后仍有可重复等待热点,再通过隔离压测评估 2/4 worker 的有界准备;提交保持串行 | 001D | `not_activated2026-09-01` |
| `SUB-GOV-002A` | 将 24 小时兜底搜索移出 Match/提交长锁,建立批次、订阅 single-flight、可取消等待和恢复游标 | 001C | `completed2026-09-01` | | `SUB-GOV-002A` | 将 24 小时兜底搜索移出 Match/提交长锁,建立批次、订阅 single-flight、可取消等待和恢复游标 | 001C | `completed2026-09-01` |
| `SUB-GOV-002B` | 为兜底搜索建立每站点并发、间隔和错误冷却预算;RSS/Spider 保持基线压力 | 002A | `completed2026-09-01` | | `SUB-GOV-002B` | 为兜底搜索建立每站点并发、间隔和错误冷却预算;RSS/Spider 保持基线压力 | 002A | `completed2026-09-01` |
| `SUB-GOV-002C` | 联合运行日常 Match、手工搜索和兜底搜索,验证公平性、锁等待、站点压力和失败恢复 | 001D, 002B | `in_progress` | | `SUB-GOV-002C` | 联合运行日常 Match、手工搜索和兜底搜索,验证公平性、锁等待、站点压力和失败恢复 | 001D, 002B | `completed2026-09-01` |
| `SUB-GOV-003A` | 建立跨入口的订阅级下载幂等、下载器不确定终态和取消补偿 | 001C, 002A | `pending` | | `SUB-GOV-003A` | 建立跨入口的订阅级下载幂等、下载器不确定终态和取消补偿 | 001C, 002A | `in_progress` |
| `SUB-GOV-003B` | 后端提供订阅及批次业务状态,前端展示排队、匹配、搜索、提交、完成、失败和取消 | 002A, 003A | `pending` | | `SUB-GOV-003B` | 后端提供订阅及批次业务状态,前端展示排队、匹配、搜索、提交、完成、失败和取消 | 002A, 003A | `pending` |
| `SUB-GOV-004` | 多条订阅记录指向同一媒体时的跨记录季集去重和产品规则 | 003A | `pending(最低优先级)` | | `SUB-GOV-004` | 多条订阅记录指向同一媒体时的跨记录季集去重和产品规则 | 003A | `pending(最低优先级)` |
当前只激活 `SUB-GOV-002C`。001BD 是日常 Match 主线;002A–C 是 24 小时兜底搜索主线;003A/B 当前只激活 `SUB-GOV-003A`。001BD 是日常 Match 主线;002A–C 是 24 小时兜底搜索主线;003A/B
只有在前两条链路的身份和终态稳定后实施。每个叶子完成验收并更新本表后,才激活下一个满足依赖的叶子。 只有在前两条链路的身份和终态稳定后实施。每个叶子完成验收并更新本表后,才激活下一个满足依赖的叶子。
### 6.1 SUB-GOV-001A 验收证据 ### 6.1 SUB-GOV-001A 验收证据
@@ -330,6 +330,19 @@ RSS/Spider 保持现有逐站点串行刷新,不通过并发化换取几秒收
`37 passed, 1 skipped`,架构合同组合 `149 passed`,错误级 Pylint 为 0Alembic 唯一 head 为 `37 passed, 1 skipped`,架构合同组合 `149 passed`,错误级 Pylint 为 0Alembic 唯一 head 为
`d2a7c5e9f1b4`Host 架构基线与 `git diff --check` 通过。 `d2a7c5e9f1b4`Host 架构基线与 `git diff --check` 通过。
### 6.8 SUB-GOV-002C 验收证据
- 队列认领增加 15 分钟持久老化窗口:新手工任务继续优先,但等待超过窗口的 fallback 任务先于后续高优先级
新任务执行,防止持续手工流量造成后台饥饿;
- 正式搜索队列继续完全不访问 Match 类级锁;联合测试在同一 provider fan-out 中证明冷却站点返回可聚合失败,
健康站点仍独立完成并释放预算,慢站点不拖住整轮;
- 同站点第二个 worker 无法认领,另一站点可同时认领;任务租约与站点租约均验证过期后以新 token 恢复,
单订阅失败继续后续任务且批次保持聚合 failed;
- 手工入口仅改变队列优先级和启动时机,不改变站点预算、失败冷却或 single-flightRSS/Spider 仍未接入该预算,
因此其请求模型保持基线;
- 验证:联合队列/站点预算/搜索状态 `35 passed`,订阅/搜索/调度架构 `58 passed`,错误级 Pylint 为 0
Host 架构基线与 `git diff --check` 通过。
## 7. 上线前验证与验收 ## 7. 上线前验证与验收
### 7.1 场景 ### 7.1 场景
@@ -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
+20
View File
@@ -130,3 +130,23 @@ def test_search_queue_finishes_batch_with_aggregated_failure(tmp_path):
assert batch.finished_count == 1 assert batch.finished_count == 1
assert batch.failed_count == 1 assert batch.failed_count == 1
assert batch.last_error == "site timeout" 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"
+19
View File
@@ -44,6 +44,25 @@ def test_site_budget_allows_one_inflight_per_site_and_independent_sites(tmp_path
assert other_site.acquired is True 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): def test_site_budget_applies_error_cooldown_and_gradual_success_recovery(tmp_path):
"""失败增加冷却计数,后续成功每次只恢复一级而非直接清零。""" """失败增加冷却计数,后续成功每次只恢复一级而非直接清零。"""
repository, engine = _repository(tmp_path) repository, engine = _repository(tmp_path)