gendesign/tradein-mvp/backend/tests/services/test_proxy_pool.py
bot-backend b89788ee99 feat(tradein/proxy): прогон знает свой узел, а снятый бан перестаёт стирать историю (#3404)
Выбор оператора мобильного прокси опирался на две ненадёжные опоры.

Первая: `scrape_runs` не знала, через какой узел шёл прогон — колонка `proxy_id`
была только у банов и ротаций. «Какой узел собрал 5 карточек из 21» не выяснялось
ни одним запросом.

Вторая: `clear_source_bans` делала DELETE, а зовётся она после КАЖДОЙ успешной
ротации exit-IP. У #540723 (МегаФон) 23 успешные ротации и ноль строк банов,
у #540722 (Tele2) ротаций почти не было и 7 банов. «7 против 0» читалось как
«Tele2 хуже», хотя в той же мере это «у МегаФона историю стёрли 23 раза».

Теперь:
- `scrape_runs.proxy_id` — последний выданный прогону узел; полная цепочка
  (если узел менялся mid-run) копится в `counters.proxy_ids`. Пишет
  `proxy_pool.attribute_run_proxy` из единственной точки — сразу после выдачи
  лиза в `acquire()`, поэтому curl-путь, браузерный sticky lease и ре-acquire
  при ротации покрыты одинаково. `run_id` доходит до адаптера через ContextVar
  (`scraper_kit.orchestration.run_context`): протокол `ProxyProvider.acquire`
  его не несёт, а `RealProxyProvider` живёт одним объектом на весь планировщик.
  Best-effort: `lock_timeout` 2с и проглоченное исключение — диагностика не
  вправе ронять выдачу прокси или ждать на блокировке строки прогона.
- `clear_source_bans` гасит строку (`banned_until = now()`, `ban_count = 0`,
  `cleared_at`/`cleared_reason`) вместо удаления. Эскалация сохраняется 1:1:
  формула в `mark_banned` берёт ПРЕДЫДУЩИЙ `ban_count` показателем степени, при
  нуле это ровно `SOURCE_BAN_BASE_HOURS` — как после DELETE. Строка доживает до
  штатного purge по `SOURCE_BAN_PURGE_DAYS`.

Для всех читателей `scrape_proxy_source_bans` погашенная строка неотличима от
отсутствующей: acquire, оба guard-подзапроса `mark_banned`, `proxy_egress`
(ранжирование по `ban_count` даёт 0, как у узла без истории), admin `_active_ban` —
все гейтятся по `banned_until > now()`.

Ничего не бэкфиллится: связать прошедшие прогоны с узлами нечем (`leased_by`
исторически = NON_RUN_LEASE_MARKER), врать восстановленным значением нельзя.

Миграция 287. Тесты: 9 новых на обе части (главный — эскалация после гашения даёт
базовые 6ч, а не удвоенные) + 14 существующих переведены с DELETE-семантики на
гашение, включая проверку, что секрет ротации не утекает в новое `cleared_reason`.
Полный прогон бэкенда: 5600 passed, 37 skipped.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011WHFxVPWoBnSZihkdH1Uou
2026-09-06 13:05:20 +03:00

1529 lines
79 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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
# Зеркалит COALESCE(:param, <колонка>) реального SQL (#3283):
# mark_health(lease, ok) без exit_ip/latency_ms приходит из
# _report_fetch_result на каждый успешный /fetch, и None не должен
# затирать адрес, записанный ротацией.
if p["exit_ip"] is not None:
row["exit_ip"] = p["exit_ip"]
if p["latency_ms"] is not None:
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 "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql:
# clear_source_bans (#3404): гасим строку (banned_until=now(), ban_count=0,
# cleared_at/cleared_reason) вместо DELETE — трассируемость снятия бана, строка
# доживает до штатного purge. Фильтр по 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
now = datetime.now(UTC)
cleared: list[dict[str, Any]] = []
for b in self.bans:
if b["proxy_id"] != p["proxy_id"]:
continue
if p["source"] is not None and b["source"] != p["source"]:
continue
if only_reason is not None and b.get("reason") != only_reason:
continue
# Гейт покоя (реальный WHERE): уже погашенную строку повторно не трогаем —
# без него повторный вызов сдвигал бы banned_until вперёд.
if (
b.get("cleared_at") is not None
and b["ban_count"] == 0
and b["banned_until"] <= now
):
continue
b["banned_until"] = now
b["ban_count"] = 0
b["cleared_at"] = now
b["cleared_reason"] = p["reason"]
b["updated_at"] = now
cleared.append(b)
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_without_exit_ip_keeps_stored_value() -> None:
"""#3283: успешный /fetch зовёт mark_health(lease, ok) БЕЗ exit_ip/latency_ms.
До фикса None затирал адрес, записанный proxy_rotation._update_exit_ip, и
scrape_proxies.exit_ip жил до первого же успешного запроса (прод 31.08: у обоих
живых узлов NULL, при том что в логе ротации new_ip=31.173.86.74).
"""
db = FakeSession([_proxy(1)])
mark_health(db, 1, ok=True, exit_ip="31.173.86.74", latency_ms=120) # type: ignore[arg-type]
mark_health(db, 1, ok=True) # type: ignore[arg-type] # как из _report_fetch_result
row = db._by_id(1)
assert row["exit_ip"] == "31.173.86.74"
assert row["latency_ms"] == 120
def test_mark_health_ok_with_exit_ip_overwrites_stored_value() -> None:
"""Явное значение по-прежнему пишется — COALESCE не превращает поле в append-only."""
db = FakeSession([_proxy(1)])
mark_health(db, 1, ok=True, exit_ip="1.2.3.4", latency_ms=88) # type: ignore[arg-type]
mark_health(db, 1, ok=True, exit_ip="5.6.7.8", latency_ms=99) # type: ignore[arg-type]
row = db._by_id(1)
assert row["exit_ip"] == "5.6.7.8"
assert row["latency_ms"] == 99
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.
#
# #3404: снятие гасит строку (banned_until=now(), ban_count=0, cleared_at/cleared_reason),
# а не удаляет её — строка живёт для трассируемости до штатного purge, но для выдачи и
# для эскалации следующего бана неотличима от прежнего DELETE.
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}",
"cleared_at": None,
"cleared_reason": None,
}
def test_clear_source_bans_gates_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
# #3404: строки НЕ удаляются — все три остаются в таблице (трассируемость).
assert {(b["proxy_id"], b["source"]) for b in db.bans} == {
(1, "avito"),
(1, "cian"),
(2, "avito"),
}
now = datetime.now(UTC)
node1_bans = [b for b in db.bans if b["proxy_id"] == 1]
assert len(node1_bans) == 2
for b in node1_bans:
assert b["banned_until"] <= now
assert b["ban_count"] == 0
assert b["cleared_at"] is not None
assert b["cleared_reason"] == "manual enable"
# чужой бан (proxy_id=2) не тронут — остаётся активным
other_ban = next(b for b in db.bans if b["proxy_id"] == 2)
assert other_ban["banned_until"] > now
assert other_ban["ban_count"] == 1
assert other_ban["cleared_at"] is None
# узел снова выдаётся источнику, который его банил
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
# обе строки остаются (#3404), но только "avito" погашена
avito_ban = db._ban(1, "avito")
cian_ban = db._ban(1, "cian")
now = datetime.now(UTC)
assert avito_ban["banned_until"] <= now
assert avito_ban["ban_count"] == 0
assert avito_ban["cleared_reason"] == "ip rotated"
assert cian_ban["banned_until"] > now
assert cian_ban["ban_count"] == 1
assert cian_ban["cleared_at"] is None
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:
"""Гашение (banned_until=now(), ban_count=0), а не удаление строки: следующий бан той
же пары начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию (#3404)."""
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