feat(proxy-pool): pool service (acquire/release/health) + healthcheck scheduler task (#2162)
Some checks failed
Deploy Trade-In / changes (push) Successful in 16s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Failing after 1m36s
Deploy Trade-In / build-backend (push) Has been skipped
Deploy Trade-In / deploy (push) Has been skipped
Deploy Trade-In / build-frontend (push) Successful in 33s
Some checks failed
Deploy Trade-In / changes (push) Successful in 16s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Failing after 1m36s
Deploy Trade-In / build-backend (push) Has been skipped
Deploy Trade-In / deploy (push) Has been skipped
Deploy Trade-In / build-frontend (push) Successful in 33s
This commit is contained in:
parent
62fa1f732a
commit
d7f4873a53
4 changed files with 763 additions and 0 deletions
321
tradein-mvp/backend/app/services/proxy_pool.py
Normal file
321
tradein-mvp/backend/app/services/proxy_pool.py
Normal file
|
|
@ -0,0 +1,321 @@
|
||||||
|
"""Пул прокси: подбор (lease), освобождение, health-трекинг (#2162).
|
||||||
|
|
||||||
|
АДДИТИВНО. Этот модуль реализует pick/lease/release/health-механику поверх таблицы
|
||||||
|
scrape_proxies (миграция 157, #2161). Ни один боевой скрейпер здесь НЕ подключается —
|
||||||
|
интеграция pick-из-пула вместо env-прокси это отдельные шаги P3/P4. Пока модуль
|
||||||
|
используется только health-checker'ом (run_proxy_healthcheck), который просто гоняет
|
||||||
|
ipify-пробу через каждый прокси и обновляет health-поля.
|
||||||
|
|
||||||
|
Семантика lease:
|
||||||
|
- scrape_proxies.leased_by IS NULL → прокси свободен.
|
||||||
|
- leased_by = <run_id> → занят run'ом.
|
||||||
|
- leased_by = NON_RUN_LEASE_MARKER → занят не-run вызовом (health-checker и т.п.),
|
||||||
|
когда run_id не применим.
|
||||||
|
acquire() берёт строку через FOR UPDATE SKIP LOCKED (конкурентные acquire не дерутся
|
||||||
|
за одну строку — второй параллельный вызов пропустит залоченную и возьмёт следующую).
|
||||||
|
|
||||||
|
Health:
|
||||||
|
- mark_health(ok=True) → consecutive_fails=0, last_ok_at/last_check_at, exit_ip, latency.
|
||||||
|
- mark_health(ok=False) → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
|
||||||
|
авто-disable (enabled=false), чтобы битый узел выпал из пула.
|
||||||
|
- acquire отфильтровывает enabled=false И consecutive_fails >= MAX_FAILS.
|
||||||
|
|
||||||
|
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import time
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"DISABLE_THRESHOLD",
|
||||||
|
"MAX_CONSECUTIVE_FAILS",
|
||||||
|
"NON_RUN_LEASE_MARKER",
|
||||||
|
"STALE_LEASE_MINUTES",
|
||||||
|
"ProxyLease",
|
||||||
|
"acquire",
|
||||||
|
"mark_health",
|
||||||
|
"reap_stale_leases",
|
||||||
|
"release",
|
||||||
|
"run_proxy_healthcheck",
|
||||||
|
]
|
||||||
|
|
||||||
|
# ── Пороги ───────────────────────────────────────────────────────────────────
|
||||||
|
# Прокси с >= MAX_CONSECUTIVE_FAILS подряд-фейлами не выдаётся acquire'ом (даже если
|
||||||
|
# ещё enabled) — «карантин» до первого успешного health-check'а (mark_health сбросит
|
||||||
|
# счётчик в 0). Мягче, чем disable: узел может ожить.
|
||||||
|
MAX_CONSECUTIVE_FAILS = 3
|
||||||
|
|
||||||
|
# При достижении этого порога подряд-фейлов прокси авто-disable (enabled=false) —
|
||||||
|
# оператор включит вручную после разбора. Строго >= MAX_CONSECUTIVE_FAILS.
|
||||||
|
DISABLE_THRESHOLD = 5
|
||||||
|
|
||||||
|
# Lease старше этого времени считается протухшим (упавший sweep не вызвал release) и
|
||||||
|
# освобождается reap_stale_leases — иначе прокси навсегда «занят» мёртвым run'ом.
|
||||||
|
STALE_LEASE_MINUTES = 30
|
||||||
|
|
||||||
|
# Маркер lease для не-run вызовов (leased_by NOT NULL = занят, но это не id из scrape_runs).
|
||||||
|
NON_RUN_LEASE_MARKER = -1
|
||||||
|
|
||||||
|
# URL для health-пробы: возвращает exit-IP JSON'ом. Тот же эндпоинт, что и admin
|
||||||
|
# /scraper/health (_probe_current_ip).
|
||||||
|
_HEALTH_PROBE_URL = "https://api.ipify.org"
|
||||||
|
_HEALTH_PROBE_TIMEOUT_S = 10.0
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class ProxyLease:
|
||||||
|
"""Арендованный прокси. url несёт схему (http:// / socks5://) — готов для httpx proxy=."""
|
||||||
|
|
||||||
|
id: int
|
||||||
|
url: str
|
||||||
|
kind: str
|
||||||
|
rotate_url: str | None
|
||||||
|
|
||||||
|
|
||||||
|
def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLease | None:
|
||||||
|
"""Взять свободный здоровый прокси под провайдера (avito/cian/yandex/generic/any).
|
||||||
|
|
||||||
|
SELECT ... FOR UPDATE SKIP LOCKED LIMIT 1 отбирает enabled-прокси с приемлемым
|
||||||
|
health (consecutive_fails < MAX_CONSECUTIVE_FAILS), affinity=provider ИЛИ 'any',
|
||||||
|
ещё не арендованный (leased_by IS NULL), предпочитая давно не проверенные
|
||||||
|
(last_ok_at NULLS LAST). Затем помечает строку leased_by=run_id (или
|
||||||
|
NON_RUN_LEASE_MARKER если run_id не задан) и коммитит.
|
||||||
|
|
||||||
|
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
||||||
|
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
||||||
|
|
||||||
|
Returns ProxyLease или None если свободных здоровых прокси нет.
|
||||||
|
"""
|
||||||
|
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
|
||||||
|
|
||||||
|
row = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT id, url, kind, rotate_url
|
||||||
|
FROM scrape_proxies
|
||||||
|
WHERE enabled
|
||||||
|
AND consecutive_fails < CAST(:max_fails AS integer)
|
||||||
|
AND provider_affinity IN (:provider, 'any')
|
||||||
|
AND leased_by IS NULL
|
||||||
|
ORDER BY last_ok_at NULLS LAST, id
|
||||||
|
FOR UPDATE SKIP LOCKED
|
||||||
|
LIMIT 1
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"max_fails": MAX_CONSECUTIVE_FAILS, "provider": provider},
|
||||||
|
)
|
||||||
|
.mappings()
|
||||||
|
.fetchone()
|
||||||
|
)
|
||||||
|
if row is None:
|
||||||
|
db.rollback() # снять FOR UPDATE-транзакцию (ничего не залочено, но чисто)
|
||||||
|
return None
|
||||||
|
|
||||||
|
proxy_id = int(row["id"])
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET leased_by = CAST(:run_id AS bigint), leased_at = now()
|
||||||
|
WHERE id = CAST(:id AS bigint)
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"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
|
||||||
|
)
|
||||||
|
return ProxyLease(
|
||||||
|
id=proxy_id,
|
||||||
|
url=str(row["url"]),
|
||||||
|
kind=str(row["kind"]),
|
||||||
|
rotate_url=row["rotate_url"],
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def release(db: Session, proxy_id: int) -> None:
|
||||||
|
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET leased_by = NULL, leased_at = NULL
|
||||||
|
WHERE id = CAST(:id AS bigint)
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"id": proxy_id},
|
||||||
|
)
|
||||||
|
db.commit()
|
||||||
|
logger.info("proxy_pool: released proxy id=%d", proxy_id)
|
||||||
|
|
||||||
|
|
||||||
|
def mark_health(
|
||||||
|
db: Session,
|
||||||
|
proxy_id: int,
|
||||||
|
ok: bool,
|
||||||
|
*,
|
||||||
|
exit_ip: str | None = None,
|
||||||
|
latency_ms: int | None = None,
|
||||||
|
) -> None:
|
||||||
|
"""Записать результат health-check'а прокси.
|
||||||
|
|
||||||
|
ok=True → consecutive_fails обнуляется, обновляются last_ok_at/last_check_at/
|
||||||
|
exit_ip/latency_ms.
|
||||||
|
ok=False → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
|
||||||
|
авто-disable (enabled=false). last_check_at обновляется в любом случае.
|
||||||
|
"""
|
||||||
|
if ok:
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET consecutive_fails = 0,
|
||||||
|
last_ok_at = now(),
|
||||||
|
last_check_at = now(),
|
||||||
|
exit_ip = CAST(:exit_ip AS text),
|
||||||
|
latency_ms = CAST(:latency_ms AS integer),
|
||||||
|
updated_at = now()
|
||||||
|
WHERE id = CAST(:id AS bigint)
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"exit_ip": exit_ip, "latency_ms": latency_ms, "id": proxy_id},
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
# consecutive_fails+1 >= порог → enabled=false (авто-вывод битого узла).
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET consecutive_fails = consecutive_fails + 1,
|
||||||
|
last_check_at = now(),
|
||||||
|
enabled = CASE
|
||||||
|
WHEN consecutive_fails + 1 >= CAST(:disable_threshold AS integer)
|
||||||
|
THEN false ELSE enabled
|
||||||
|
END,
|
||||||
|
updated_at = now()
|
||||||
|
WHERE id = CAST(:id AS bigint)
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"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)
|
||||||
|
|
||||||
|
|
||||||
|
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
|
||||||
|
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
|
||||||
|
|
||||||
|
Без этого прокси навсегда «занят» мёртвым run'ом и выпадает из пула. Returns число
|
||||||
|
освобождённых прокси.
|
||||||
|
"""
|
||||||
|
rows = db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET leased_by = NULL, leased_at = NULL
|
||||||
|
WHERE leased_by IS NOT NULL
|
||||||
|
AND leased_at < now() - make_interval(mins => CAST(:mins AS integer))
|
||||||
|
RETURNING id
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"mins": older_than_minutes},
|
||||||
|
).fetchall()
|
||||||
|
db.commit()
|
||||||
|
if rows:
|
||||||
|
logger.warning("proxy_pool: reaped %d stale lease(s): %s", len(rows), [r.id for r in rows])
|
||||||
|
return len(rows)
|
||||||
|
|
||||||
|
|
||||||
|
async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None]:
|
||||||
|
"""GET ipify через прокси (timeout _HEALTH_PROBE_TIMEOUT_S).
|
||||||
|
|
||||||
|
Returns (ok, exit_ip, latency_ms). ok=False + (None, None) при любой ошибке.
|
||||||
|
url несёт схему (http:// / socks5://) — httpx[socks] обрабатывает оба.
|
||||||
|
"""
|
||||||
|
started = time.monotonic()
|
||||||
|
try:
|
||||||
|
async with httpx.AsyncClient(proxy=url, timeout=_HEALTH_PROBE_TIMEOUT_S) as client:
|
||||||
|
resp = await client.get(_HEALTH_PROBE_URL, params={"format": "json"})
|
||||||
|
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
|
||||||
|
except Exception:
|
||||||
|
logger.warning("proxy_pool: health probe failed proxy=%s", _mask(url), exc_info=True)
|
||||||
|
return False, None, None
|
||||||
|
|
||||||
|
|
||||||
|
def _mask(url: str) -> str:
|
||||||
|
"""Скрыть пароль в proxy-url для логов (scheme://user:***@host)."""
|
||||||
|
if "@" not in url or "//" not in url:
|
||||||
|
return url
|
||||||
|
scheme, rest = url.split("//", 1)
|
||||||
|
creds, host = rest.split("@", 1)
|
||||||
|
if ":" in creds:
|
||||||
|
user, _pwd = creds.split(":", 1)
|
||||||
|
creds = f"{user}:***"
|
||||||
|
return f"{scheme}//{creds}@{host}"
|
||||||
|
|
||||||
|
|
||||||
|
async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
|
"""Периодический health-check всех enabled-прокси пула (#2162).
|
||||||
|
|
||||||
|
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем для каждого
|
||||||
|
enabled-прокси гоняет ipify-пробу через сам прокси и пишет результат через
|
||||||
|
mark_health (успех → сброс fails + свежий exit_ip/latency; фейл → инкремент,
|
||||||
|
авто-disable при DISABLE_THRESHOLD).
|
||||||
|
|
||||||
|
Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на
|
||||||
|
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
|
||||||
|
{reaped, checked, ok, failed}.
|
||||||
|
"""
|
||||||
|
reaped = reap_stale_leases(db)
|
||||||
|
|
||||||
|
proxies = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT id, url, kind
|
||||||
|
FROM scrape_proxies
|
||||||
|
WHERE enabled
|
||||||
|
ORDER BY id
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
)
|
||||||
|
.mappings()
|
||||||
|
.all()
|
||||||
|
)
|
||||||
|
|
||||||
|
checked = 0
|
||||||
|
ok_count = 0
|
||||||
|
failed = 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)
|
||||||
|
checked += 1
|
||||||
|
if ok:
|
||||||
|
ok_count += 1
|
||||||
|
else:
|
||||||
|
failed += 1
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d",
|
||||||
|
reaped,
|
||||||
|
checked,
|
||||||
|
ok_count,
|
||||||
|
failed,
|
||||||
|
)
|
||||||
|
return {"reaped": reaped, "checked": checked, "ok": ok_count, "failed": failed}
|
||||||
|
|
@ -1765,6 +1765,64 @@ async def trigger_house_dedup_merge_run(db: Session, schedule_row: dict[str, Any
|
||||||
return run_id
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
|
async def trigger_proxy_healthcheck_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
|
||||||
|
"""Создать scrape_runs + launch run_proxy_healthcheck в asyncio.create_task (#2162).
|
||||||
|
|
||||||
|
Периодический health-check пула scrape_proxies: reap_stale_leases + ipify-проба через
|
||||||
|
каждый enabled-прокси → mark_health (сброс fails/exit_ip/latency при успехе; инкремент
|
||||||
|
+ авто-disable при DISABLE_THRESHOLD на фейле). АДДИТИВНО — не трогает боевые скрейперы.
|
||||||
|
|
||||||
|
Зеркало trigger_yandex_newbuilding_sweep_run: у run_proxy_healthcheck нет tasks/*-обёртки,
|
||||||
|
владеющей lifecycle, поэтому финализацию scrape_runs (mark_done/mark_failed) делаем здесь.
|
||||||
|
|
||||||
|
Каденс: _claim_run ставит next_run_at через compute_next_run_at (суточная гранулярность),
|
||||||
|
но health-check хочется чаще. Поэтому СРАЗУ после claim переопределяем next_run_at на
|
||||||
|
now() + interval_minutes (default 30) — source-специфичный sub-hourly re-schedule, не
|
||||||
|
трогая shared compute_next_run_at. Если run ещё идёт на следующем тике — has_running_run
|
||||||
|
в _claim_run вернёт None (skip) и next_run_at не сбросится, т.е. авто-throttle по факту.
|
||||||
|
|
||||||
|
Returns run_id или None (skip — already running).
|
||||||
|
"""
|
||||||
|
run_id = _claim_run(db, schedule_row)
|
||||||
|
if run_id is None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
params = schedule_row.get("default_params") or {}
|
||||||
|
interval_minutes = int(params.get("interval_minutes", 30))
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
UPDATE scrape_schedules
|
||||||
|
SET next_run_at = now() + make_interval(mins => CAST(:mins AS integer)),
|
||||||
|
updated_at = now()
|
||||||
|
WHERE source = 'proxy_healthcheck'
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"mins": interval_minutes},
|
||||||
|
)
|
||||||
|
db.commit()
|
||||||
|
|
||||||
|
async def _run() -> None:
|
||||||
|
run_db = SessionLocal()
|
||||||
|
try:
|
||||||
|
from app.services.proxy_pool import run_proxy_healthcheck
|
||||||
|
|
||||||
|
counters = await run_proxy_healthcheck(run_db)
|
||||||
|
runs_mod.mark_done(run_db, run_id, counters)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.exception("scheduler: run_proxy_healthcheck crashed run_id=%d", run_id)
|
||||||
|
try:
|
||||||
|
runs_mod.mark_failed(run_db, run_id, str(exc)[:1000], {})
|
||||||
|
except Exception:
|
||||||
|
logger.exception("scheduler: mark_failed crashed run_id=%d", run_id)
|
||||||
|
finally:
|
||||||
|
run_db.close()
|
||||||
|
|
||||||
|
_spawn_tracked(_run())
|
||||||
|
logger.info("scheduler: triggered proxy_healthcheck run_id=%d", run_id)
|
||||||
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
def get_due_schedules(db: Session) -> list[dict[str, Any]]:
|
def get_due_schedules(db: Session) -> list[dict[str, Any]]:
|
||||||
"""SELECT scrape_schedules WHERE enabled AND (next_run_at IS NULL OR next_run_at <= NOW())."""
|
"""SELECT scrape_schedules WHERE enabled AND (next_run_at IS NULL OR next_run_at <= NOW())."""
|
||||||
rows = (
|
rows = (
|
||||||
|
|
@ -1855,6 +1913,8 @@ async def scheduler_loop() -> None:
|
||||||
await trigger_house_imv_backfill_run(db, sch)
|
await trigger_house_imv_backfill_run(db, sch)
|
||||||
elif source == "house_dedup_merge":
|
elif source == "house_dedup_merge":
|
||||||
await trigger_house_dedup_merge_run(db, sch)
|
await trigger_house_dedup_merge_run(db, sch)
|
||||||
|
elif source == "proxy_healthcheck":
|
||||||
|
await trigger_proxy_healthcheck_run(db, sch)
|
||||||
else:
|
else:
|
||||||
logger.warning("scheduler: unknown source=%s, skip", source)
|
logger.warning("scheduler: unknown source=%s, skip", source)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,61 @@
|
||||||
|
-- 158_seed_proxy_healthcheck_schedule.sql
|
||||||
|
-- Scheduler seed for the proxy-pool health-checker (#2162). АДДИТИВНО.
|
||||||
|
--
|
||||||
|
-- WHAT (source='proxy_healthcheck'):
|
||||||
|
-- trigger_proxy_healthcheck_run (scheduler.py) → run_proxy_healthcheck
|
||||||
|
-- (services/proxy_pool.py):
|
||||||
|
-- 1. reap_stale_leases — освобождает lease'ы старше STALE_LEASE_MINUTES (упавший
|
||||||
|
-- sweep не вызвал release), иначе прокси навсегда «занят» мёртвым run'ом.
|
||||||
|
-- 2. для каждого enabled-прокси гоняет ipify-пробу ЧЕРЕЗ сам прокси и пишет
|
||||||
|
-- результат через mark_health: успех → consecutive_fails=0 + свежий exit_ip/
|
||||||
|
-- latency; фейл → consecutive_fails += 1, авто-disable при DISABLE_THRESHOLD.
|
||||||
|
-- Чистый health-refresh пула — НЕ трогает боевые скрейперы (pick/lease/rotate из
|
||||||
|
-- пула вместо env — отдельные шаги P3/P4).
|
||||||
|
--
|
||||||
|
-- КАДЕНС (sub-hourly, не суточный):
|
||||||
|
-- Штатный scheduler-каденс суточный (compute_next_run_at пикает next_run_at в окне
|
||||||
|
-- [window_start_hour, window_end_hour) на interval_days суток вперёд). Health-check
|
||||||
|
-- хочется чаще, поэтому trigger_proxy_healthcheck_run СРАЗУ после claim переопределяет
|
||||||
|
-- next_run_at на now() + interval_minutes (default 30) — source-специфичный re-schedule,
|
||||||
|
-- не трогая shared compute_next_run_at. window_start/end_hour для этого source не важны
|
||||||
|
-- (next_run_at выставляется напрямую), но CHECK 0-23 требует валидных значений → 0..23
|
||||||
|
-- = «в любой час». next_run_at при сиде = now() → первый health-check вскоре после
|
||||||
|
-- деплоя (в пределах SCHEDULER_TICK_SEC).
|
||||||
|
--
|
||||||
|
-- default_params:
|
||||||
|
-- interval_minutes -- 30: период между health-check'ами (trigger читает это).
|
||||||
|
--
|
||||||
|
-- *** enabled=true — БЕЗОПАСНО ***
|
||||||
|
-- Пул scrape_proxies по умолчанию ПУСТ (прокси грузятся оператором через admin-ручку
|
||||||
|
-- POST /api/v1/admin/proxies/bulk). Пустой пул → run_proxy_healthcheck просто ничего не
|
||||||
|
-- проверяет (checked=0), deploy нейтрален. Как только оператор зальёт прокси —
|
||||||
|
-- health-check начнёт их пинговать. Внешние вызовы идут ТОЛЬКО на ipify через сами
|
||||||
|
-- прокси (не на боевые площадки), анти-бот-рисков нет.
|
||||||
|
--
|
||||||
|
-- DEPENDENCIES: 052_scrape_schedules.sql (table + UNIQUE(source)),
|
||||||
|
-- 157_scrape_proxies.sql (пул, который проверяем).
|
||||||
|
-- Idempotent: ON CONFLICT (source) DO NOTHING.
|
||||||
|
-- Runner applies migrations WITHOUT --single-transaction, so wrap in explicit BEGIN/COMMIT.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
INSERT INTO scrape_schedules (
|
||||||
|
source,
|
||||||
|
enabled,
|
||||||
|
window_start_hour,
|
||||||
|
window_end_hour,
|
||||||
|
next_run_at,
|
||||||
|
default_params
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
'proxy_healthcheck',
|
||||||
|
true,
|
||||||
|
0,
|
||||||
|
23,
|
||||||
|
now(),
|
||||||
|
'{"interval_minutes": 30}'::jsonb
|
||||||
|
)
|
||||||
|
ON CONFLICT (source) DO NOTHING;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
321
tradein-mvp/backend/tests/services/test_proxy_pool.py
Normal file
321
tradein-mvp/backend/tests/services/test_proxy_pool.py
Normal file
|
|
@ -0,0 +1,321 @@
|
||||||
|
"""Offline-тесты пула прокси (#2162).
|
||||||
|
|
||||||
|
Покрытие БЕЗ live-сети/БД: stateful FakeSession эмулирует таблицу scrape_proxies и
|
||||||
|
интерпретирует SQL по ключевым фрагментам, так что acquire/release/mark_health/
|
||||||
|
reap_stale_leases проверяются по фактическому изменению состояния строк.
|
||||||
|
|
||||||
|
- acquire: возвращает свободный+здоровый прокси нужного affinity; ставит lease.
|
||||||
|
- два acquire подряд → РАЗНЫЕ прокси (первый лизнут → выпал из выборки второго).
|
||||||
|
- release освобождает (leased_by → NULL), прокси снова acquire-абелен.
|
||||||
|
- mark_health fail → инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
|
||||||
|
- mark_health ok → сброс fails + exit_ip/latency.
|
||||||
|
- reap_stale_leases освобождает старый lease, свежий не трогает.
|
||||||
|
- affinity-фильтр: acquire('avito') не берёт cian-only прокси.
|
||||||
|
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
|
||||||
|
- run_proxy_healthcheck: reap + проба каждого enabled + mark_health (проба замокана).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
from datetime import UTC, datetime, timedelta
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from app.services import proxy_pool
|
||||||
|
from app.services.proxy_pool import (
|
||||||
|
DISABLE_THRESHOLD,
|
||||||
|
MAX_CONSECUTIVE_FAILS,
|
||||||
|
acquire,
|
||||||
|
mark_health,
|
||||||
|
reap_stale_leases,
|
||||||
|
release,
|
||||||
|
)
|
||||||
|
|
||||||
|
# ── stateful fake session ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeResult:
|
||||||
|
def __init__(self, rows: list[dict[str, Any]]):
|
||||||
|
self._rows = rows
|
||||||
|
|
||||||
|
def mappings(self) -> _FakeResult:
|
||||||
|
return self
|
||||||
|
|
||||||
|
def fetchone(self) -> dict[str, Any] | None:
|
||||||
|
return self._rows[0] if self._rows else None
|
||||||
|
|
||||||
|
def all(self) -> list[dict[str, Any]]:
|
||||||
|
return list(self._rows)
|
||||||
|
|
||||||
|
def fetchall(self) -> list[Any]:
|
||||||
|
# для RETURNING id: код делает r.id → нужен attribute-access
|
||||||
|
return [type("Row", (), r)() for r in self._rows]
|
||||||
|
|
||||||
|
|
||||||
|
class FakeSession:
|
||||||
|
"""Эмуляция Session поверх in-memory списка scrape_proxies-строк."""
|
||||||
|
|
||||||
|
def __init__(self, rows: list[dict[str, Any]]):
|
||||||
|
self.rows = rows
|
||||||
|
|
||||||
|
# helpers
|
||||||
|
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
||||||
|
return next((r for r in self.rows if r["id"] == pid), None)
|
||||||
|
|
||||||
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||||
|
sql = str(stmt)
|
||||||
|
p = params or {}
|
||||||
|
|
||||||
|
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT
|
||||||
|
provider = p["provider"]
|
||||||
|
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
|
||||||
|
]
|
||||||
|
# ORDER BY last_ok_at NULLS LAST, id
|
||||||
|
cands.sort(
|
||||||
|
key=lambda r: (
|
||||||
|
r["last_ok_at"] is None,
|
||||||
|
r["last_ok_at"] or datetime.min.replace(tzinfo=UTC),
|
||||||
|
r["id"],
|
||||||
|
)
|
||||||
|
)
|
||||||
|
return _FakeResult(cands[:1])
|
||||||
|
|
||||||
|
if "SET leased_by = CAST(:run_id" in sql: # acquire lease UPDATE
|
||||||
|
row = self._by_id(p["id"])
|
||||||
|
if row is not None:
|
||||||
|
row["leased_by"] = p["run_id"]
|
||||||
|
row["leased_at"] = datetime.now(UTC)
|
||||||
|
return _FakeResult([])
|
||||||
|
|
||||||
|
if "make_interval(mins =>" in sql and "leased_by IS NOT NULL" in sql: # reap
|
||||||
|
cutoff = datetime.now(UTC) - timedelta(minutes=p["mins"])
|
||||||
|
freed: list[dict[str, Any]] = []
|
||||||
|
for r in self.rows:
|
||||||
|
if (
|
||||||
|
r["leased_by"] is not None
|
||||||
|
and r["leased_at"] is not None
|
||||||
|
and r["leased_at"] < cutoff
|
||||||
|
):
|
||||||
|
r["leased_by"] = None
|
||||||
|
r["leased_at"] = None
|
||||||
|
freed.append({"id": r["id"]})
|
||||||
|
return _FakeResult(freed)
|
||||||
|
|
||||||
|
if "SET leased_by = NULL" in sql: # release (WHERE id)
|
||||||
|
row = self._by_id(p["id"])
|
||||||
|
if row is not None:
|
||||||
|
row["leased_by"] = None
|
||||||
|
row["leased_at"] = None
|
||||||
|
return _FakeResult([])
|
||||||
|
|
||||||
|
if "SET consecutive_fails = 0" in sql: # mark_health ok
|
||||||
|
row = self._by_id(p["id"])
|
||||||
|
if row is not None:
|
||||||
|
row["consecutive_fails"] = 0
|
||||||
|
row["exit_ip"] = p["exit_ip"]
|
||||||
|
row["latency_ms"] = p["latency_ms"]
|
||||||
|
row["last_ok_at"] = datetime.now(UTC)
|
||||||
|
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
|
||||||
|
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"])
|
||||||
|
return _FakeResult([dict(r) for r in rows])
|
||||||
|
|
||||||
|
raise AssertionError(f"unhandled SQL: {sql}")
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
pass
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def _proxy(
|
||||||
|
pid: int,
|
||||||
|
*,
|
||||||
|
affinity: str = "any",
|
||||||
|
enabled: bool = True,
|
||||||
|
fails: int = 0,
|
||||||
|
leased_by: int | None = None,
|
||||||
|
leased_at: datetime | None = None,
|
||||||
|
last_ok_at: datetime | None = None,
|
||||||
|
kind: str = "http",
|
||||||
|
rotate_url: str | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"id": pid,
|
||||||
|
"url": f"http://u:p@h{pid}:8080",
|
||||||
|
"kind": kind,
|
||||||
|
"rotate_url": rotate_url,
|
||||||
|
"provider_affinity": affinity,
|
||||||
|
"enabled": enabled,
|
||||||
|
"consecutive_fails": fails,
|
||||||
|
"leased_by": leased_by,
|
||||||
|
"leased_at": leased_at,
|
||||||
|
"last_ok_at": last_ok_at,
|
||||||
|
"exit_ip": None,
|
||||||
|
"latency_ms": None,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
# ── acquire ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_returns_free_healthy_of_affinity() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="cian")])
|
||||||
|
lease = acquire(db, "avito", run_id=100) # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert lease.id == 1
|
||||||
|
assert lease.url == "http://u:p@h1:8080"
|
||||||
|
assert db._by_id(1)["leased_by"] == 100 # lease проставлен
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_includes_any_affinity() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="any")])
|
||||||
|
lease = acquire(db, "avito", run_id=5) # type: ignore[arg-type]
|
||||||
|
assert lease is not None and lease.id == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_two_acquire_return_distinct_proxies() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="avito")])
|
||||||
|
first = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
second = acquire(db, "avito", run_id=2) # type: ignore[arg-type]
|
||||||
|
assert first is not None and second is not None
|
||||||
|
assert first.id != second.id # первый лизнут → второй берёт другой
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_empty_pool_returns_none() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito", leased_by=99)]) # единственный занят
|
||||||
|
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_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]
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_skips_unhealthy() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito", fails=MAX_CONSECUTIVE_FAILS)])
|
||||||
|
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_without_run_id_uses_marker() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito")])
|
||||||
|
lease = acquire(db, "avito") # run_id=None # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert db._by_id(1)["leased_by"] == proxy_pool.NON_RUN_LEASE_MARKER
|
||||||
|
|
||||||
|
|
||||||
|
# ── release ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_release_frees_proxy() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito")])
|
||||||
|
lease = acquire(db, "avito", run_id=7) # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert acquire(db, "avito", run_id=8) is None # type: ignore[arg-type] # занят
|
||||||
|
release(db, 1) # type: ignore[arg-type]
|
||||||
|
assert db._by_id(1)["leased_by"] is None
|
||||||
|
assert acquire(db, "avito", run_id=9) is not None # type: ignore[arg-type] # снова свободен
|
||||||
|
|
||||||
|
|
||||||
|
# ── mark_health ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_health_fail_increments() -> None:
|
||||||
|
db = FakeSession([_proxy(1, fails=0)])
|
||||||
|
mark_health(db, 1, ok=False) # type: ignore[arg-type]
|
||||||
|
assert db._by_id(1)["consecutive_fails"] == 1
|
||||||
|
assert db._by_id(1)["enabled"] is True # ещё не порог
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_health_fail_disables_at_threshold() -> None:
|
||||||
|
db = FakeSession([_proxy(1, fails=DISABLE_THRESHOLD - 1)])
|
||||||
|
mark_health(db, 1, ok=False) # type: ignore[arg-type]
|
||||||
|
assert db._by_id(1)["consecutive_fails"] == DISABLE_THRESHOLD
|
||||||
|
assert db._by_id(1)["enabled"] is False # авто-disable
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_health_ok_resets_and_records() -> None:
|
||||||
|
db = FakeSession([_proxy(1, fails=4)])
|
||||||
|
mark_health(db, 1, ok=True, exit_ip="1.2.3.4", latency_ms=88) # type: ignore[arg-type]
|
||||||
|
row = db._by_id(1)
|
||||||
|
assert row["consecutive_fails"] == 0
|
||||||
|
assert row["exit_ip"] == "1.2.3.4"
|
||||||
|
assert row["latency_ms"] == 88
|
||||||
|
assert row["last_ok_at"] is not None
|
||||||
|
|
||||||
|
|
||||||
|
# ── reap_stale_leases ────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_reap_frees_stale_lease_keeps_fresh() -> None:
|
||||||
|
old = datetime.now(UTC) - timedelta(minutes=120)
|
||||||
|
fresh = datetime.now(UTC) - timedelta(minutes=1)
|
||||||
|
db = FakeSession(
|
||||||
|
[
|
||||||
|
_proxy(1, leased_by=50, leased_at=old),
|
||||||
|
_proxy(2, leased_by=51, leased_at=fresh),
|
||||||
|
]
|
||||||
|
)
|
||||||
|
freed = reap_stale_leases(db, older_than_minutes=30) # type: ignore[arg-type]
|
||||||
|
assert freed == 1
|
||||||
|
assert db._by_id(1)["leased_by"] is None # протухший освобождён
|
||||||
|
assert db._by_id(2)["leased_by"] == 51 # свежий не тронут
|
||||||
|
|
||||||
|
|
||||||
|
# ── run_proxy_healthcheck ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_healthcheck_probes_enabled_and_marks_health(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
db = FakeSession(
|
||||||
|
[
|
||||||
|
_proxy(1, fails=2),
|
||||||
|
_proxy(2, enabled=False), # disabled — не проверяется
|
||||||
|
_proxy(3, fails=0),
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None]:
|
||||||
|
# прокси 1 «жив», прокси 3 «мёртв»
|
||||||
|
if "h1:" in url:
|
||||||
|
return True, "9.9.9.9", 42
|
||||||
|
return False, None, None
|
||||||
|
|
||||||
|
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["ok"] == 1
|
||||||
|
assert counters["failed"] == 1
|
||||||
|
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 → инкремент
|
||||||
Loading…
Add table
Reference in a new issue