fix(tradein/proxy): самовосстановление пула и запасной прокси чужой affinity (#2600) #2609

Merged
lekss361 merged 2 commits from fix/tradein-proxy-pool-self-healing into main 2026-08-01 18:54:02 +00:00
2 changed files with 374 additions and 48 deletions

View file

@ -15,11 +15,25 @@ ipify-пробу через каждый прокси и обновляет heal
за одну строку второй параллельный вызов пропустит залоченную и возьмёт следующую).
Health:
- mark_health(ok=True) consecutive_fails=0, last_ok_at/last_check_at, exit_ip, latency.
- mark_health(ok=True) consecutive_fails=0, enabled=true, last_ok_at/last_check_at,
exit_ip, latency. enabled=true реанимация: узел, выключенный
ранее авто-disable'ом, возвращается в строй первой же успешной
пробой (см. run_proxy_healthcheck).
- mark_health(ok=False) consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
авто-disable (enabled=false), чтобы битый узел выпал из пула.
- acquire отфильтровывает enabled=false И consecutive_fails >= MAX_FAILS.
Self-healing (#2600):
- run_proxy_healthcheck проверяет не только enabled-узлы, но и disabled реже, раз в
DISABLED_RECHECK_MINUTES (или если ни разу не проверялся). Успешная проба выключенного
узла реанимирует его (enabled=true), инкрементит счётчик `revived` и пишет INFO-лог.
Без этого auto-disable необратим: транзиентный сбой = вечный приговор узлу.
- acquire, не найдя свободного здорового узла нужной provider_affinity, вторым заходом
берёт любой свободный здоровый узел ЛЮБОЙ affinity (WARNING-лог) иначе источник
голодает при живых свободных узлах чужой affinity. Fallback НЕ забирает последний
enabled-узел выделенной affinity (пример domclick, один узел на всё, см. acquire
docstring) иначе чинили бы один источник ценой полной поломки другого.
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
"""
@ -36,6 +50,7 @@ from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
__all__ = [
"DISABLED_RECHECK_MINUTES",
"DISABLE_THRESHOLD",
"MAX_CONSECUTIVE_FAILS",
"NON_RUN_LEASE_MARKER",
@ -62,6 +77,12 @@ DISABLE_THRESHOLD = 5
# освобождается reap_stale_leases — иначе прокси навсегда «занят» мёртвым run'ом.
STALE_LEASE_MINUTES = 30
# Disabled-узлы перепроверяются не каждый прогон (это долбёж по мёртвому/дорогому
# провайдеру), а раз в это число минут — либо если ни разу не проверялся. Успешная
# проба реанимирует узел (см. run_proxy_healthcheck). Без recheck'а auto-disable
# необратим: транзиентный сбой = вечный приговор (#2600).
DISABLED_RECHECK_MINUTES = 60
# Маркер lease для не-run вызовов (leased_by NOT NULL = занят, но это не id из scrape_runs).
NON_RUN_LEASE_MARKER = -1
@ -90,10 +111,25 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
(last_ok_at NULLS LAST). Затем помечает строку leased_by=run_id (или
NON_RUN_LEASE_MARKER если run_id не задан) и коммитит.
Если свободных здоровых узлов нужной affinity (provider/'any') нет вторым заходом
берётся любой свободный здоровый узел ЛЮБОЙ affinity (тот же ORDER BY/FOR UPDATE SKIP
LOCKED), с WARNING-логом. Приоритет не меняется: своя affinity всегда предпочтительнее,
чужая только запасной вариант, чтобы источник не голодал при живых свободных узлах
чужой affinity (#2600).
Fallback НЕ трогает последний enabled-узел выделенной (не-'any') affinity см.
173_scrape_proxies_add_domclick_affinity.sql: у domclick ровно один узел (id=1),
намеренно вырезанный из общего пула, потому что QRATOR банит все прокси кроме этого
одного чистого residential-адреса. Если fallback заберёт его под avito/cian/yandex,
domclick останется без прокси вообще хуже, чем голодание исходного источника,
которое фикс призван устранить. Кандидат участвует в fallback, только если его
affinity='any' ИЛИ у этой affinity есть ДРУГОЙ enabled-узел (EXISTS-подзапрос)
т.е. выдача не обнулит доступность выделенной affinity целиком.
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
Returns ProxyLease или None если свободных здоровых прокси нет.
Returns ProxyLease или None если свободных здоровых прокси нет вообще.
"""
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
@ -117,6 +153,44 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
.mappings()
.fetchone()
)
fallback_used = False
if row is None:
# Нет своих (provider/'any') — запасной заход: любой свободный здоровый узел
# ЛЮБОЙ affinity, кроме последнего enabled-узла выделенной affinity (domclick и
# т.п.) — EXISTS-подзапрос требует хотя бы ОДИН ДРУГОЙ enabled-узел той же
# affinity, иначе affinity='any' достаточно.
row = (
db.execute(
text(
"""
SELECT sp.id, sp.url, sp.kind, sp.rotate_url
FROM scrape_proxies AS sp
WHERE sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer)
AND sp.leased_by IS NULL
AND (
sp.provider_affinity = 'any'
OR EXISTS (
SELECT 1
FROM scrape_proxies AS other
WHERE other.provider_affinity = sp.provider_affinity
AND other.enabled
AND other.id <> sp.id
)
)
ORDER BY sp.last_ok_at NULLS LAST, sp.id
FOR UPDATE SKIP LOCKED
LIMIT 1
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS},
)
.mappings()
.fetchone()
)
fallback_used = row is not None
if row is None:
db.rollback() # снять FOR UPDATE-транзакцию (ничего не залочено, но чисто)
return None
@ -133,9 +207,18 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
{"run_id": lease_marker, "id": proxy_id},
)
db.commit()
logger.info(
"proxy_pool: leased proxy id=%d provider=%s by=%s", proxy_id, provider, lease_marker
)
if fallback_used:
logger.warning(
"proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity "
"(no free healthy proxy of matching affinity, issuing proxy of other affinity)",
proxy_id,
provider,
lease_marker,
)
else:
logger.info(
"proxy_pool: leased proxy id=%d provider=%s by=%s", proxy_id, provider, lease_marker
)
return ProxyLease(
id=proxy_id,
url=str(row["url"]),
@ -167,13 +250,24 @@ def mark_health(
*,
exit_ip: str | None = None,
latency_ms: int | None = None,
fail_kind: str | None = None,
) -> None:
"""Записать результат health-check'а прокси.
ok=True consecutive_fails обнуляется, обновляются last_ok_at/last_check_at/
exit_ip/latency_ms.
ok=True consecutive_fails обнуляется, enabled=true, обновляются last_ok_at/
last_check_at/exit_ip/latency_ms. enabled=true безусловно это реанимация:
узел, ранее выключенный auto-disable'ом, возвращается в строй первой же
успешной пробой (см. run_proxy_healthcheck, #2600 п.1).
ok=False consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
авто-disable (enabled=false). last_check_at обновляется в любом случае.
fail_kind необязательная классификация неуспеха ("timeout" / "connect_error" /
"http_error" / "other", см. _probe_proxy), используется ТОЛЬКО для логирования.
Счётчик consecutive_fails/порог disable инкрементится одинаково для любого fail_kind
аккуратное разделение "транзиентный сбой vs перманентный бан" (разные пороги/скорость
инкремента по типу ошибки) требует более глубокой переработки модуля (отдельный
трекинг по типам ошибок, вероятно per-fail_kind счётчики) и намеренно НЕ сделано в
рамках #2600 п.2 — см. обоснование в PR. fail_kind — задел под это на будущее.
"""
if ok:
db.execute(
@ -185,6 +279,7 @@ def mark_health(
last_check_at = now(),
exit_ip = CAST(:exit_ip AS text),
latency_ms = CAST(:latency_ms AS integer),
enabled = true,
updated_at = now()
WHERE id = CAST(:id AS bigint)
"""
@ -210,7 +305,13 @@ def mark_health(
{"disable_threshold": DISABLE_THRESHOLD, "id": proxy_id},
)
db.commit()
logger.info("proxy_pool: mark_health id=%d ok=%s exit_ip=%s", proxy_id, ok, exit_ip)
logger.info(
"proxy_pool: mark_health id=%d ok=%s exit_ip=%s fail_kind=%s",
proxy_id,
ok,
exit_ip,
fail_kind,
)
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
@ -237,10 +338,17 @@ def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES
return len(rows)
async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None]:
async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None, str | None]:
"""GET ipify через прокси (timeout _HEALTH_PROBE_TIMEOUT_S).
Returns (ok, exit_ip, latency_ms). ok=False + (None, None) при любой ошибке.
Returns (ok, exit_ip, latency_ms, fail_kind). При успехе fail_kind=None. При неуспехе
exit_ip/latency_ms=None, а fail_kind классифицирует что случилось (#2600 п.2 —
транзиентный сбой узла перманентный бан, используется пока только для логов):
- "timeout" сеть недоступна/медленная (httpx.TimeoutException)
- "connect_error" прокси не поднят/не слушает/DNS (httpx.ConnectError)
- "http_error" ipify ответил ошибкой через прокси (auth/upstream)
- "other" прочее
url несёт схему (http:// / socks5://) httpx[socks] обрабатывает оба.
"""
started = time.monotonic()
@ -250,10 +358,23 @@ async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None]:
resp.raise_for_status()
ip = resp.json().get("ip")
latency_ms = int((time.monotonic() - started) * 1000)
return True, (str(ip) if ip else None), latency_ms
return True, (str(ip) if ip else None), latency_ms, None
except httpx.TimeoutException:
logger.warning("proxy_pool: health probe timeout proxy=%s", _mask(url))
return False, None, None, "timeout"
except httpx.ConnectError:
logger.warning("proxy_pool: health probe connect_error proxy=%s", _mask(url))
return False, None, None, "connect_error"
except httpx.HTTPStatusError as exc:
logger.warning(
"proxy_pool: health probe http_error proxy=%s status=%s",
_mask(url),
exc.response.status_code,
)
return False, None, None, "http_error"
except Exception:
logger.warning("proxy_pool: health probe failed proxy=%s", _mask(url), exc_info=True)
return False, None, None
return False, None, None, "other"
def _mask(url: str) -> str:
@ -269,16 +390,23 @@ def _mask(url: str) -> str:
async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
"""Периодический health-check всех enabled-прокси пула (#2162).
"""Периодический health-check прокси пула — enabled каждый прогон, disabled реже (#2162, #2600).
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем для каждого
enabled-прокси гоняет ipify-пробу через сам прокси и пишет результат через
mark_health (успех сброс fails + свежий exit_ip/latency; фейл инкремент,
авто-disable при DISABLE_THRESHOLD).
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем гоняет ipify-пробу
через каждый кандидат и пишет результат через mark_health (успех сброс fails +
enabled=true + свежий exit_ip/latency; фейл инкремент, авто-disable при
DISABLE_THRESHOLD).
Кандидаты: ВСЕ enabled-узлы (как раньше) + disabled-узлы, которые ни разу не
проверялись (last_check_at IS NULL) или проверялись давнее DISABLED_RECHECK_MINUTES
назад. Без этого auto-disable необратим узел, ушедший в disable из-за транзиентного
сбоя, никогда больше не проверяется и не может вернуться (#2600 п.1). Успешная проба
disabled-узла реанимирует его (enabled=true через mark_health) инкрементит `revived`
и пишет отдельный INFO-лог.
Пробы идут последовательно пул небольшой (десятки узлов), а параллельный залп на
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
{reaped, checked, ok, failed}.
{reaped, checked, ok, failed, revived}.
"""
reaped = reap_stale_leases(db)
@ -286,12 +414,17 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
db.execute(
text(
"""
SELECT id, url, kind
SELECT id, url, kind, enabled
FROM scrape_proxies
WHERE enabled
OR last_check_at IS NULL
OR last_check_at < now() - make_interval(
mins => CAST(:disabled_recheck_minutes AS integer)
)
ORDER BY id
"""
)
),
{"disabled_recheck_minutes": DISABLED_RECHECK_MINUTES},
)
.mappings()
.all()
@ -300,22 +433,38 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
checked = 0
ok_count = 0
failed = 0
revived = 0
for row in proxies:
proxy_id = int(row["id"])
url = str(row["url"])
ok, exit_ip, latency_ms = await _probe_proxy(url)
mark_health(db, proxy_id, ok, exit_ip=exit_ip, latency_ms=latency_ms)
was_disabled = not bool(row["enabled"])
ok, exit_ip, latency_ms, fail_kind = await _probe_proxy(url)
mark_health(db, proxy_id, ok, exit_ip=exit_ip, latency_ms=latency_ms, fail_kind=fail_kind)
checked += 1
if ok:
ok_count += 1
if was_disabled:
revived += 1
logger.info(
"proxy_pool: REVIVED proxy id=%d — successful probe of a disabled node, "
"returned to service (enabled=true, consecutive_fails=0)",
proxy_id,
)
else:
failed += 1
logger.info(
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d",
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d revived=%d",
reaped,
checked,
ok_count,
failed,
revived,
)
return {"reaped": reaped, "checked": checked, "ok": ok_count, "failed": failed}
return {
"reaped": reaped,
"checked": checked,
"ok": ok_count,
"failed": failed,
"revived": revived,
}

View file

@ -1,4 +1,4 @@
"""Offline-тесты пула прокси (#2162).
"""Offline-тесты пула прокси (#2162, #2600).
Покрытие БЕЗ live-сети/БД: stateful FakeSession эмулирует таблицу scrape_proxies и
интерпретирует SQL по ключевым фрагментам, так что acquire/release/mark_health/
@ -8,11 +8,15 @@ reap_stale_leases проверяются по фактическому изме
- два acquire подряд РАЗНЫЕ прокси (первый лизнут выпал из выборки второго).
- release освобождает (leased_by NULL), прокси снова acquire-абелен.
- mark_health fail инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
- mark_health ok сброс fails + exit_ip/latency.
- 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
@ -29,6 +33,7 @@ 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,
@ -71,17 +76,44 @@ class FakeSession:
sql = str(stmt)
p = params or {}
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT
provider = p["provider"]
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
max_fails = p["max_fails"]
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
]
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: (
@ -127,18 +159,33 @@ class FakeSession:
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
rows = sorted((r for r in self.rows if r["enabled"]), key=lambda r: r["id"])
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}")
@ -159,6 +206,7 @@ def _proxy(
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]:
@ -173,6 +221,7 @@ def _proxy(
"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,
}
@ -209,11 +258,6 @@ def test_acquire_empty_pool_returns_none() -> None:
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
def test_acquire_affinity_filter_excludes_other_provider() -> None:
db = FakeSession([_proxy(1, affinity="cian")])
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]
@ -231,6 +275,62 @@ def test_acquire_without_run_id_uses_marker() -> 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 ──────────────────────────────────────────────────────────────────
@ -271,6 +371,15 @@ def test_mark_health_ok_resets_and_records() -> None:
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 ────────────────────────────────────────────────────────
@ -295,27 +404,95 @@ def test_reap_frees_stale_lease_keeps_fresh() -> None:
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),
_proxy(2, enabled=False), # disabled — не проверяется
# 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]:
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
return False, None, None
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)
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