All checks were successful
Deploy Trade-In / changes (push) Successful in 13s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m14s
Deploy Trade-In / build-backend (push) Successful in 1m16s
Deploy Trade-In / deploy (push) Successful in 7m32s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 12s
1448 lines
75 KiB
Python
1448 lines
75 KiB
Python
"""Offline-тесты пула прокси (#2162, #2600).
|
||
|
||
Покрытие БЕЗ live-сети/БД: stateful FakeSession эмулирует таблицу scrape_proxies и
|
||
интерпретирует SQL по ключевым фрагментам, так что acquire/release/mark_health/
|
||
reap_stale_leases проверяются по фактическому изменению состояния строк.
|
||
|
||
- acquire: возвращает свободный+здоровый прокси нужного affinity; ставит lease.
|
||
- два acquire подряд → РАЗНЫЕ прокси (первый лизнут → выпал из выборки второго).
|
||
- release освобождает (leased_by → NULL), прокси снова acquire-абелен.
|
||
- mark_health fail → инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
|
||
- mark_health ok → сброс fails + exit_ip/latency + enabled=true (реанимация).
|
||
- reap_stale_leases освобождает старый lease, свежий не трогает.
|
||
- touch (#2164 sticky-session fix, 2026-08): heartbeat продлевает leased_at активного
|
||
lease, no-op на свободном узле; повторный touch перед reap не даёт reap_stale_leases
|
||
отобрать многочасовую browser-сессию, отсутствие touch — реапится как раньше.
|
||
- affinity-фильтр: acquire('avito') не берёт cian-only прокси.
|
||
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
|
||
- acquire без своих/any свободных → берёт свободный чужой affinity (fallback, #2600 п.3).
|
||
- run_proxy_healthcheck: reap + проба каждого enabled + mark_health (проба замокана).
|
||
- run_proxy_healthcheck: disabled-узлы — самовосстановление (#2600 п.1):
|
||
* успешная проба выключенного узла возвращает его в строй + revived++;
|
||
* недавно проверенный выключенный узел повторно не проверяется (не долбим провайдера).
|
||
- mark_health / disabled_reason (#2610 — ручное vs авто-выключение):
|
||
* авто-выключенный узел (disabled_reason IS NULL) по-прежнему воскресает
|
||
успешной пробой — поведение #2609 не сломано;
|
||
* ручно-выключенный узел (disabled_reason НЕ NULL) НЕ воскресает даже при
|
||
ok=True, и это логируется (WARNING);
|
||
* run_proxy_healthcheck не считает ручно-выключенный узел в revived.
|
||
- бан по паре «узел × источник» (#2600 п.2, scrape_proxy_source_bans):
|
||
* mark_banned пишет строку бана и НЕ выключает узел глобально;
|
||
* повторный бан той же пары эскалирует срок (ban_count растёт);
|
||
* acquire не выдаёт узел с активным баном по ЭТОМУ source, но выдаёт по ДРУГОМУ
|
||
(суть задачи: Авито забанил — Яндекс продолжает ходить через тот же узел);
|
||
* истёкший бан снова не мешает выдаче;
|
||
* защита последнего узла: бан НЕ записывается, если у acquire(source) не
|
||
останется кандидатов;
|
||
* run_proxy_healthcheck сносит бан-строки, истёкшие дольше SOURCE_BAN_PURGE_DAYS,
|
||
и НЕ трогает истёкшие недавно (в них живёт ban_count для эскалации);
|
||
* backup выделенной affinity засчитывается, только если он сам не забанен своим
|
||
источником (иначе fallback уводил бы последний рабочий узел, deep-review);
|
||
* clear_source_bans снимает баны узла (все / один source) и обнуляет эскалацию.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import os
|
||
|
||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||
|
||
from datetime import UTC, datetime, timedelta
|
||
from typing import Any
|
||
|
||
import pytest
|
||
|
||
from app.services import proxy_pool
|
||
from app.services.proxy_pool import (
|
||
DISABLE_THRESHOLD,
|
||
DISABLED_RECHECK_MINUTES,
|
||
MAX_CONSECUTIVE_FAILS,
|
||
SOURCE_BAN_BASE_HOURS,
|
||
SOURCE_BAN_MAX_HOURS,
|
||
SOURCE_BAN_PURGE_DAYS,
|
||
STALE_LEASE_MINUTES,
|
||
acquire,
|
||
mark_health,
|
||
reap_stale_leases,
|
||
release,
|
||
)
|
||
|
||
# mark_banned() не существовал до #2600 п.1 (2026-08) — тот же паттерн отложенного
|
||
# импорта, что touch() выше: AttributeError падает ТОЛЬКО в этих тестах, не роняя
|
||
# коллекцию всего файла на pre-fix коде.
|
||
try:
|
||
from app.services.proxy_pool import mark_banned
|
||
except ImportError: # pragma: no cover — pre-fix guard, см. комментарий выше
|
||
mark_banned = None # type: ignore[assignment]
|
||
|
||
# touch() не существовал до sticky-session фикса (#2164, 2026-08) — тесты ниже
|
||
# обращаются к нему через `proxy_pool.touch(...)` (module attribute), а не прямым
|
||
# top-level импортом, чтобы отсутствие функции в pre-fix коде падало ТОЛЬКО в этих
|
||
# тестах (AttributeError — capability реально не существовала), а не ронуло
|
||
# коллекцию всего файла и не маскировало value-based тесты acquire/release/
|
||
# mark_health/reap выше, которые этим фиксом не менялись.
|
||
|
||
# ── stateful fake session ────────────────────────────────────────────────────
|
||
|
||
|
||
class _FakeResult:
|
||
def __init__(self, rows: list[dict[str, Any]]):
|
||
self._rows = rows
|
||
|
||
def mappings(self) -> _FakeResult:
|
||
return self
|
||
|
||
def fetchone(self) -> dict[str, Any] | None:
|
||
return self._rows[0] if self._rows else None
|
||
|
||
def all(self) -> list[dict[str, Any]]:
|
||
return list(self._rows)
|
||
|
||
def fetchall(self) -> list[Any]:
|
||
# для RETURNING id: код делает r.id → нужен attribute-access
|
||
return [type("Row", (), r)() for r in self._rows]
|
||
|
||
|
||
class FakeSession:
|
||
"""Эмуляция Session поверх in-memory списков scrape_proxies + scrape_proxy_source_bans."""
|
||
|
||
def __init__(self, rows: list[dict[str, Any]], bans: list[dict[str, Any]] | None = None):
|
||
self.rows = rows
|
||
# #2600 п.2: строки scrape_proxy_source_bans (proxy_id, source, banned_until,
|
||
# ban_count) — бан теперь свойство ПАРЫ «узел × источник», не узла.
|
||
self.bans: list[dict[str, Any]] = bans or []
|
||
# deep-review fix 2 (#2600): вызовы pg_advisory_xact_lock — для теста
|
||
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
|
||
# проверить — фейксессия однопоточна).
|
||
self.advisory_lock_calls: list[int] = []
|
||
|
||
# helpers
|
||
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
||
return next((r for r in self.rows if r["id"] == pid), None)
|
||
|
||
def _ban(self, pid: int, source: str) -> dict[str, Any] | None:
|
||
return next((b for b in self.bans if b["proxy_id"] == pid and b["source"] == source), None)
|
||
|
||
def _has_active_ban(self, pid: int, source: str) -> bool:
|
||
ban = self._ban(pid, source)
|
||
return ban is not None and ban["banned_until"] > datetime.now(UTC)
|
||
|
||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||
sql = str(stmt)
|
||
p = params or {}
|
||
|
||
if "expires_at <= now()" in sql and "FOR UPDATE SKIP LOCKED" not in sql:
|
||
# acquire() отдельный WARNING-запрос (#3287): просроченные enabled+свободные
|
||
# узлы, вне зависимости от affinity/consecutive_fails — они всё равно не
|
||
# попадут ни в основной, ни в fallback-отбор ниже.
|
||
expired = [
|
||
r
|
||
for r in self.rows
|
||
if r["enabled"]
|
||
and r["leased_by"] is None
|
||
and r.get("expires_at") is not None
|
||
and r["expires_at"] <= datetime.now(UTC)
|
||
]
|
||
return _FakeResult([{"id": r["id"]} for r in expired])
|
||
|
||
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
|
||
max_fails = p["max_fails"]
|
||
provider = p["provider"]
|
||
# #2600 п.2: узел с активным баном по ЭТОМУ source не выдаётся. Как и с
|
||
# protects_last_node ниже — фильтруем ТОЛЬКО если сам SQL реально содержит
|
||
# NOT EXISTS по scrape_proxy_source_bans, иначе мок реализовывал бы логику
|
||
# независимо от проверяемого кода и не отличил бы старый запрос от нового.
|
||
filters_bans = "scrape_proxy_source_bans" in sql
|
||
# #3287: истёкшая аренда порта не выдаётся — гейтим по подстроке фильтра,
|
||
# тот же принцип, что и у filters_bans выше.
|
||
filters_expiry = "expires_at" in sql
|
||
|
||
def _not_banned(row: dict[str, Any]) -> bool:
|
||
return not filters_bans or not self._has_active_ban(row["id"], provider)
|
||
|
||
def _not_expired(row: dict[str, Any]) -> bool:
|
||
if not filters_expiry:
|
||
return True
|
||
exp = row.get("expires_at")
|
||
return exp is None or exp > datetime.now(UTC)
|
||
|
||
if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any'
|
||
cands = [
|
||
r
|
||
for r in self.rows
|
||
if r["enabled"]
|
||
and r["consecutive_fails"] < max_fails
|
||
and r["provider_affinity"] in (provider, "any")
|
||
and r["leased_by"] is None
|
||
and _not_banned(r)
|
||
and _not_expired(r)
|
||
]
|
||
else: # fallback: любая affinity, но не последний узел выделенной affinity
|
||
# (domclick и т.п. — #2600 review). ВАЖНО: применяем эту фильтрацию,
|
||
# только если сама SQL реально содержит защиту (EXISTS-подзапрос) —
|
||
# иначе мок реализовывал бы бизнес-логику независимо от проверяемого
|
||
# кода и не смог бы отличить старый (незащищённый) fallback-запрос от
|
||
# нового. Тот же класс бага, что был с "enabled" в mark_health-моке.
|
||
protects_last_node = "EXISTS" in sql
|
||
# backup засчитывается, только если он ПРИГОДЕН для своей affinity —
|
||
# тоже гейтим по подстроке (b2-подзапрос), иначе мок «чинил» бы
|
||
# незащищённый SQL сам.
|
||
backup_must_be_usable = "b2.banned_until > now()" in sql
|
||
|
||
def _has_backup(row: dict[str, Any]) -> bool:
|
||
if row["provider_affinity"] == "any":
|
||
return True
|
||
return any(
|
||
other["provider_affinity"] == row["provider_affinity"]
|
||
and other["enabled"]
|
||
and other["id"] != row["id"]
|
||
and not (
|
||
backup_must_be_usable
|
||
and self._has_active_ban(other["id"], other["provider_affinity"])
|
||
)
|
||
for other in self.rows
|
||
)
|
||
|
||
cands = [
|
||
r
|
||
for r in self.rows
|
||
if r["enabled"]
|
||
and r["consecutive_fails"] < max_fails
|
||
and r["leased_by"] is None
|
||
and _not_banned(r)
|
||
and _not_expired(r)
|
||
and (not protects_last_node or _has_backup(r))
|
||
]
|
||
# ORDER BY (browser_unfit_since IS NOT NULL), last_ok_at NULLS LAST, id.
|
||
# Первый ключ гейтим по подстроке самого SQL (как ban-фильтры выше): иначе
|
||
# мок сортировал бы «правильно» независимо от боевого запроса и не отличил
|
||
# бы код до #2723 от кода после.
|
||
deprioritises_unfit = "browser_unfit_since IS NOT NULL" in sql
|
||
cands.sort(
|
||
key=lambda r: (
|
||
bool(deprioritises_unfit and r.get("browser_unfit_since") is not None),
|
||
r["last_ok_at"] is None,
|
||
r["last_ok_at"] or datetime.min.replace(tzinfo=UTC),
|
||
r["id"],
|
||
)
|
||
)
|
||
return _FakeResult(cands[:1])
|
||
|
||
if "SET leased_by = CAST(:run_id" in sql: # acquire lease UPDATE
|
||
row = self._by_id(p["id"])
|
||
if row is not None:
|
||
row["leased_by"] = p["run_id"]
|
||
row["leased_at"] = datetime.now(UTC)
|
||
return _FakeResult([])
|
||
|
||
if "make_interval(mins =>" in sql and "leased_by IS NOT NULL" in sql: # reap
|
||
cutoff = datetime.now(UTC) - timedelta(minutes=p["mins"])
|
||
freed: list[dict[str, Any]] = []
|
||
for r in self.rows:
|
||
if (
|
||
r["leased_by"] is not None
|
||
and r["leased_at"] is not None
|
||
and r["leased_at"] < cutoff
|
||
):
|
||
r["leased_by"] = None
|
||
r["leased_at"] = None
|
||
freed.append({"id": r["id"]})
|
||
return _FakeResult(freed)
|
||
|
||
if "SET leased_by = NULL" in sql: # release (WHERE id)
|
||
row = self._by_id(p["id"])
|
||
if row is not None:
|
||
row["leased_by"] = None
|
||
row["leased_at"] = None
|
||
return _FakeResult([])
|
||
|
||
if "SET leased_at = now()" in sql: # touch heartbeat (#2164 sticky-session fix)
|
||
row = self._by_id(p["id"])
|
||
if row is not None and row["leased_by"] is not None:
|
||
row["leased_at"] = datetime.now(UTC)
|
||
return _FakeResult([])
|
||
|
||
if "SET consecutive_fails = 0" in sql: # mark_health ok
|
||
row = self._by_id(p["id"])
|
||
if row is not None:
|
||
row["consecutive_fails"] = 0
|
||
row["exit_ip"] = p["exit_ip"]
|
||
row["latency_ms"] = p["latency_ms"]
|
||
row["last_ok_at"] = datetime.now(UTC)
|
||
row["last_check_at"] = datetime.now(UTC)
|
||
# Ручное выключение (#2610): реальный SQL реанимирует (enabled=true)
|
||
# ТОЛЬКО если disabled_reason IS NULL — мок обязан честно это
|
||
# воспроизвести, иначе тест не поймает регресс "снова безусловно
|
||
# enabled=true" (тот же класс бага, что был с "enabled" in sql до
|
||
# #2600 review).
|
||
if "disabled_reason IS NULL" in sql:
|
||
if row.get("disabled_reason") is None:
|
||
row["enabled"] = True
|
||
elif "enabled" in sql:
|
||
row["enabled"] = True
|
||
# #2723: если боевой mark_health когда-нибудь снова начнёт обнулять
|
||
# ещё и браузерный вердикт (как делал до фикса — тот жил в общем
|
||
# consecutive_fails), мок обязан это воспроизвести, иначе
|
||
# test_ipify_success_does_not_erase_browser_verdict останется зелёным
|
||
# на сломанном коде.
|
||
if "browser_fail_streak = 0" in sql:
|
||
row["browser_fail_streak"] = 0
|
||
row["browser_unfit_since"] = None
|
||
if "RETURNING disabled_reason" in sql:
|
||
return _FakeResult([{"disabled_reason": row.get("disabled_reason")}])
|
||
return _FakeResult([])
|
||
|
||
if "consecutive_fails = consecutive_fails + 1" in sql: # mark_health fail
|
||
row = self._by_id(p["id"])
|
||
if row is not None:
|
||
row["consecutive_fails"] += 1
|
||
row["last_check_at"] = datetime.now(UTC)
|
||
if row["consecutive_fails"] >= p["disable_threshold"]:
|
||
row["enabled"] = False
|
||
return _FakeResult([])
|
||
|
||
if "WHERE enabled" in sql and "ORDER BY id" in sql: # healthcheck SELECT (#2600 п.1)
|
||
recheck_minutes = p["disabled_recheck_minutes"]
|
||
cutoff = datetime.now(UTC) - timedelta(minutes=recheck_minutes)
|
||
cands = [
|
||
r
|
||
for r in self.rows
|
||
if r["enabled"] or r.get("last_check_at") is None or r["last_check_at"] < cutoff
|
||
]
|
||
rows = sorted(cands, key=lambda r: r["id"])
|
||
# #2723: браузерная проба со своим тактом. Признак считаем, только если
|
||
# боевой SQL его реально запрашивает (см. гейты по подстрокам выше).
|
||
if "browser_probe_due" in sql:
|
||
b_cutoff = datetime.now(UTC) - timedelta(minutes=p["browser_probe_minutes"])
|
||
return _FakeResult(
|
||
[
|
||
dict(
|
||
r,
|
||
browser_probe_due=(
|
||
r.get("browser_check_at") is None
|
||
or r["browser_check_at"] < b_cutoff
|
||
),
|
||
)
|
||
for r in rows
|
||
]
|
||
)
|
||
return _FakeResult([dict(r) for r in rows])
|
||
|
||
if "SET browser_fail_streak = 0" in sql: # mark_browser_health ok (#2723)
|
||
row = self._by_id(p["id"])
|
||
if row is None:
|
||
return _FakeResult([])
|
||
was_unfit = row.get("browser_unfit_since") is not None
|
||
row["browser_fail_streak"] = 0
|
||
row["browser_unfit_since"] = None
|
||
row["browser_check_at"] = datetime.now(UTC)
|
||
return _FakeResult([{"was_unfit": was_unfit}])
|
||
|
||
if "browser_fail_streak = browser_fail_streak + 1" in sql: # mark_browser_health fail
|
||
row = self._by_id(p["id"])
|
||
if row is None:
|
||
return _FakeResult([])
|
||
row["browser_fail_streak"] = row.get("browser_fail_streak", 0) + 1
|
||
if row["browser_fail_streak"] >= p["threshold"]:
|
||
if row.get("browser_unfit_since") is None:
|
||
row["browser_unfit_since"] = datetime.now(UTC)
|
||
# такт двигаем ТОЛЬКО на подтверждённом провале — иначе неподтверждённое
|
||
# подозрение ждало бы полный BROWSER_PROBE_MINUTES (#2723)
|
||
row["browser_check_at"] = datetime.now(UTC)
|
||
return _FakeResult(
|
||
[
|
||
{
|
||
"browser_fail_streak": row["browser_fail_streak"],
|
||
"browser_unfit_since": row.get("browser_unfit_since"),
|
||
}
|
||
]
|
||
)
|
||
|
||
if "pg_advisory_xact_lock" in sql: # deep-review fix 2 (#2600) — mark_banned serialize
|
||
self.advisory_lock_calls.append(p["key"])
|
||
return _FakeResult([])
|
||
|
||
if "INSERT INTO scrape_proxy_source_bans" in sql: # mark_banned UPSERT (#2600 п.2)
|
||
proxy_id = p["proxy_id"]
|
||
source = p["source"]
|
||
max_fails = p["max_fails"]
|
||
if self._by_id(proxy_id) is None:
|
||
return _FakeResult([]) # узла нет — no-op
|
||
|
||
# Оба ban-предиката гейтим по подстрокам самого SQL (как в acquire-ветке):
|
||
# иначе мок реализовывал бы защиту сам и тест оставался бы зелёным даже
|
||
# после удаления NOT EXISTS из боевого запроса.
|
||
filters_bans = "b.banned_until > now()" in sql
|
||
backup_must_be_usable = "b2.banned_until > now()" in sql
|
||
|
||
def _is_candidate(sp: dict[str, Any]) -> bool:
|
||
if not (sp["enabled"] and sp["consecutive_fails"] < max_fails):
|
||
return False
|
||
if filters_bans and self._has_active_ban(sp["id"], source):
|
||
return False # уже забанен этим же источником — не кандидат
|
||
if sp["provider_affinity"] in (source, "any"):
|
||
return True
|
||
# fallback-safe: другой ПРИГОДНЫЙ узел ТОЙ ЖЕ affinity (банимый узел
|
||
# остаётся enabled и тоже считается — бан теперь per-source; а вот
|
||
# забаненный своим же источником backup'ом не считается).
|
||
return any(
|
||
other["provider_affinity"] == sp["provider_affinity"]
|
||
and other["enabled"]
|
||
and other["id"] != sp["id"]
|
||
and not (
|
||
backup_must_be_usable
|
||
and self._has_active_ban(other["id"], other["provider_affinity"])
|
||
)
|
||
for other in self.rows
|
||
)
|
||
|
||
still_available = any(r["id"] != proxy_id and _is_candidate(r) for r in self.rows)
|
||
if not still_available:
|
||
return _FakeResult([]) # protected — последний узел для source, бан не пишем
|
||
|
||
now = datetime.now(UTC)
|
||
ban = self._ban(proxy_id, source)
|
||
if ban is None:
|
||
ban = {
|
||
"proxy_id": proxy_id,
|
||
"source": source,
|
||
"ban_count": 1,
|
||
"banned_until": now + timedelta(hours=p["base_hours"]),
|
||
"reason": p["reason"],
|
||
}
|
||
self.bans.append(ban)
|
||
else:
|
||
# #2803-follow-up: активную строку чужого владельца не перехватываем.
|
||
# Гейтим по подстроке боевого SQL (как ban-предикаты выше) — иначе мок
|
||
# реализовал бы защиту сам и тест был бы зелёным на сломанном коде.
|
||
defends_owner = "scrape_proxy_source_bans.banned_until <= now()" in sql
|
||
if (
|
||
defends_owner
|
||
and ban["banned_until"] > now
|
||
and ban.get("reason") != p["reason"]
|
||
and p["reason"] != p.get("live_reason")
|
||
):
|
||
return _FakeResult([]) # владельца активной строки не меняем
|
||
# эскалация: срок = base * 2^(новый ban_count - 1), потолок max_hours
|
||
ban["ban_count"] += 1
|
||
hours = min(p["base_hours"] * 2 ** (ban["ban_count"] - 1), p["max_hours"])
|
||
ban["banned_until"] = now + timedelta(hours=hours)
|
||
ban["reason"] = p["reason"]
|
||
return _FakeResult(
|
||
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
|
||
)
|
||
|
||
if "DELETE FROM scrape_proxy_source_bans" in sql and "proxy_id = CAST" in sql:
|
||
# clear_source_bans: снять баны узла (все либо один source), #2600 п.2.
|
||
# Фильтр по reason (#2800) гейтим по подстроке боевого SQL — как ban-фильтры
|
||
# в acquire-ветке: иначе мок «чинил» бы код, который фильтра не содержит, и
|
||
# тест на «успешная проба не гасит чужой бан» остался бы зелёным на сломанном.
|
||
filters_reason = "reason = CAST(:only_reason AS text)" in sql
|
||
only_reason = p.get("only_reason") if filters_reason else None
|
||
cleared = [
|
||
b
|
||
for b in self.bans
|
||
if b["proxy_id"] == p["proxy_id"]
|
||
and (p["source"] is None or b["source"] == p["source"])
|
||
and (only_reason is None or b.get("reason") == only_reason)
|
||
]
|
||
self.bans = [b for b in self.bans if b not in cleared]
|
||
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||
|
||
if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2)
|
||
cutoff = datetime.now(UTC) - timedelta(days=p["days"])
|
||
purged = [{"proxy_id": b["proxy_id"]} for b in self.bans if b["banned_until"] < cutoff]
|
||
self.bans = [b for b in self.bans if b["banned_until"] >= cutoff]
|
||
return _FakeResult(purged)
|
||
|
||
if "SELECT reason, banned_until" in sql: # mark_banned: кто держит активный бан пары
|
||
ban = self._ban(p["proxy_id"], p["source"])
|
||
if ban is None or ban["banned_until"] <= datetime.now(UTC):
|
||
return _FakeResult([])
|
||
return _FakeResult([{"reason": ban.get("reason"), "banned_until": ban["banned_until"]}])
|
||
|
||
if "SELECT enabled, disabled_reason FROM scrape_proxies" in sql: # mark_banned diag read
|
||
row = self._by_id(p["id"])
|
||
if row is None:
|
||
return _FakeResult([])
|
||
return _FakeResult(
|
||
[{"enabled": row["enabled"], "disabled_reason": row["disabled_reason"]}]
|
||
)
|
||
|
||
raise AssertionError(f"unhandled SQL: {sql}")
|
||
|
||
def commit(self) -> None:
|
||
pass
|
||
|
||
def rollback(self) -> None:
|
||
pass
|
||
|
||
|
||
def _proxy(
|
||
pid: int,
|
||
*,
|
||
affinity: str = "any",
|
||
enabled: bool = True,
|
||
fails: int = 0,
|
||
leased_by: int | None = None,
|
||
leased_at: datetime | None = None,
|
||
last_ok_at: datetime | None = None,
|
||
last_check_at: datetime | None = None,
|
||
kind: str = "http",
|
||
rotate_url: str | None = None,
|
||
disabled_reason: str | None = None,
|
||
browser_unfit_since: datetime | None = None,
|
||
browser_fail_streak: int = 0,
|
||
browser_check_at: datetime | None = None,
|
||
expires_at: datetime | None = None,
|
||
) -> dict[str, Any]:
|
||
return {
|
||
"id": pid,
|
||
"url": f"http://u:p@h{pid}:8080",
|
||
"kind": kind,
|
||
"rotate_url": rotate_url,
|
||
"provider_affinity": affinity,
|
||
"enabled": enabled,
|
||
"disabled_reason": disabled_reason,
|
||
"consecutive_fails": fails,
|
||
"leased_by": leased_by,
|
||
"leased_at": leased_at,
|
||
"last_ok_at": last_ok_at,
|
||
"last_check_at": last_check_at,
|
||
"exit_ip": None,
|
||
"latency_ms": None,
|
||
# #2723: здоровье браузерного тракта — отдельные поля, миграция 228.
|
||
"browser_unfit_since": browser_unfit_since,
|
||
"browser_fail_streak": browser_fail_streak,
|
||
"browser_check_at": browser_check_at,
|
||
# срок аренды порта у провайдера — #3287, None = "не отслеживается"
|
||
"expires_at": expires_at,
|
||
}
|
||
|
||
|
||
# ── acquire ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
def test_acquire_returns_free_healthy_of_affinity() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="cian")])
|
||
lease = acquire(db, "avito", run_id=100) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 1
|
||
assert lease.url == "http://u:p@h1:8080"
|
||
assert db._by_id(1)["leased_by"] == 100 # lease проставлен
|
||
|
||
|
||
def test_acquire_includes_any_affinity() -> None:
|
||
db = FakeSession([_proxy(1, affinity="any")])
|
||
lease = acquire(db, "avito", run_id=5) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 1
|
||
|
||
|
||
def test_two_acquire_return_distinct_proxies() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="avito")])
|
||
first = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
second = acquire(db, "avito", run_id=2) # type: ignore[arg-type]
|
||
assert first is not None and second is not None
|
||
assert first.id != second.id # первый лизнут → второй берёт другой
|
||
|
||
|
||
def test_acquire_empty_pool_returns_none() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito", leased_by=99)]) # единственный занят
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
def test_acquire_skips_disabled() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito", enabled=False)])
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
def test_acquire_skips_unhealthy() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito", fails=MAX_CONSECUTIVE_FAILS)])
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
def test_acquire_without_run_id_uses_marker() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito")])
|
||
lease = acquire(db, "avito") # run_id=None # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert db._by_id(1)["leased_by"] == proxy_pool.NON_RUN_LEASE_MARKER
|
||
|
||
|
||
# ── acquire: fallback affinity (#2600 п.3 — не морить источник голодом) ────────
|
||
|
||
|
||
def test_acquire_prefers_own_affinity_when_available() -> None:
|
||
"""Своих (affinity=avito) хватает — приоритет не сломан, чужой (cian) не берём."""
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="cian")])
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 1
|
||
|
||
|
||
def test_acquire_falls_back_to_other_affinity_when_no_own_free() -> None:
|
||
"""Свободных avito/any нет, но есть свободный здоровый cian с бэкапом → fallback, а не None.
|
||
|
||
Два cian-узла — забрать один через fallback безопасно: у cian остаётся другой
|
||
enabled-узел (protection на "последний узел affinity" не срабатывает).
|
||
"""
|
||
db = FakeSession([_proxy(1, affinity="cian"), _proxy(2, affinity="cian")])
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 1
|
||
assert db._by_id(1)["leased_by"] == 1
|
||
|
||
|
||
def test_acquire_no_fallback_when_nothing_free_at_all() -> None:
|
||
"""Fallback не выдумывает прокси из воздуха — если свободных нет вообще, None."""
|
||
db = FakeSession([_proxy(1, affinity="cian", leased_by=99)]) # занят
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
# ── acquire: fallback НЕ забирает последний узел выделенной affinity (review #2609) ──
|
||
#
|
||
# domclick — ровно один узел (прод scrape_proxies.id=1), намеренно вырезанный из общего
|
||
# пула через provider_affinity='domclick': QRATOR банит всё, кроме этого одного чистого
|
||
# residential-адреса (см. 173_scrape_proxies_add_domclick_affinity.sql). Если fallback
|
||
# заберёт его под avito/cian/yandex — domclick (сейчас исправно собирает: 6501 активных
|
||
# объявлений, 368/сутки) останется без прокси вообще. Починка одного источника ценой
|
||
# полной поломки другого недопустима.
|
||
|
||
|
||
def test_acquire_fallback_protects_last_node_of_dedicated_affinity() -> None:
|
||
"""Единственный enabled-узел domclick НЕ отдаётся avito через fallback — None."""
|
||
db = FakeSession([_proxy(1, affinity="domclick")])
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
assert db._by_id(1)["leased_by"] is None # узел не тронут
|
||
|
||
|
||
def test_acquire_fallback_allows_when_dedicated_affinity_has_backup() -> None:
|
||
"""Второй enabled-узел domclick есть → fallback как и раньше отдаёт свободный."""
|
||
db = FakeSession([_proxy(1, affinity="domclick"), _proxy(2, affinity="domclick")])
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 1
|
||
assert db._by_id(2)["leased_by"] is None # у domclick остался живой запасной узел
|
||
|
||
|
||
def test_acquire_fallback_backup_must_be_usable_for_its_own_source() -> None:
|
||
"""deep-review #2600 п.2: «backup» — это ПРИГОДНЫЙ узел, а не просто enabled.
|
||
|
||
Два узла domclick: node1 забанен САМИМ domclick'ом (законно — node2 тогда был жив) и
|
||
сейчас занят чужим прогоном, node2 свободен. Для domclick node2 — последний рабочий.
|
||
Раньше EXISTS видел node1 как backup (он ведь enabled) и разрешал fallback увести
|
||
node2 под avito — domclick оставался бы без прокси вообще при двух включённых узлах.
|
||
"""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="domclick", leased_by=99), _proxy(2, affinity="domclick")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "domclick",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:domclick",
|
||
}
|
||
],
|
||
)
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
assert db._by_id(2)["leased_by"] is None # последний рабочий узел domclick не тронут
|
||
# сам domclick при этом обслуживается: node2 свободен и не забанен
|
||
lease = acquire(db, "domclick", run_id=2) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 2
|
||
|
||
|
||
# ── acquire × expires_at — просроченная аренда порта не выдаётся (#3287) ───────
|
||
|
||
|
||
def test_acquire_skips_expired_lease() -> None:
|
||
"""expires_at в прошлом → узел не выдаётся, даже если enabled/здоров/свободен."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="avito", expires_at=datetime.now(UTC) - timedelta(minutes=1))]
|
||
)
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
assert db._by_id(1)["leased_by"] is None # узел не тронут
|
||
|
||
|
||
def test_acquire_issues_proxy_with_null_expires_at() -> None:
|
||
"""expires_at IS NULL — "срок не отслеживается", выдаче не мешает (как раньше)."""
|
||
db = FakeSession([_proxy(1, affinity="avito", expires_at=None)])
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 1
|
||
|
||
|
||
def test_acquire_issues_proxy_with_future_expires_at() -> None:
|
||
"""expires_at в будущем — аренда ещё жива, узел выдаётся."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="avito", expires_at=datetime.now(UTC) + timedelta(hours=1))]
|
||
)
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 1
|
||
|
||
|
||
def test_acquire_fallback_skips_expired_lease() -> None:
|
||
"""Fallback-ветка (чужая affinity) тоже не выдаёт просроченный узел."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="cian", expires_at=datetime.now(UTC) - timedelta(minutes=1))]
|
||
)
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
def test_acquire_fallback_prefers_non_expired_over_expired() -> None:
|
||
"""Просроченный узел пропускается, живой той же чужой affinity — выдан fallback'ом."""
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, affinity="cian", expires_at=datetime.now(UTC) - timedelta(minutes=1)),
|
||
_proxy(2, affinity="cian", expires_at=datetime.now(UTC) + timedelta(hours=1)),
|
||
]
|
||
)
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 2
|
||
|
||
|
||
def test_acquire_warns_on_expired_proxy(caplog: pytest.LogCaptureFixture) -> None:
|
||
"""Просроченный узел логируется WARNING'ом — оператору нужен явный сигнал (#3287)."""
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, affinity="avito", expires_at=datetime.now(UTC) - timedelta(minutes=1)),
|
||
_proxy(2, affinity="avito", expires_at=datetime.now(UTC) + timedelta(hours=1)),
|
||
]
|
||
)
|
||
with caplog.at_level("WARNING", logger="app.services.proxy_pool"):
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 2 # живой узел всё равно выдан
|
||
assert any(
|
||
"expired" in r.message and "id=1" in r.message for r in caplog.records
|
||
), "ожидался WARNING про просроченный proxy id=1"
|
||
|
||
|
||
# ── release ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
def test_release_frees_proxy() -> None:
|
||
db = FakeSession([_proxy(1, affinity="avito")])
|
||
lease = acquire(db, "avito", run_id=7) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert acquire(db, "avito", run_id=8) is None # type: ignore[arg-type] # занят
|
||
release(db, 1) # type: ignore[arg-type]
|
||
assert db._by_id(1)["leased_by"] is None
|
||
assert acquire(db, "avito", run_id=9) is not None # type: ignore[arg-type] # снова свободен
|
||
|
||
|
||
# ── mark_health ──────────────────────────────────────────────────────────────
|
||
|
||
|
||
def test_mark_health_fail_increments() -> None:
|
||
db = FakeSession([_proxy(1, fails=0)])
|
||
mark_health(db, 1, ok=False) # type: ignore[arg-type]
|
||
assert db._by_id(1)["consecutive_fails"] == 1
|
||
assert db._by_id(1)["enabled"] is True # ещё не порог
|
||
|
||
|
||
def test_mark_health_fail_disables_at_threshold() -> None:
|
||
db = FakeSession([_proxy(1, fails=DISABLE_THRESHOLD - 1)])
|
||
mark_health(db, 1, ok=False) # type: ignore[arg-type]
|
||
assert db._by_id(1)["consecutive_fails"] == DISABLE_THRESHOLD
|
||
assert db._by_id(1)["enabled"] is False # авто-disable
|
||
|
||
|
||
def test_mark_health_ok_resets_and_records() -> None:
|
||
db = FakeSession([_proxy(1, fails=4)])
|
||
mark_health(db, 1, ok=True, exit_ip="1.2.3.4", latency_ms=88) # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["consecutive_fails"] == 0
|
||
assert row["exit_ip"] == "1.2.3.4"
|
||
assert row["latency_ms"] == 88
|
||
assert row["last_ok_at"] is not None
|
||
|
||
|
||
def test_mark_health_ok_revives_disabled_proxy() -> None:
|
||
"""Успешная проба реанимирует выключенный узел (#2600 п.1) — enabled=true, fails=0."""
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD)])
|
||
mark_health(db, 1, ok=True) # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is True
|
||
assert row["consecutive_fails"] == 0
|
||
|
||
|
||
# ── mark_health / disabled_reason (#2610 — ручное vs авто-выключение) ──────────
|
||
#
|
||
# red/green контракт issue: (a) авто-выключенный узел воскресает по успешной пробе
|
||
# (#2609 не сломан — дублирует test_mark_health_ok_revives_disabled_proxy выше, но
|
||
# явно рядом с (b)/(c) для контраста), (b) ручно-выключенный НЕ воскресает + лог,
|
||
# (c) ручное включение сбрасывает флаг → узел снова авто-восстанавливаем (уровень
|
||
# admin API — см. tests/test_admin_proxies.py, mark_health сам флаг не трогает).
|
||
|
||
|
||
def test_mark_health_ok_revives_auto_disabled_proxy_a() -> None:
|
||
"""(a) Авто-выключенный (disabled_reason=NULL) воскресает — поведение #2609 сохранено."""
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD, disabled_reason=None)])
|
||
mark_health(db, 1, ok=True) # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is True
|
||
assert row["consecutive_fails"] == 0
|
||
|
||
|
||
def test_mark_health_ok_does_not_revive_manually_disabled_proxy_b(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""(b) Ручно-выключенный (disabled_reason НЕ NULL) НЕ воскресает — и это логируется."""
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD, disabled_reason="manual")])
|
||
with caplog.at_level("WARNING"):
|
||
mark_health(db, 1, ok=True) # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is False # НЕ реанимирован, несмотря на ok=True
|
||
assert row["consecutive_fails"] == 0 # fails всё равно сбрасывается пробой
|
||
assert any("manually disabled" in rec.message for rec in caplog.records)
|
||
|
||
|
||
# ── reap_stale_leases ────────────────────────────────────────────────────────
|
||
|
||
|
||
def test_reap_frees_stale_lease_keeps_fresh() -> None:
|
||
old = datetime.now(UTC) - timedelta(minutes=120)
|
||
fresh = datetime.now(UTC) - timedelta(minutes=1)
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, leased_by=50, leased_at=old),
|
||
_proxy(2, leased_by=51, leased_at=fresh),
|
||
]
|
||
)
|
||
freed = reap_stale_leases(db, older_than_minutes=30) # type: ignore[arg-type]
|
||
assert freed == 1
|
||
assert db._by_id(1)["leased_by"] is None # протухший освобождён
|
||
assert db._by_id(2)["leased_by"] == 51 # свежий не тронут
|
||
|
||
|
||
# ── touch (heartbeat) — sticky browser-session lease fix (#2164, 2026-08) ─────
|
||
#
|
||
# Живая регрессия: BrowserFetcher раньше acquire/release-ил прокси НА КАЖДЫЙ /fetch —
|
||
# при N>=2 живых узлах пула это гарантированно меняло прокси между соседними запросами
|
||
# и гоняло camoufox relaunch на каждый /fetch (17 relaunch'ей за 15 минут в проде).
|
||
# Фикс — один lease на весь жизненный цикл BrowserFetcher (сессия, часы). Коллизия:
|
||
# reap_stale_leases отбирает lease старше STALE_LEASE_MINUTES=30, а прогоны бывают
|
||
# ДОЛЬШЕ (полная загрузка Циана — часами). touch() — heartbeat, которым BrowserFetcher
|
||
# продлевает leased_at на каждый /fetch, пока сессия жива; без touch (мёртвый/зависший
|
||
# run) lease по-прежнему реапится штатно — semantics краш-recovery не ослаблена.
|
||
|
||
|
||
def test_touch_refreshes_leased_at_of_active_lease() -> None:
|
||
old = datetime.now(UTC) - timedelta(minutes=45)
|
||
db = FakeSession([_proxy(1, leased_by=100, leased_at=old)])
|
||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||
assert db._by_id(1)["leased_at"] > old
|
||
|
||
|
||
def test_touch_noop_when_not_leased() -> None:
|
||
"""Прокси свободен (leased_by=NULL) — touch не должен «арендовывать» его тайком."""
|
||
db = FakeSession([_proxy(1, leased_by=None, leased_at=None)])
|
||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["leased_by"] is None
|
||
assert row["leased_at"] is None # touch не проставил leased_at свободному узлу
|
||
|
||
|
||
def test_touch_prevents_reap_of_long_running_session() -> None:
|
||
"""КЛЮЧЕВАЯ коллизия из PR: lease взят 45 минут назад (> STALE_LEASE_MINUTES=30 —
|
||
reap_stale_leases его бы отобрал по «возрасту acquire»), но BrowserFetcher вызывал
|
||
touch() на каждый /fetch все эти 45 минут — leased_at всегда свежий. reap НЕ должен
|
||
освободить активную многочасовую сессию."""
|
||
old_acquire = datetime.now(UTC) - timedelta(minutes=45)
|
||
db = FakeSession([_proxy(1, leased_by=proxy_pool.NON_RUN_LEASE_MARKER, leased_at=old_acquire)])
|
||
|
||
# heartbeat только что прошёл (BrowserFetcher вызвал touch на последнем /fetch)
|
||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||
|
||
freed = reap_stale_leases(db, older_than_minutes=STALE_LEASE_MINUTES) # type: ignore[arg-type]
|
||
|
||
assert freed == 0
|
||
assert db._by_id(1)["leased_by"] == proxy_pool.NON_RUN_LEASE_MARKER # НЕ отобран
|
||
|
||
|
||
def test_touch_absence_still_reaps_dead_session() -> None:
|
||
"""Без heartbeat'а (упавший/зависший прогон, ни разу не сходивший в touch) reaper
|
||
по-прежнему освобождает протухший lease — семантика краш-recovery не ослаблена
|
||
фиксом (мы НЕ увеличивали STALE_LEASE_MINUTES, чтобы не притупить эту защиту)."""
|
||
old_acquire = datetime.now(UTC) - timedelta(minutes=45)
|
||
db = FakeSession([_proxy(1, leased_by=proxy_pool.NON_RUN_LEASE_MARKER, leased_at=old_acquire)])
|
||
|
||
freed = reap_stale_leases(db, older_than_minutes=STALE_LEASE_MINUTES) # type: ignore[arg-type]
|
||
|
||
assert freed == 1
|
||
assert db._by_id(1)["leased_by"] is None # мёртвая сессия реапится как раньше
|
||
|
||
|
||
# ── run_proxy_healthcheck ────────────────────────────────────────────────────
|
||
|
||
|
||
async def test_healthcheck_probes_enabled_and_marks_health(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
recently_checked = datetime.now(UTC) - timedelta(minutes=5) # < DISABLED_RECHECK_MINUTES
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, fails=2),
|
||
# disabled, recheck ещё не наступил (недавно проверен) — не проверяется в этот прогон
|
||
_proxy(2, enabled=False, last_check_at=recently_checked),
|
||
_proxy(3, fails=0),
|
||
]
|
||
)
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
# прокси 1 «жив», прокси 3 «мёртв»
|
||
if "h1:" in url:
|
||
return True, "9.9.9.9", 42, None
|
||
return False, None, None, "other"
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["checked"] == 2 # только enabled (1 и 3), disabled recheck не наступил
|
||
assert counters["ok"] == 1
|
||
assert counters["failed"] == 1
|
||
assert counters["revived"] == 0
|
||
assert db._by_id(1)["consecutive_fails"] == 0 # ok → сброс
|
||
assert db._by_id(1)["exit_ip"] == "9.9.9.9"
|
||
assert db._by_id(3)["consecutive_fails"] == 1 # fail → инкремент
|
||
|
||
|
||
# ── run_proxy_healthcheck: self-healing disabled-узлов (#2600 п.1) ─────────────
|
||
|
||
|
||
async def test_healthcheck_revives_disabled_proxy_on_success(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Выключенный узел с успешной пробой возвращается в строй, revived++."""
|
||
stale_check = datetime.now(UTC) - timedelta(minutes=DISABLED_RECHECK_MINUTES + 5)
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD, last_check_at=stale_check)])
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
return True, "5.5.5.5", 30, None
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["checked"] == 1
|
||
assert counters["ok"] == 1
|
||
assert counters["revived"] == 1
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is True
|
||
assert row["consecutive_fails"] == 0
|
||
|
||
|
||
async def test_healthcheck_does_not_revive_manually_disabled_proxy(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""(b)/(с) на уровне healthcheck: ручно-выключенный узел пробуется (recheck наступил),
|
||
но остаётся disabled даже при успешной пробе — revived НЕ растёт (#2610)."""
|
||
stale_check = datetime.now(UTC) - timedelta(minutes=DISABLED_RECHECK_MINUTES + 5)
|
||
db = FakeSession(
|
||
[
|
||
_proxy(
|
||
1,
|
||
enabled=False,
|
||
fails=DISABLE_THRESHOLD,
|
||
last_check_at=stale_check,
|
||
disabled_reason="manual",
|
||
)
|
||
]
|
||
)
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
return True, "5.5.5.5", 30, None
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["checked"] == 1
|
||
assert counters["ok"] == 1
|
||
assert counters["revived"] == 0 # ручной флаг — mark_health не воскресил
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is False # остался выключенным
|
||
assert row["consecutive_fails"] == 0 # проба всё равно сбросила счётчик fails
|
||
|
||
|
||
async def test_healthcheck_skips_recently_checked_disabled_proxy(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Выключенный узел, проверенный недавно, повторно не проверяется в этот прогон."""
|
||
fresh_check = datetime.now(UTC) - timedelta(minutes=5) # < DISABLED_RECHECK_MINUTES
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD, last_check_at=fresh_check)])
|
||
probed: list[str] = []
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
probed.append(url) # не должно вызваться
|
||
return True, "5.5.5.5", 30, None
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["checked"] == 0
|
||
assert counters["revived"] == 0
|
||
assert probed == [] # провайдер не долбим каждый тик
|
||
assert db._by_id(1)["enabled"] is False # остался выключенным
|
||
|
||
|
||
async def test_healthcheck_checks_disabled_proxy_never_checked_before(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Выключенный узел без last_check_at (никогда не проверялся) — проверяется сразу."""
|
||
db = FakeSession([_proxy(1, enabled=False, fails=DISABLE_THRESHOLD, last_check_at=None)])
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
return False, None, None, "timeout"
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["checked"] == 1
|
||
assert counters["revived"] == 0 # неуспех — не реанимируем
|
||
assert db._by_id(1)["enabled"] is False
|
||
|
||
|
||
# ── mark_banned (#2600 п.2 — бан по паре «узел × источник») ────────────────────
|
||
#
|
||
# Red/green контракт issue: (a) распознанный бан → строка в scrape_proxy_source_bans,
|
||
# узел НЕ выключен глобально (в п.1 было enabled=false — площадка забанила IP, а не
|
||
# сломала прокси; узел обязан остаться живым для остальных источников); (b) повторный
|
||
# бан той же пары эскалирует срок; (c) последний достижимый для source узел НЕ банится
|
||
# (только лог); (d) сетевой сбой по-прежнему идёт через mark_health(ok=False), НЕ через
|
||
# mark_banned (проверяется на уровне browser_fetcher/curl_proxy_url тестов — здесь
|
||
# mark_banned сам по себе не участвует в различении причин, это забота caller'а).
|
||
|
||
|
||
def test_mark_banned_writes_source_ban_row() -> None:
|
||
"""Бан пишется в scrape_proxy_source_bans: source, ban_count=1, срок = base-часы."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
ban = db._ban(1, "avito")
|
||
assert ban is not None
|
||
assert ban["ban_count"] == 1
|
||
assert ban["reason"] == "banned:avito"
|
||
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS)
|
||
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||
|
||
|
||
def test_mark_banned_does_not_disable_node_globally() -> None:
|
||
"""ГЛАВНОЕ отличие от #2600 п.1: узел остаётся enabled и без disabled_reason —
|
||
глобальное выключение теперь только за оператором/авто-disable'ом (#2610)."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
row = db._by_id(1)
|
||
assert row["enabled"] is True
|
||
assert row["disabled_reason"] is None
|
||
|
||
|
||
def test_mark_banned_repeat_escalates_ban_count_and_duration() -> None:
|
||
"""Повторный бан той же пары: ban_count растёт, срок удваивается (base * 2^(N-1))."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
first_until = db._ban(1, "avito")["banned_until"]
|
||
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
|
||
ban = db._ban(1, "avito")
|
||
assert ban["ban_count"] == 2
|
||
assert len(db.bans) == 1 # PK (proxy_id, source) — дубля нет
|
||
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS * 2)
|
||
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||
assert ban["banned_until"] > first_until
|
||
|
||
|
||
def test_mark_banned_escalation_capped_at_max_hours() -> None:
|
||
"""Эскалация упирается в SOURCE_BAN_MAX_HOURS, а не растёт до бесконечности."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
for _ in range(8):
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
ban = db._ban(1, "avito")
|
||
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS)
|
||
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||
|
||
|
||
def test_mark_banned_different_sources_are_independent_rows() -> None:
|
||
"""Бан Авито и бан Циана на одном узле — две независимые строки, не перезапись."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
mark_banned(db, 1, source="cian") # type: ignore[arg-type]
|
||
assert {b["source"] for b in db.bans} == {"avito", "cian"}
|
||
assert all(b["ban_count"] == 1 for b in db.bans)
|
||
|
||
|
||
def test_mark_banned_unknown_id_is_noop() -> None:
|
||
"""Несуществующий proxy_id — бан не пишется, падать не должно (best-effort caller)."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 999, source="avito") # type: ignore[arg-type]
|
||
assert db.bans == []
|
||
|
||
|
||
def test_mark_banned_protects_last_live_node_own_affinity() -> None:
|
||
"""Единственный узел avito, других (свободных/'any') нет вообще — бан НЕ пишется,
|
||
только лог (issue #2600: бан не должен обрушить единственный источник целиком —
|
||
ходить через забаненный узел лучше, чем не ходить вообще)."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db.bans == []
|
||
assert db._by_id(1)["enabled"] is True # и глобально не тронут
|
||
|
||
|
||
def test_mark_banned_protects_last_live_node_logs_warning(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito")])
|
||
with caplog.at_level("WARNING"):
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert any("бан не записан" in rec.message for rec in caplog.records)
|
||
|
||
|
||
def test_mark_banned_protection_counts_active_bans_of_other_nodes() -> None:
|
||
"""Второй узел формально жив, но уже забанен ЭТИМ ЖЕ источником — кандидатом для
|
||
source он не является, значит бан первого узла оставил бы avito без прокси вообще.
|
||
Защита обязана сработать (иначе оба узла разом выпадут для avito)."""
|
||
assert mark_banned is not None
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="avito"), _proxy(2, affinity="any")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 2,
|
||
"source": "avito",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:avito",
|
||
}
|
||
],
|
||
)
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db._ban(1, "avito") is None # бан первого узла не записан
|
||
|
||
|
||
def test_mark_banned_writes_when_any_affinity_backup_exists() -> None:
|
||
"""Забанен единственный avito-специфичный узел, но есть 'any' — 'any' закрывает
|
||
availability для avito (та же семантика, что acquire()'s primary IN (provider,
|
||
'any')) → бан записывается безопасно."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db._ban(1, "avito") is not None
|
||
|
||
|
||
# ── mark_banned: affinity-aware last-node (orchestrator follow-up, свежий прод-факт) ──
|
||
#
|
||
# COUNT(*) WHERE enabled наивно посчитал бы "живых узлов много" даже когда
|
||
# конкретно для `source` не осталось НИ ОДНОГО — если единственные оставшиеся
|
||
# enabled-узлы это domclick (выделенная affinity, fallback НЕ имеет права её
|
||
# забрать при отсутствии backup — #2609). Проверяем ТОЧНО это расхождение.
|
||
|
||
|
||
def test_mark_banned_dedicated_affinity_alone_does_not_count_as_backup_for_other_source() -> None:
|
||
"""avito банится; в пуле остаётся только один domclick-узел (affinity выделенная,
|
||
БЕЗ backup) — для avito это НЕ доступный узел (acquire('avito') не взял бы его через
|
||
fallback, #2609 protects last node of domclick). Защита должна сработать — бан
|
||
НЕ записывается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="domclick")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db.bans == [] # domclick-узел НЕ считается доступной заменой
|
||
assert db._by_id(2)["enabled"] is True # и сам не тронут
|
||
|
||
|
||
def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None:
|
||
"""Та же ситуация, но у domclick есть ВТОРОЙ узел (backup) — тогда fallback может
|
||
забрать ОДИН из них под avito (acquire()'s EXISTS-правило #2609), доступность для
|
||
avito сохраняется через fallback → бан avito-узла записывается безопасно."""
|
||
assert mark_banned is not None
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, affinity="avito"),
|
||
_proxy(2, affinity="domclick"),
|
||
_proxy(3, affinity="domclick"),
|
||
]
|
||
)
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db._ban(1, "avito") is not None
|
||
|
||
|
||
def test_mark_banned_backup_of_dedicated_affinity_must_be_usable() -> None:
|
||
"""Зеркало acquire-правила в защите (deep-review #2600 п.2).
|
||
|
||
Узел 3 (domclick) забанен САМИМ domclick'ом и вдобавок в карантине по fails, т.е.
|
||
сам заменой для avito быть не может. Узел 2 — последний РАБОЧИЙ узел domclick,
|
||
fallback не имеет права его забрать. Значит замены для avito нет вообще → бан
|
||
avito-узла НЕ пишется. Со старым правилом («backup = любой enabled той же affinity»)
|
||
узел 3 засчитался бы бэкапом, узел 2 стал бы «доступным» и бан бы записался.
|
||
"""
|
||
assert mark_banned is not None
|
||
db = FakeSession(
|
||
[
|
||
_proxy(1, affinity="avito"),
|
||
_proxy(2, affinity="domclick"),
|
||
_proxy(3, affinity="domclick", fails=MAX_CONSECUTIVE_FAILS),
|
||
],
|
||
bans=[
|
||
{
|
||
"proxy_id": 3,
|
||
"source": "domclick",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:domclick",
|
||
}
|
||
],
|
||
)
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db._ban(1, "avito") is None
|
||
|
||
|
||
def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
|
||
"""Кандидат формально enabled, но consecutive_fails>=MAX_CONSECUTIVE_FAILS (карантин,
|
||
acquire() его не выдаёт) — НЕ считается доступной заменой, защита срабатывает."""
|
||
assert mark_banned is not None
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="avito"), _proxy(2, affinity="any", fails=MAX_CONSECUTIVE_FAILS)]
|
||
)
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db.bans == [] # карантинный узел не спасает
|
||
|
||
|
||
# ── mark_banned: TOCTOU-защита (deep-review fix 2, #2600) ───────────────────────
|
||
#
|
||
# Реальную гонку (два ПАРАЛЛЕЛЬНЫХ mark_banned на РАЗНЫХ proxy_id) честно юнитом не
|
||
# проверить — FakeSession однопоточна, а pg_advisory_xact_lock — свойство реальной
|
||
# СУБД (сериализация конкурентных транзакций), не что-то, что можно воспроизвести
|
||
# in-memory. Проверяем то, что юнитом ПРОВЕРИТЬ можно: лок реально берётся, с
|
||
# фиксированным ключом, ДО check+update (по SQL-подстроке — тот же паттерн, что и
|
||
# остальные тесты этого файла различают SQL веток).
|
||
|
||
|
||
def test_mark_banned_takes_advisory_xact_lock_with_fixed_key() -> None:
|
||
assert mark_banned is not None
|
||
from app.services.proxy_pool import _MARK_BANNED_ADVISORY_LOCK_KEY
|
||
|
||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db.advisory_lock_calls == [_MARK_BANNED_ADVISORY_LOCK_KEY]
|
||
|
||
|
||
def test_mark_banned_takes_advisory_lock_even_when_protected() -> None:
|
||
"""Лок берётся ПЕРЕД проверкой доступности — даже когда защита последнего узла
|
||
в итоге отменяет запись бана, лок всё равно взят (сериализация check+decide, не
|
||
только сам INSERT)."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="avito")]) # единственный узел — protected
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
assert db.advisory_lock_calls # лок взят, хотя бан не записан
|
||
assert db.bans == []
|
||
|
||
|
||
# ── acquire × бан по источнику (#2600 п.2 — суть задачи) ───────────────────────
|
||
#
|
||
# Прод-факт, из-за которого всё затевалось: Авито банит IP, Яндекс через тот же IP
|
||
# ходит чисто. Глобальный enabled=false выкидывал живой узел из пула для ВСЕХ
|
||
# источников; per-source бан обязан выключать выдачу ровно одному.
|
||
|
||
|
||
def test_acquire_skips_node_banned_for_this_source() -> None:
|
||
"""Единственный узел забанен ЭТИМ источником → acquire(source) возвращает None."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "avito",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:avito",
|
||
}
|
||
],
|
||
)
|
||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||
|
||
|
||
def test_acquire_still_issues_node_banned_by_other_source() -> None:
|
||
"""ГЛАВНЫЙ тест задачи: узел забанен Авито — Яндексу он выдаётся как ни в чём не
|
||
бывало (бан — свойство пары, а не узла)."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "avito",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:avito",
|
||
}
|
||
],
|
||
)
|
||
lease = acquire(db, "yandex", run_id=2) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 1
|
||
|
||
|
||
def test_acquire_ignores_expired_ban() -> None:
|
||
"""Истёкшая бан-строка (banned_until в прошлом) выдаче не мешает — она ещё лежит
|
||
только ради ban_count (purge снесёт её позже, SOURCE_BAN_PURGE_DAYS)."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="avito")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "avito",
|
||
"ban_count": 2,
|
||
"banned_until": datetime.now(UTC) - timedelta(hours=1),
|
||
"reason": "banned:avito",
|
||
}
|
||
],
|
||
)
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 1
|
||
|
||
|
||
def test_acquire_fallback_also_respects_source_ban() -> None:
|
||
"""Fallback-заход (чужая affinity) тоже отсекает забаненные для source узлы: два
|
||
cian-узла, один забанен avito → avito достаётся ВТОРОЙ, а не забаненный."""
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="cian"), _proxy(2, affinity="cian")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "avito",
|
||
"ban_count": 1,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||
"reason": "banned:avito",
|
||
}
|
||
],
|
||
)
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None
|
||
assert lease.id == 2
|
||
|
||
|
||
def test_ban_then_acquire_end_to_end() -> None:
|
||
"""Сквозной сценарий: mark_banned('avito') на узле 1 → avito получает узел 2,
|
||
yandex по-прежнему может получить узел 1."""
|
||
assert mark_banned is not None
|
||
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
|
||
avito_lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert avito_lease is not None and avito_lease.id == 2
|
||
|
||
release(db, 2) # type: ignore[arg-type]
|
||
yandex_lease = acquire(db, "yandex", run_id=2) # type: ignore[arg-type]
|
||
assert yandex_lease is not None and yandex_lease.id == 1 # забанен только для avito
|
||
|
||
|
||
# ── purge истёкших бан-строк в run_proxy_healthcheck (#2600 п.2) ───────────────
|
||
|
||
|
||
async def test_healthcheck_purges_long_expired_bans_only(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Сносим строки, истёкшие дольше SOURCE_BAN_PURGE_DAYS назад; истёкшую вчера
|
||
оставляем — в ней живёт ban_count (память об эскалации, см. комментарий у DELETE)."""
|
||
now = datetime.now(UTC)
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any")],
|
||
bans=[
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "avito",
|
||
"ban_count": 3,
|
||
"banned_until": now - timedelta(days=SOURCE_BAN_PURGE_DAYS + 1),
|
||
"reason": "banned:avito",
|
||
},
|
||
{
|
||
"proxy_id": 1,
|
||
"source": "cian",
|
||
"ban_count": 2,
|
||
"banned_until": now - timedelta(days=1),
|
||
"reason": "banned:cian",
|
||
},
|
||
],
|
||
)
|
||
|
||
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||
return True, "1.2.3.4", 10, None
|
||
|
||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||
|
||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||
|
||
assert counters["bans_purged"] == 1
|
||
assert [b["source"] for b in db.bans] == ["cian"]
|
||
|
||
|
||
# ── clear_source_bans (#2600 п.2 — рычаг оператора против ложного бана) ────────
|
||
#
|
||
# До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан
|
||
# в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой
|
||
# ручки ложное срабатывание детектора капчи (#2642) снималось бы только руками в SQL.
|
||
|
||
|
||
def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
|
||
return {
|
||
"proxy_id": pid,
|
||
"source": source,
|
||
"ban_count": ban_count,
|
||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS),
|
||
"reason": f"banned:{source}",
|
||
}
|
||
|
||
|
||
def test_clear_source_bans_removes_all_bans_of_node() -> None:
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||
bans=[_active_ban(1, "avito"), _active_ban(1, "cian"), _active_ban(2, "avito")],
|
||
)
|
||
cleared = proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
|
||
assert cleared == 2
|
||
assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел
|
||
# узел снова выдаётся источнику, который его банил
|
||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||
assert lease is not None and lease.id == 1
|
||
|
||
|
||
def test_clear_source_bans_single_source_keeps_others() -> None:
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any")],
|
||
bans=[_active_ban(1, "avito"), _active_ban(1, "cian")],
|
||
)
|
||
cleared = proxy_pool.clear_source_bans( # type: ignore[arg-type]
|
||
db, 1, source="avito", reason="ip rotated"
|
||
)
|
||
assert cleared == 1
|
||
assert [b["source"] for b in db.bans] == ["cian"]
|
||
|
||
|
||
def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
|
||
db = FakeSession([_proxy(1, affinity="any")])
|
||
assert proxy_pool.clear_source_bans(db, 1, reason="manual enable") == 0 # type: ignore[arg-type]
|
||
|
||
|
||
def test_clear_source_bans_resets_escalation() -> None:
|
||
"""DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан
|
||
начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию."""
|
||
assert mark_banned is not None
|
||
db = FakeSession(
|
||
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||
bans=[_active_ban(1, "avito", ban_count=4)],
|
||
)
|
||
proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
|
||
|
||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||
|
||
ban = db._ban(1, "avito")
|
||
assert ban["ban_count"] == 1
|
||
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS)
|
||
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|