gendesign/tradein-mvp/backend/tests/services/test_proxy_pool.py
bot-backend 876b666424
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m44s
fix(tradein/proxy): не отдавать в fallback последний узел выделенной affinity (#2600)
Ревью PR #2609: domclick — ровно один узел (прод scrape_proxies.id=1),
намеренно вырезанный из общего пула через provider_affinity='domclick'
(см. 173_scrape_proxies_add_domclick_affinity.sql) — QRATOR банит всё,
кроме этого одного чистого residential-адреса. Fallback-запрос из
предыдущего коммита мог законно забрать его под avito/cian/yandex,
оставив domclick (сейчас исправно собирает: 6501 активных объявлений,
368/сутки) без прокси вообще — чинили бы один источник ценой полной
поломки другого.

- acquire(): fallback-SELECT дополнен условием "affinity='any' ИЛИ есть
  ДРУГОЙ enabled-узел той же affinity" через коррелированный EXISTS-
  подзапрос (WHERE + FOR UPDATE SKIP LOCKED + ORDER BY last_ok_at NULLS
  LAST, id — сохранены). Кандидат с единственным enabled-узлом своей
  выделенной affinity в fallback не участвует.
- Тесты: единственный domclick-узел → acquire('avito') возвращает None;
  второй enabled domclick-узел появляется — fallback снова срабатывает.
- Починен мок FakeSession (tests/services/test_proxy_pool.py):
  ветка "mark_health ok" раньше ставила enabled=True безусловно по
  совпадению общей подстроки "SET consecutive_fails = 0" (одинаковой в
  старом и новом SQL) — test_mark_health_ok_revives_disabled_proxy
  проходил бы и против кода без реанимации. Теперь ставит enabled=True
  только если в тексте SQL реально есть "enabled". Та же проблема была
  и в fallback-ветке (protects_last_node переопределял логику в Python
  независимо от SQL) — исправлено аналогично: применяется, только если
  в SQL реально есть EXISTS-подзапрос.
2026-08-01 21:36:02 +03:00

498 lines
23 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, свежий не трогает.
- 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++;
* недавно проверенный выключенный узел повторно не проверяется (не долбим провайдера).
"""
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,
acquire,
mark_health,
reap_stale_leases,
release,
)
# ── 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-строк."""
def __init__(self, rows: list[dict[str, Any]]):
self.rows = rows
# 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 execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
p = params or {}
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
max_fails = p["max_fails"]
if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any'
provider = p["provider"]
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
]
else: # fallback: любая affinity, но не последний узел выделенной affinity
# (domclick и т.п. — #2600 review). ВАЖНО: применяем эту фильтрацию,
# только если сама SQL реально содержит защиту (EXISTS-подзапрос) —
# иначе мок реализовывал бы бизнес-логику независимо от проверяемого
# кода и не смог бы отличить старый (незащищённый) fallback-запрос от
# нового. Тот же класс бага, что был с "enabled" в mark_health-моке.
protects_last_node = "EXISTS" 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"]
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 protects_last_node or _has_backup(r))
]
# ORDER BY last_ok_at NULLS LAST, id
cands.sort(
key=lambda r: (
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 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)
# "SET consecutive_fails = 0" — общая подстрока старого И нового SQL,
# НЕ различает их сама по себе. Реанимация (enabled=true) — только если
# в тексте запроса реально есть присвоение enabled (#2600 review: старый
# мок ставил enabled=True безусловно и не ловил регресс).
if "enabled" in sql:
row["enabled"] = True
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"])
return _FakeResult([dict(r) for r in rows])
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,
) -> 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,
"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,
}
# ── 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 остался живой запасной узел
# ── 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
# ── 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 # свежий не тронут
# ── 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_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