fix(tradein/proxy): самовосстановление пула и запасной прокси чужой affinity (#2600)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 10s
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 2m42s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 10s
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 2m42s
Прод-замер: disabled-узлы никогда не перепроверялись (WHERE enabled в run_proxy_healthcheck) — auto-disable по DISABLE_THRESHOLD необратим, транзиентный сбой = вечный приговор (id 11 сгорел за ночь, будучи физически исправным). acquire() при пустой выборке по provider_affinity падал в None, морив источник голодом при живых свободных узлах чужой affinity. - run_proxy_healthcheck: disabled-узлы проверяются реже (DISABLED_RECHECK_MINUTES=60 либо last_check_at IS NULL); успешная проба реанимирует узел (enabled=true через mark_health) и инкрементит новый счётчик revived. - mark_health(ok=True) теперь безусловно ставит enabled=true (реанимация). - acquire: вторым заходом при пустой выборке своей affinity берёт любой свободный здоровый узел любой affinity (WARNING-лог), приоритет своих сохранён. - _probe_proxy классифицирует неуспех (timeout/connect_error/http_error/other) в fail_kind — прокидывается в mark_health только для логирования; полноценное разделение порогов транзиент/бан отложено (см. docstring mark_health).
This commit is contained in:
parent
34b346b097
commit
ad753c6a87
2 changed files with 309 additions and 45 deletions
|
|
@ -15,11 +15,23 @@ ipify-пробу через каждый прокси и обновляет heal
|
||||||
за одну строку — второй параллельный вызов пропустит залоченную и возьмёт следующую).
|
за одну строку — второй параллельный вызов пропустит залоченную и возьмёт следующую).
|
||||||
|
|
||||||
Health:
|
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 прокси
|
- mark_health(ok=False) → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
|
||||||
авто-disable (enabled=false), чтобы битый узел выпал из пула.
|
авто-disable (enabled=false), чтобы битый узел выпал из пула.
|
||||||
- acquire отфильтровывает enabled=false И consecutive_fails >= MAX_FAILS.
|
- 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.
|
||||||
|
|
||||||
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
|
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
|
@ -36,6 +48,7 @@ from sqlalchemy.orm import Session
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
|
"DISABLED_RECHECK_MINUTES",
|
||||||
"DISABLE_THRESHOLD",
|
"DISABLE_THRESHOLD",
|
||||||
"MAX_CONSECUTIVE_FAILS",
|
"MAX_CONSECUTIVE_FAILS",
|
||||||
"NON_RUN_LEASE_MARKER",
|
"NON_RUN_LEASE_MARKER",
|
||||||
|
|
@ -62,6 +75,12 @@ DISABLE_THRESHOLD = 5
|
||||||
# освобождается reap_stale_leases — иначе прокси навсегда «занят» мёртвым run'ом.
|
# освобождается reap_stale_leases — иначе прокси навсегда «занят» мёртвым run'ом.
|
||||||
STALE_LEASE_MINUTES = 30
|
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).
|
# Маркер lease для не-run вызовов (leased_by NOT NULL = занят, но это не id из scrape_runs).
|
||||||
NON_RUN_LEASE_MARKER = -1
|
NON_RUN_LEASE_MARKER = -1
|
||||||
|
|
||||||
|
|
@ -90,10 +109,16 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
(last_ok_at NULLS LAST). Затем помечает строку leased_by=run_id (или
|
(last_ok_at NULLS LAST). Затем помечает строку leased_by=run_id (или
|
||||||
NON_RUN_LEASE_MARKER если run_id не задан) и коммитит.
|
NON_RUN_LEASE_MARKER если run_id не задан) и коммитит.
|
||||||
|
|
||||||
|
Если свободных здоровых узлов нужной affinity (provider/'any') нет — вторым заходом
|
||||||
|
берётся любой свободный здоровый узел ЛЮБОЙ affinity (тот же ORDER BY/FOR UPDATE SKIP
|
||||||
|
LOCKED), с WARNING-логом. Приоритет не меняется: своя affinity всегда предпочтительнее,
|
||||||
|
чужая — только запасной вариант, чтобы источник не голодал при живых свободных узлах
|
||||||
|
чужой affinity (#2600).
|
||||||
|
|
||||||
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
||||||
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
||||||
|
|
||||||
Returns ProxyLease или None если свободных здоровых прокси нет.
|
Returns ProxyLease или None если свободных здоровых прокси нет вообще.
|
||||||
"""
|
"""
|
||||||
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
|
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
|
||||||
|
|
||||||
|
|
@ -117,6 +142,33 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
.mappings()
|
.mappings()
|
||||||
.fetchone()
|
.fetchone()
|
||||||
)
|
)
|
||||||
|
|
||||||
|
fallback_used = False
|
||||||
|
if row is None:
|
||||||
|
# Нет своих (provider/'any') — запасной заход: любой свободный здоровый узел,
|
||||||
|
# affinity не важна. Лучше выдать источнику чужой прокси, чем оставить его без
|
||||||
|
# прокси при живых свободных узлах.
|
||||||
|
row = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT id, url, kind, rotate_url
|
||||||
|
FROM scrape_proxies
|
||||||
|
WHERE enabled
|
||||||
|
AND consecutive_fails < CAST(:max_fails AS integer)
|
||||||
|
AND leased_by IS NULL
|
||||||
|
ORDER BY last_ok_at NULLS LAST, id
|
||||||
|
FOR UPDATE SKIP LOCKED
|
||||||
|
LIMIT 1
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"max_fails": MAX_CONSECUTIVE_FAILS},
|
||||||
|
)
|
||||||
|
.mappings()
|
||||||
|
.fetchone()
|
||||||
|
)
|
||||||
|
fallback_used = row is not None
|
||||||
|
|
||||||
if row is None:
|
if row is None:
|
||||||
db.rollback() # снять FOR UPDATE-транзакцию (ничего не залочено, но чисто)
|
db.rollback() # снять FOR UPDATE-транзакцию (ничего не залочено, но чисто)
|
||||||
return None
|
return None
|
||||||
|
|
@ -133,9 +185,18 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
{"run_id": lease_marker, "id": proxy_id},
|
{"run_id": lease_marker, "id": proxy_id},
|
||||||
)
|
)
|
||||||
db.commit()
|
db.commit()
|
||||||
logger.info(
|
if fallback_used:
|
||||||
"proxy_pool: leased proxy id=%d provider=%s by=%s", proxy_id, provider, lease_marker
|
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(
|
return ProxyLease(
|
||||||
id=proxy_id,
|
id=proxy_id,
|
||||||
url=str(row["url"]),
|
url=str(row["url"]),
|
||||||
|
|
@ -167,13 +228,24 @@ def mark_health(
|
||||||
*,
|
*,
|
||||||
exit_ip: str | None = None,
|
exit_ip: str | None = None,
|
||||||
latency_ms: int | None = None,
|
latency_ms: int | None = None,
|
||||||
|
fail_kind: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Записать результат health-check'а прокси.
|
"""Записать результат health-check'а прокси.
|
||||||
|
|
||||||
ok=True → consecutive_fails обнуляется, обновляются last_ok_at/last_check_at/
|
ok=True → consecutive_fails обнуляется, enabled=true, обновляются last_ok_at/
|
||||||
exit_ip/latency_ms.
|
last_check_at/exit_ip/latency_ms. enabled=true безусловно — это реанимация:
|
||||||
|
узел, ранее выключенный auto-disable'ом, возвращается в строй первой же
|
||||||
|
успешной пробой (см. run_proxy_healthcheck, #2600 п.1).
|
||||||
ok=False → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
|
ok=False → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
|
||||||
авто-disable (enabled=false). last_check_at обновляется в любом случае.
|
авто-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:
|
if ok:
|
||||||
db.execute(
|
db.execute(
|
||||||
|
|
@ -185,6 +257,7 @@ def mark_health(
|
||||||
last_check_at = now(),
|
last_check_at = now(),
|
||||||
exit_ip = CAST(:exit_ip AS text),
|
exit_ip = CAST(:exit_ip AS text),
|
||||||
latency_ms = CAST(:latency_ms AS integer),
|
latency_ms = CAST(:latency_ms AS integer),
|
||||||
|
enabled = true,
|
||||||
updated_at = now()
|
updated_at = now()
|
||||||
WHERE id = CAST(:id AS bigint)
|
WHERE id = CAST(:id AS bigint)
|
||||||
"""
|
"""
|
||||||
|
|
@ -210,7 +283,13 @@ def mark_health(
|
||||||
{"disable_threshold": DISABLE_THRESHOLD, "id": proxy_id},
|
{"disable_threshold": DISABLE_THRESHOLD, "id": proxy_id},
|
||||||
)
|
)
|
||||||
db.commit()
|
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:
|
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
|
||||||
|
|
@ -237,10 +316,17 @@ def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES
|
||||||
return len(rows)
|
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).
|
"""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] обрабатывает оба.
|
url несёт схему (http:// / socks5://) — httpx[socks] обрабатывает оба.
|
||||||
"""
|
"""
|
||||||
started = time.monotonic()
|
started = time.monotonic()
|
||||||
|
|
@ -250,10 +336,23 @@ async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None]:
|
||||||
resp.raise_for_status()
|
resp.raise_for_status()
|
||||||
ip = resp.json().get("ip")
|
ip = resp.json().get("ip")
|
||||||
latency_ms = int((time.monotonic() - started) * 1000)
|
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:
|
except Exception:
|
||||||
logger.warning("proxy_pool: health probe failed proxy=%s", _mask(url), exc_info=True)
|
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:
|
def _mask(url: str) -> str:
|
||||||
|
|
@ -269,16 +368,23 @@ def _mask(url: str) -> str:
|
||||||
|
|
||||||
|
|
||||||
async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
"""Периодический health-check всех enabled-прокси пула (#2162).
|
"""Периодический health-check прокси пула — enabled каждый прогон, disabled реже (#2162, #2600).
|
||||||
|
|
||||||
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем для каждого
|
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем гоняет ipify-пробу
|
||||||
enabled-прокси гоняет ipify-пробу через сам прокси и пишет результат через
|
через каждый кандидат и пишет результат через mark_health (успех → сброс fails +
|
||||||
mark_health (успех → сброс fails + свежий exit_ip/latency; фейл → инкремент,
|
enabled=true + свежий exit_ip/latency; фейл → инкремент, авто-disable при
|
||||||
авто-disable при DISABLE_THRESHOLD).
|
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
|
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
|
||||||
{reaped, checked, ok, failed}.
|
{reaped, checked, ok, failed, revived}.
|
||||||
"""
|
"""
|
||||||
reaped = reap_stale_leases(db)
|
reaped = reap_stale_leases(db)
|
||||||
|
|
||||||
|
|
@ -286,12 +392,17 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
SELECT id, url, kind
|
SELECT id, url, kind, enabled
|
||||||
FROM scrape_proxies
|
FROM scrape_proxies
|
||||||
WHERE enabled
|
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
|
ORDER BY id
|
||||||
"""
|
"""
|
||||||
)
|
),
|
||||||
|
{"disabled_recheck_minutes": DISABLED_RECHECK_MINUTES},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.all()
|
.all()
|
||||||
|
|
@ -300,22 +411,38 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
checked = 0
|
checked = 0
|
||||||
ok_count = 0
|
ok_count = 0
|
||||||
failed = 0
|
failed = 0
|
||||||
|
revived = 0
|
||||||
for row in proxies:
|
for row in proxies:
|
||||||
proxy_id = int(row["id"])
|
proxy_id = int(row["id"])
|
||||||
url = str(row["url"])
|
url = str(row["url"])
|
||||||
ok, exit_ip, latency_ms = await _probe_proxy(url)
|
was_disabled = not bool(row["enabled"])
|
||||||
mark_health(db, proxy_id, ok, exit_ip=exit_ip, latency_ms=latency_ms)
|
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
|
checked += 1
|
||||||
if ok:
|
if ok:
|
||||||
ok_count += 1
|
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:
|
else:
|
||||||
failed += 1
|
failed += 1
|
||||||
|
|
||||||
logger.info(
|
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,
|
reaped,
|
||||||
checked,
|
checked,
|
||||||
ok_count,
|
ok_count,
|
||||||
failed,
|
failed,
|
||||||
|
revived,
|
||||||
)
|
)
|
||||||
return {"reaped": reaped, "checked": checked, "ok": ok_count, "failed": failed}
|
return {
|
||||||
|
"reaped": reaped,
|
||||||
|
"checked": checked,
|
||||||
|
"ok": ok_count,
|
||||||
|
"failed": failed,
|
||||||
|
"revived": revived,
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
"""Offline-тесты пула прокси (#2162).
|
"""Offline-тесты пула прокси (#2162, #2600).
|
||||||
|
|
||||||
Покрытие БЕЗ live-сети/БД: stateful FakeSession эмулирует таблицу scrape_proxies и
|
Покрытие БЕЗ live-сети/БД: stateful FakeSession эмулирует таблицу scrape_proxies и
|
||||||
интерпретирует SQL по ключевым фрагментам, так что acquire/release/mark_health/
|
интерпретирует SQL по ключевым фрагментам, так что acquire/release/mark_health/
|
||||||
|
|
@ -8,11 +8,15 @@ reap_stale_leases проверяются по фактическому изме
|
||||||
- два acquire подряд → РАЗНЫЕ прокси (первый лизнут → выпал из выборки второго).
|
- два acquire подряд → РАЗНЫЕ прокси (первый лизнут → выпал из выборки второго).
|
||||||
- release освобождает (leased_by → NULL), прокси снова acquire-абелен.
|
- release освобождает (leased_by → NULL), прокси снова acquire-абелен.
|
||||||
- mark_health fail → инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
|
- 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, свежий не трогает.
|
- reap_stale_leases освобождает старый lease, свежий не трогает.
|
||||||
- affinity-фильтр: acquire('avito') не берёт cian-only прокси.
|
- affinity-фильтр: acquire('avito') не берёт cian-only прокси.
|
||||||
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
|
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
|
||||||
|
- acquire без своих/any свободных → берёт свободный чужой affinity (fallback, #2600 п.3).
|
||||||
- run_proxy_healthcheck: reap + проба каждого enabled + mark_health (проба замокана).
|
- run_proxy_healthcheck: reap + проба каждого enabled + mark_health (проба замокана).
|
||||||
|
- run_proxy_healthcheck: disabled-узлы — самовосстановление (#2600 п.1):
|
||||||
|
* успешная проба выключенного узла возвращает его в строй + revived++;
|
||||||
|
* недавно проверенный выключенный узел повторно не проверяется (не долбим провайдера).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -29,6 +33,7 @@ import pytest
|
||||||
from app.services import proxy_pool
|
from app.services import proxy_pool
|
||||||
from app.services.proxy_pool import (
|
from app.services.proxy_pool import (
|
||||||
DISABLE_THRESHOLD,
|
DISABLE_THRESHOLD,
|
||||||
|
DISABLED_RECHECK_MINUTES,
|
||||||
MAX_CONSECUTIVE_FAILS,
|
MAX_CONSECUTIVE_FAILS,
|
||||||
acquire,
|
acquire,
|
||||||
mark_health,
|
mark_health,
|
||||||
|
|
@ -71,17 +76,26 @@ class FakeSession:
|
||||||
sql = str(stmt)
|
sql = str(stmt)
|
||||||
p = params or {}
|
p = params or {}
|
||||||
|
|
||||||
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT
|
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
|
||||||
provider = p["provider"]
|
|
||||||
max_fails = p["max_fails"]
|
max_fails = p["max_fails"]
|
||||||
cands = [
|
if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any'
|
||||||
r
|
provider = p["provider"]
|
||||||
for r in self.rows
|
cands = [
|
||||||
if r["enabled"]
|
r
|
||||||
and r["consecutive_fails"] < max_fails
|
for r in self.rows
|
||||||
and r["provider_affinity"] in (provider, "any")
|
if r["enabled"]
|
||||||
and r["leased_by"] is None
|
and r["consecutive_fails"] < max_fails
|
||||||
]
|
and r["provider_affinity"] in (provider, "any")
|
||||||
|
and r["leased_by"] is None
|
||||||
|
]
|
||||||
|
else: # fallback: любая affinity (#2600 п.3)
|
||||||
|
cands = [
|
||||||
|
r
|
||||||
|
for r in self.rows
|
||||||
|
if r["enabled"]
|
||||||
|
and r["consecutive_fails"] < max_fails
|
||||||
|
and r["leased_by"] is None
|
||||||
|
]
|
||||||
# ORDER BY last_ok_at NULLS LAST, id
|
# ORDER BY last_ok_at NULLS LAST, id
|
||||||
cands.sort(
|
cands.sort(
|
||||||
key=lambda r: (
|
key=lambda r: (
|
||||||
|
|
@ -127,18 +141,28 @@ class FakeSession:
|
||||||
row["exit_ip"] = p["exit_ip"]
|
row["exit_ip"] = p["exit_ip"]
|
||||||
row["latency_ms"] = p["latency_ms"]
|
row["latency_ms"] = p["latency_ms"]
|
||||||
row["last_ok_at"] = datetime.now(UTC)
|
row["last_ok_at"] = datetime.now(UTC)
|
||||||
|
row["last_check_at"] = datetime.now(UTC)
|
||||||
|
row["enabled"] = True # реанимация выключенного узла (#2600 п.1)
|
||||||
return _FakeResult([])
|
return _FakeResult([])
|
||||||
|
|
||||||
if "consecutive_fails = consecutive_fails + 1" in sql: # mark_health fail
|
if "consecutive_fails = consecutive_fails + 1" in sql: # mark_health fail
|
||||||
row = self._by_id(p["id"])
|
row = self._by_id(p["id"])
|
||||||
if row is not None:
|
if row is not None:
|
||||||
row["consecutive_fails"] += 1
|
row["consecutive_fails"] += 1
|
||||||
|
row["last_check_at"] = datetime.now(UTC)
|
||||||
if row["consecutive_fails"] >= p["disable_threshold"]:
|
if row["consecutive_fails"] >= p["disable_threshold"]:
|
||||||
row["enabled"] = False
|
row["enabled"] = False
|
||||||
return _FakeResult([])
|
return _FakeResult([])
|
||||||
|
|
||||||
if "WHERE enabled" in sql and "ORDER BY id" in sql: # healthcheck SELECT
|
if "WHERE enabled" in sql and "ORDER BY id" in sql: # healthcheck SELECT (#2600 п.1)
|
||||||
rows = sorted((r for r in self.rows if r["enabled"]), key=lambda r: r["id"])
|
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])
|
return _FakeResult([dict(r) for r in rows])
|
||||||
|
|
||||||
raise AssertionError(f"unhandled SQL: {sql}")
|
raise AssertionError(f"unhandled SQL: {sql}")
|
||||||
|
|
@ -159,6 +183,7 @@ def _proxy(
|
||||||
leased_by: int | None = None,
|
leased_by: int | None = None,
|
||||||
leased_at: datetime | None = None,
|
leased_at: datetime | None = None,
|
||||||
last_ok_at: datetime | None = None,
|
last_ok_at: datetime | None = None,
|
||||||
|
last_check_at: datetime | None = None,
|
||||||
kind: str = "http",
|
kind: str = "http",
|
||||||
rotate_url: str | None = None,
|
rotate_url: str | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
|
|
@ -173,6 +198,7 @@ def _proxy(
|
||||||
"leased_by": leased_by,
|
"leased_by": leased_by,
|
||||||
"leased_at": leased_at,
|
"leased_at": leased_at,
|
||||||
"last_ok_at": last_ok_at,
|
"last_ok_at": last_ok_at,
|
||||||
|
"last_check_at": last_check_at,
|
||||||
"exit_ip": None,
|
"exit_ip": None,
|
||||||
"latency_ms": None,
|
"latency_ms": None,
|
||||||
}
|
}
|
||||||
|
|
@ -209,9 +235,17 @@ def test_acquire_empty_pool_returns_none() -> None:
|
||||||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
def test_acquire_affinity_filter_excludes_other_provider() -> None:
|
def test_acquire_affinity_filter_falls_back_instead_of_none() -> None:
|
||||||
|
"""До #2600 такой сетап возвращал None (голодный источник); теперь — fallback-выдача.
|
||||||
|
|
||||||
|
Поведение намеренно изменено п.3 issue #2600: чужой прокси лучше, чем никакого при
|
||||||
|
живом свободном узле. Дублирующее покрытие того же сценария —
|
||||||
|
test_acquire_falls_back_to_other_affinity_when_no_own_free.
|
||||||
|
"""
|
||||||
db = FakeSession([_proxy(1, affinity="cian")])
|
db = FakeSession([_proxy(1, affinity="cian")])
|
||||||
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert lease.id == 1
|
||||||
|
|
||||||
|
|
||||||
def test_acquire_skips_disabled() -> None:
|
def test_acquire_skips_disabled() -> None:
|
||||||
|
|
@ -231,6 +265,32 @@ def test_acquire_without_run_id_uses_marker() -> None:
|
||||||
assert db._by_id(1)["leased_by"] == proxy_pool.NON_RUN_LEASE_MARKER
|
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."""
|
||||||
|
db = FakeSession([_proxy(1, 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]
|
||||||
|
|
||||||
|
|
||||||
# ── release ──────────────────────────────────────────────────────────────────
|
# ── release ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -271,6 +331,15 @@ def test_mark_health_ok_resets_and_records() -> None:
|
||||||
assert row["last_ok_at"] is not 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 ────────────────────────────────────────────────────────
|
# ── reap_stale_leases ────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -295,27 +364,95 @@ def test_reap_frees_stale_lease_keeps_fresh() -> None:
|
||||||
async def test_healthcheck_probes_enabled_and_marks_health(
|
async def test_healthcheck_probes_enabled_and_marks_health(
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
recently_checked = datetime.now(UTC) - timedelta(minutes=5) # < DISABLED_RECHECK_MINUTES
|
||||||
db = FakeSession(
|
db = FakeSession(
|
||||||
[
|
[
|
||||||
_proxy(1, fails=2),
|
_proxy(1, fails=2),
|
||||||
_proxy(2, enabled=False), # disabled — не проверяется
|
# disabled, recheck ещё не наступил (недавно проверен) — не проверяется в этот прогон
|
||||||
|
_proxy(2, enabled=False, last_check_at=recently_checked),
|
||||||
_proxy(3, fails=0),
|
_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 «мёртв»
|
# прокси 1 «жив», прокси 3 «мёртв»
|
||||||
if "h1:" in url:
|
if "h1:" in url:
|
||||||
return True, "9.9.9.9", 42
|
return True, "9.9.9.9", 42, None
|
||||||
return False, None, None
|
return False, None, None, "other"
|
||||||
|
|
||||||
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||||||
|
|
||||||
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
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["ok"] == 1
|
||||||
assert counters["failed"] == 1
|
assert counters["failed"] == 1
|
||||||
|
assert counters["revived"] == 0
|
||||||
assert db._by_id(1)["consecutive_fails"] == 0 # ok → сброс
|
assert db._by_id(1)["consecutive_fails"] == 0 # ok → сброс
|
||||||
assert db._by_id(1)["exit_ip"] == "9.9.9.9"
|
assert db._by_id(1)["exit_ip"] == "9.9.9.9"
|
||||||
assert db._by_id(3)["consecutive_fails"] == 1 # fail → инкремент
|
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
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue