Merge pull request 'fix(tradein/proxy): самовосстановление пула и запасной прокси чужой affinity (#2600)' (#2609) from fix/tradein-proxy-pool-self-healing into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 16s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m40s
Deploy Trade-In / build-backend (push) Successful in 1m8s
Deploy Trade-In / deploy (push) Successful in 1m17s
All checks were successful
Deploy Trade-In / changes (push) Successful in 16s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m40s
Deploy Trade-In / build-backend (push) Successful in 1m8s
Deploy Trade-In / deploy (push) Successful in 1m17s
This commit is contained in:
commit
5b2d6807bc
2 changed files with 374 additions and 48 deletions
|
|
@ -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,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue