Merge remote-tracking branch 'origin/main' into fix/2603-geocode-city-hint-tails
All checks were successful
CI / changes (pull_request) Successful in 7s
CI Trade-In / changes (pull_request) Successful in 8s
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 2m41s
All checks were successful
CI / changes (pull_request) Successful in 7s
CI Trade-In / changes (pull_request) Successful in 8s
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 2m41s
# Conflicts: # tradein-mvp/backend/app/api/v1/admin.py
This commit is contained in:
commit
17fcf746f7
10 changed files with 1060 additions and 150 deletions
|
|
@ -73,6 +73,7 @@ from app.services import domclick_session as domclick_session_svc
|
||||||
from app.services import proxy_rotation as proxy_rotation_svc
|
from app.services import proxy_rotation as proxy_rotation_svc
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
from app.services.geocoder import geocode, known_city_hint
|
from app.services.geocoder import geocode, known_city_hint
|
||||||
|
from app.services.proxy_pool import clear_source_bans
|
||||||
from app.services.scheduler import has_running_run
|
from app.services.scheduler import has_running_run
|
||||||
from app.services.scraper_adapters import (
|
from app.services.scraper_adapters import (
|
||||||
RealEnrichmentJobs,
|
RealEnrichmentJobs,
|
||||||
|
|
@ -2645,6 +2646,52 @@ class ProxyBulkResponse(BaseModel):
|
||||||
updated: int
|
updated: int
|
||||||
|
|
||||||
|
|
||||||
|
class ProxySourceBan(BaseModel):
|
||||||
|
"""Активный бан узла КОНКРЕТНОЙ площадкой (#2600 п.2, scrape_proxy_source_bans)."""
|
||||||
|
|
||||||
|
source: str
|
||||||
|
banned_until: str
|
||||||
|
ban_count: int
|
||||||
|
|
||||||
|
|
||||||
|
def _fetch_source_bans(db: Session, proxy_ids: list[int]) -> dict[int, list[ProxySourceBan]]:
|
||||||
|
"""Активные (banned_until > now()) баны по источникам для указанных узлов.
|
||||||
|
|
||||||
|
Без этого оператор видит `enabled=true` и не понимает, почему узел не выдаётся
|
||||||
|
конкретному источнику (#2600 п.2 — бан теперь по паре «узел × источник», а не
|
||||||
|
глобальное выключение). Истёкшие строки не показываем: они ни на что не влияют,
|
||||||
|
живут ещё SOURCE_BAN_PURGE_DAYS только как память об эскалации.
|
||||||
|
"""
|
||||||
|
if not proxy_ids:
|
||||||
|
return {}
|
||||||
|
rows = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT proxy_id, source, banned_until, ban_count
|
||||||
|
FROM scrape_proxy_source_bans
|
||||||
|
WHERE banned_until > now()
|
||||||
|
AND proxy_id = ANY(CAST(:ids AS bigint[]))
|
||||||
|
ORDER BY proxy_id, source
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"ids": proxy_ids},
|
||||||
|
)
|
||||||
|
.mappings()
|
||||||
|
.all()
|
||||||
|
)
|
||||||
|
bans: dict[int, list[ProxySourceBan]] = {}
|
||||||
|
for r in rows:
|
||||||
|
bans.setdefault(int(r["proxy_id"]), []).append(
|
||||||
|
ProxySourceBan(
|
||||||
|
source=r["source"],
|
||||||
|
banned_until=r["banned_until"].isoformat(),
|
||||||
|
ban_count=r["ban_count"],
|
||||||
|
)
|
||||||
|
)
|
||||||
|
return bans
|
||||||
|
|
||||||
|
|
||||||
class ProxyRow(BaseModel):
|
class ProxyRow(BaseModel):
|
||||||
id: int
|
id: int
|
||||||
label: str | None
|
label: str | None
|
||||||
|
|
@ -2666,6 +2713,8 @@ class ProxyRow(BaseModel):
|
||||||
expires_at: str | None
|
expires_at: str | None
|
||||||
created_at: str | None
|
created_at: str | None
|
||||||
updated_at: str | None
|
updated_at: str | None
|
||||||
|
# #2600 п.2: активные баны площадками. Пустой список = узел выдаётся всем источникам.
|
||||||
|
source_bans: list[ProxySourceBan] = Field(default_factory=list)
|
||||||
|
|
||||||
|
|
||||||
@router.post("/proxies/bulk", response_model=ProxyBulkResponse)
|
@router.post("/proxies/bulk", response_model=ProxyBulkResponse)
|
||||||
|
|
@ -2751,6 +2800,9 @@ def list_proxies(
|
||||||
"""Список прокси со статусами. Пароли в url/rotate_url маскируются.
|
"""Список прокси со статусами. Пароли в url/rotate_url маскируются.
|
||||||
|
|
||||||
Фильтры: provider (=provider_affinity), enabled. Без фильтров — все.
|
Фильтры: provider (=provider_affinity), enabled. Без фильтров — все.
|
||||||
|
|
||||||
|
source_bans — активные баны узла площадками (#2600 п.2): узел может быть
|
||||||
|
enabled=true и при этом не выдаваться конкретному источнику.
|
||||||
"""
|
"""
|
||||||
clauses: list[str] = []
|
clauses: list[str] = []
|
||||||
params: dict[str, Any] = {}
|
params: dict[str, Any] = {}
|
||||||
|
|
@ -2785,6 +2837,8 @@ def list_proxies(
|
||||||
def _iso(v: Any) -> str | None:
|
def _iso(v: Any) -> str | None:
|
||||||
return v.isoformat() if v is not None else None
|
return v.isoformat() if v is not None else None
|
||||||
|
|
||||||
|
bans = _fetch_source_bans(db, [int(r["id"]) for r in rows])
|
||||||
|
|
||||||
return [
|
return [
|
||||||
ProxyRow(
|
ProxyRow(
|
||||||
id=r["id"],
|
id=r["id"],
|
||||||
|
|
@ -2807,6 +2861,7 @@ def list_proxies(
|
||||||
expires_at=_iso(r["expires_at"]),
|
expires_at=_iso(r["expires_at"]),
|
||||||
created_at=_iso(r["created_at"]),
|
created_at=_iso(r["created_at"]),
|
||||||
updated_at=_iso(r["updated_at"]),
|
updated_at=_iso(r["updated_at"]),
|
||||||
|
source_bans=bans.get(int(r["id"]), []),
|
||||||
)
|
)
|
||||||
for r in rows
|
for r in rows
|
||||||
]
|
]
|
||||||
|
|
@ -2875,6 +2930,11 @@ def patch_proxy(
|
||||||
if row is None:
|
if row is None:
|
||||||
raise HTTPException(status_code=404, detail=f"proxy id={proxy_id} not found")
|
raise HTTPException(status_code=404, detail=f"proxy id={proxy_id} not found")
|
||||||
db.commit()
|
db.commit()
|
||||||
|
if payload.enabled:
|
||||||
|
# Ручное включение = чистый лист, как и обнуление disabled_reason выше (#2610).
|
||||||
|
# Иначе узел вернулся бы enabled=true, но по-прежнему невыдаваемым источникам с
|
||||||
|
# активным баном — и оператор не имел бы способа снять ложный бан (#2600 п.2).
|
||||||
|
clear_source_bans(db, proxy_id, reason="manual enable via admin API")
|
||||||
if not payload.enabled:
|
if not payload.enabled:
|
||||||
logger.info(
|
logger.info(
|
||||||
"proxy_pool: proxy id=%d manually disabled via admin API (reason=%r) — "
|
"proxy_pool: proxy id=%d manually disabled via admin API (reason=%r) — "
|
||||||
|
|
@ -2907,6 +2967,7 @@ def patch_proxy(
|
||||||
expires_at=_iso(row["expires_at"]),
|
expires_at=_iso(row["expires_at"]),
|
||||||
created_at=_iso(row["created_at"]),
|
created_at=_iso(row["created_at"]),
|
||||||
updated_at=_iso(row["updated_at"]),
|
updated_at=_iso(row["updated_at"]),
|
||||||
|
source_bans=_fetch_source_bans(db, [int(row["id"])]).get(int(row["id"]), []),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -34,6 +34,28 @@ Self-healing (#2600):
|
||||||
enabled-узел выделенной affinity (пример — domclick, один узел на всё, см. acquire
|
enabled-узел выделенной affinity (пример — domclick, один узел на всё, см. acquire
|
||||||
docstring) — иначе чинили бы один источник ценой полной поломки другого.
|
docstring) — иначе чинили бы один источник ценой полной поломки другого.
|
||||||
|
|
||||||
|
Бан по паре «узел × источник» (#2600 п.2, таблица scrape_proxy_source_bans, миграция 210):
|
||||||
|
- Авито банит IP, Яндекс через тот же IP ходит чисто. Поэтому распознанный бан
|
||||||
|
площадкой (`mark_banned`) НЕ выключает узел глобально (так делал #2600 п.1), а
|
||||||
|
пишет строку (proxy_id, source, banned_until) — `acquire(source)` перестаёт
|
||||||
|
выдавать узел ЭТОМУ источнику, для остальных узел остаётся первосортным.
|
||||||
|
- Отличие от `enabled=false`: глобальное выключение — это либо решение оператора
|
||||||
|
(disabled_reason НЕ NULL, #2610), либо авто-disable по серии ТРАНСПОРТНЫХ сбоев
|
||||||
|
(mark_health, DISABLE_THRESHOLD). Бан площадкой — свойство ПАРЫ, а не узла, и
|
||||||
|
снимается сам по времени, без ручного PATCH и без ipify-пробы (ipify площадку не
|
||||||
|
эмулирует, бана не видит — ровно тот баг, из-за которого п.1 требовал ручного
|
||||||
|
вмешательства).
|
||||||
|
- Срок эскалирует на повторных банах той же пары: SOURCE_BAN_BASE_HOURS *
|
||||||
|
2^(ban_count-1), но не больше SOURCE_BAN_MAX_HOURS. Истёкшие строки сносятся
|
||||||
|
purge'ем в run_proxy_healthcheck только через SOURCE_BAN_PURGE_DAYS — это же и
|
||||||
|
механизм сброса ban_count (см. комментарий у purge, НЕ «оптимизировать»).
|
||||||
|
- Защита последнего узла сохранена, но теперь ПО ИСТОЧНИКУ: если после записи бана
|
||||||
|
у acquire(source) не останется ни одного кандидата — бан не пишется, только
|
||||||
|
WARNING (пул надо пополнять, #2638).
|
||||||
|
- Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс
|
||||||
|
автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился
|
||||||
|
IP, поэтому смена адреса делает строку недействительной.
|
||||||
|
|
||||||
Ручное выключение vs авто-выключение (#2610):
|
Ручное выключение vs авто-выключение (#2610):
|
||||||
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
|
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
|
||||||
enabled=false: пул выключил сам после серии сбоев (disabled_reason IS NULL) —
|
enabled=false: пул выключил сам после серии сбоев (disabled_reason IS NULL) —
|
||||||
|
|
@ -75,9 +97,13 @@ __all__ = [
|
||||||
"DISABLE_THRESHOLD",
|
"DISABLE_THRESHOLD",
|
||||||
"MAX_CONSECUTIVE_FAILS",
|
"MAX_CONSECUTIVE_FAILS",
|
||||||
"NON_RUN_LEASE_MARKER",
|
"NON_RUN_LEASE_MARKER",
|
||||||
|
"SOURCE_BAN_BASE_HOURS",
|
||||||
|
"SOURCE_BAN_MAX_HOURS",
|
||||||
|
"SOURCE_BAN_PURGE_DAYS",
|
||||||
"STALE_LEASE_MINUTES",
|
"STALE_LEASE_MINUTES",
|
||||||
"ProxyLease",
|
"ProxyLease",
|
||||||
"acquire",
|
"acquire",
|
||||||
|
"clear_source_bans",
|
||||||
"mark_banned",
|
"mark_banned",
|
||||||
"mark_health",
|
"mark_health",
|
||||||
"reap_stale_leases",
|
"reap_stale_leases",
|
||||||
|
|
@ -109,6 +135,21 @@ 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
|
||||||
|
|
||||||
|
# ── Бан по паре «узел × источник» (#2600 п.2) ────────────────────────────────
|
||||||
|
# Срок ПЕРВОГО бана пары (proxy_id, source). 6 часов — эмпирический компромисс:
|
||||||
|
# площадки снимают IP-баны обычно за часы, а не минуты (короче — вернём узел под тот
|
||||||
|
# же бан и потратим прогон впустую), но и не сутки (узел дефицитный, #2638).
|
||||||
|
SOURCE_BAN_BASE_HOURS = 6
|
||||||
|
|
||||||
|
# Потолок эскалации: SOURCE_BAN_BASE_HOURS * 2^(ban_count-1) обрезается этим значением
|
||||||
|
# (6 → 12 → 24 → 48 → 72 → 72 …). Дольше 3 суток держать бесполезно: либо площадка
|
||||||
|
# сняла бан, либо узел мёртв насовсем и его должен вычистить оператор.
|
||||||
|
SOURCE_BAN_MAX_HOURS = 72
|
||||||
|
|
||||||
|
# Через столько суток ПОСЛЕ истечения бана строка сносится purge'ем (см.
|
||||||
|
# run_proxy_healthcheck). Это же и сброс ban_count — см. комментарий там.
|
||||||
|
SOURCE_BAN_PURGE_DAYS = 7
|
||||||
|
|
||||||
# URL для health-пробы: возвращает exit-IP JSON'ом. Тот же эндпоинт, что и admin
|
# URL для health-пробы: возвращает exit-IP JSON'ом. Тот же эндпоинт, что и admin
|
||||||
# /scraper/health (_probe_current_ip).
|
# /scraper/health (_probe_current_ip).
|
||||||
_HEALTH_PROBE_URL = "https://api.ipify.org"
|
_HEALTH_PROBE_URL = "https://api.ipify.org"
|
||||||
|
|
@ -156,6 +197,11 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
affinity='any' ИЛИ у этой affinity есть ДРУГОЙ enabled-узел (EXISTS-подзапрос) —
|
affinity='any' ИЛИ у этой affinity есть ДРУГОЙ enabled-узел (EXISTS-подзапрос) —
|
||||||
т.е. выдача не обнулит доступность выделенной affinity целиком.
|
т.е. выдача не обнулит доступность выделенной affinity целиком.
|
||||||
|
|
||||||
|
ОБА запроса отсекают узлы с АКТИВНЫМ баном по ЭТОМУ provider'у
|
||||||
|
(scrape_proxy_source_bans.banned_until > now(), #2600 п.2) — узел, забаненный Авито,
|
||||||
|
остаётся полноценным кандидатом для Яндекса и остальных источников. Бан по чужому
|
||||||
|
source на выдачу не влияет вообще.
|
||||||
|
|
||||||
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
||||||
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
||||||
|
|
||||||
|
|
@ -173,6 +219,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
AND consecutive_fails < CAST(:max_fails AS integer)
|
AND consecutive_fails < CAST(:max_fails AS integer)
|
||||||
AND provider_affinity IN (:provider, 'any')
|
AND provider_affinity IN (:provider, 'any')
|
||||||
AND leased_by IS NULL
|
AND leased_by IS NULL
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM scrape_proxy_source_bans b
|
||||||
|
WHERE b.proxy_id = scrape_proxies.id
|
||||||
|
AND b.source = :provider
|
||||||
|
AND b.banned_until > now()
|
||||||
|
)
|
||||||
ORDER BY last_ok_at NULLS LAST, id
|
ORDER BY last_ok_at NULLS LAST, id
|
||||||
FOR UPDATE SKIP LOCKED
|
FOR UPDATE SKIP LOCKED
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
|
|
@ -199,14 +252,35 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
WHERE sp.enabled
|
WHERE sp.enabled
|
||||||
AND sp.consecutive_fails < CAST(:max_fails AS integer)
|
AND sp.consecutive_fails < CAST(:max_fails AS integer)
|
||||||
AND sp.leased_by IS NULL
|
AND sp.leased_by IS NULL
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM scrape_proxy_source_bans b
|
||||||
|
WHERE b.proxy_id = sp.id
|
||||||
|
AND b.source = :provider
|
||||||
|
AND b.banned_until > now()
|
||||||
|
)
|
||||||
AND (
|
AND (
|
||||||
sp.provider_affinity = 'any'
|
sp.provider_affinity = 'any'
|
||||||
|
-- backup обязан быть ПРИГОДЕН для своей affinity, а не просто
|
||||||
|
-- enabled (#2600 п.2 deep-review): после перехода на per-source
|
||||||
|
-- баны узел бывает enabled и одновременно забанен СВОИМ же
|
||||||
|
-- источником. Засчитывать такой как backup — значит разрешить
|
||||||
|
-- fallback увести последний реально рабочий узел выделенной
|
||||||
|
-- affinity и обрушить её (два domclick-узла, один забанен
|
||||||
|
-- domclick'ом → второй уходит под avito → domclick без прокси).
|
||||||
OR EXISTS (
|
OR EXISTS (
|
||||||
SELECT 1
|
SELECT 1
|
||||||
FROM scrape_proxies AS other
|
FROM scrape_proxies AS other
|
||||||
WHERE other.provider_affinity = sp.provider_affinity
|
WHERE other.provider_affinity = sp.provider_affinity
|
||||||
AND other.enabled
|
AND other.enabled
|
||||||
AND other.id <> sp.id
|
AND other.id <> sp.id
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM scrape_proxy_source_bans b2
|
||||||
|
WHERE b2.proxy_id = other.id
|
||||||
|
AND b2.source = other.provider_affinity
|
||||||
|
AND b2.banned_until > now()
|
||||||
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
ORDER BY sp.last_ok_at NULLS LAST, sp.id
|
ORDER BY sp.last_ok_at NULLS LAST, sp.id
|
||||||
|
|
@ -214,7 +288,7 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
{"max_fails": MAX_CONSECUTIVE_FAILS},
|
{"max_fails": MAX_CONSECUTIVE_FAILS, "provider": provider},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.fetchone()
|
.fetchone()
|
||||||
|
|
@ -407,92 +481,146 @@ def mark_health(
|
||||||
|
|
||||||
|
|
||||||
def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
|
def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
|
||||||
"""Пометить прокси НАДЁЖНО забаненным площадкой `source` (#2600 п.1).
|
"""Записать бан узла площадкой `source` — по ПАРЕ (proxy_id, source), #2600 п.2.
|
||||||
|
|
||||||
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
|
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
|
||||||
авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация —
|
авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация —
|
||||||
транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно
|
транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно
|
||||||
распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер — НЕ
|
распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер — НЕ
|
||||||
исключение транспорта, НЕ голый network-fail) — узел выключается НЕМЕДЛЕННО,
|
исключение транспорта, НЕ голый network-fail).
|
||||||
без ожидания порога: `enabled=false`, `disabled_reason='banned:<source>'`. Тот же
|
|
||||||
non-NULL `disabled_reason`, что и ручное выключение (#2610) — `mark_health(ok=True)`
|
|
||||||
больше не воскресит узел голым ipify-пробой: сам ipify площадку не эмулирует, бана
|
|
||||||
не увидит (ровно баг, который #2600 описывает как корень проблемы).
|
|
||||||
|
|
||||||
ЗАЩИТА ПОСЛЕДНЕГО ЖИВОГО УЗЛА (issue #2600 риск, переиспользует паттерн #2609):
|
ЧТО ИМЕННО ДЕЛАЕТСЯ (изменение против #2600 п.1): узел БОЛЬШЕ НЕ выключается
|
||||||
если это последний узел, из-за которого `acquire(source)` вообще способен что-то
|
глобально (`enabled=false, disabled_reason='banned:<source>'` — так было в п.1).
|
||||||
вернуть — НЕ выключаем, только WARNING-лог. Доступность считается ТЕМ ЖЕ правилом,
|
Пишется строка в `scrape_proxy_source_bans` (миграция 210): пока
|
||||||
что и acquire() (primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не
|
`banned_until > now()`, `acquire(source)` этот узел не выдаёт, а для ЛЮБОГО
|
||||||
последняя из своей) — EXISTS-проверка ниже, а не наивный `COUNT(*) WHERE enabled`.
|
другого источника он остаётся первосортным. Авито банит IP — Яндекс через тот же
|
||||||
Без этого защита не сработала бы, например, если проверять только
|
IP ходит чисто; глобальное выключение выкидывало живой узел отовсюду и худило пул
|
||||||
"affinity=proxy_id.provider_affinity" (fallback от #2609 позволяет чужой affinity
|
в разы быстрее, чем его пополняют (#2638). `enabled`/`disabled_reason` остаются
|
||||||
подменить исчерпанную) — источник мог бы остаться без единого узла, даже когда
|
исключительно за оператором (#2610) и за авто-disable'ом по транспортным сбоям.
|
||||||
формально в пуле есть живые строки другой выделенной affinity (domclick).
|
|
||||||
|
|
||||||
КОНКУРЕНТНОСТЬ (deep-review fix 2): один `UPDATE ... WHERE ... AND EXISTS(...)` —
|
ЭСКАЛАЦИЯ: первый бан пары — SOURCE_BAN_BASE_HOURS; каждый следующий удваивает
|
||||||
НЕ атомарная гарантия поперёк СТРОК. EXISTS читает состояние других строк на момент
|
срок (ban_count после инкремента N → SOURCE_BAN_BASE_HOURS * 2^(N-1)), но не выше
|
||||||
своего снапшота (READ COMMITTED), но не лочит их — два ПАРАЛЛЕЛЬНЫХ mark_banned для
|
SOURCE_BAN_MAX_HOURS. Узел, который площадка банит раз за разом, отдыхает от неё
|
||||||
РАЗНЫХ proxy_id (напр. avito банит A, cian банит B миллисекундами позже) каждый может
|
всё дольше, вместо того чтобы жечь прогоны. Сброс ban_count — только purge'ем
|
||||||
увидеть другого как "ещё живого" в своём EXISTS и оба закоммититься → пул уходит с 2
|
истёкших строк (run_proxy_healthcheck, SOURCE_BAN_PURGE_DAYS).
|
||||||
живых узлов в 0 разом. Это НЕ самолечится (#2610: disabled_reason блокирует ipify-
|
|
||||||
воскрешение, нужен ручной PATCH). Фикс: `pg_advisory_xact_lock` в начале транзакции
|
ЗАЩИТА ПОСЛЕДНЕГО УЗЛА, ТЕПЕРЬ ПО ИСТОЧНИКУ (issue #2600 риск, паттерн #2609):
|
||||||
сериализует ВСЕ mark_banned-вызовы между собой (xact-scoped — снимается сам на
|
если после записи бана у `acquire(source)` не останется НИ ОДНОГО кандидата — бан
|
||||||
commit/rollback, leak невозможен). Один глобальный ключ вместо per-affinity —
|
НЕ пишется, только WARNING. Доступность считается ТЕМ ЖЕ правилом, что и acquire()
|
||||||
сериализует и НЕ пересекающиеся по affinity баны тоже (avito vs domclick не гонятся
|
(primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не последняя из
|
||||||
за одни строки), но частота вызовов низкая (несколько банов в час, не hot-path) —
|
своей) ПЛЮС отсутствие активной бан-строки для этого source — EXISTS ниже, а не
|
||||||
цена оправдана простотой против per-row `SELECT ... FOR UPDATE` по кандидатам
|
наивный `COUNT(*) WHERE enabled`. Голодать без прокси хуже, чем ходить через
|
||||||
(потребовал бы лочить весь EXISTS-кандидат-сет заранее, выше риск deadlock между
|
забаненный: капча хотя бы иногда пропускает, отсутствие узла — нет.
|
||||||
параллельными mark_banned, лочащими пересекающиеся строки в разном порядке).
|
|
||||||
ponytail: global advisory lock, не per-affinity — переходи на составной ключ
|
`leased_by IS NULL` защита НАМЕРЕННО не проверяет (в отличие от acquire) — так было
|
||||||
|
и в п.1, и это не оплошность: lease живёт минуты-часы и снимается сам (release /
|
||||||
|
reap_stale_leases), т.е. занятый узел — это доступный узел через мгновение, а вот
|
||||||
|
отказ записать бан из-за чужого lease был бы вечным (узел так и остался бы в выдаче
|
||||||
|
забаненным). Точность здесь не бесплатна: с проверкой lease защита срабатывала бы
|
||||||
|
ложно при каждом параллельном прогоне.
|
||||||
|
|
||||||
|
КОНКУРЕНТНОСТЬ (deep-review fix 2 из #2600 п.1, сохранено): один
|
||||||
|
`INSERT ... WHERE EXISTS(...)` — НЕ атомарная гарантия поперёк СТРОК. EXISTS читает
|
||||||
|
состояние других строк на момент своего снапшота (READ COMMITTED), но не лочит их —
|
||||||
|
два ПАРАЛЛЕЛЬНЫХ mark_banned для РАЗНЫХ proxy_id (напр. avito банит A, cian банит B
|
||||||
|
миллисекундами позже) каждый может увидеть другого как "ещё живого" в своём EXISTS и
|
||||||
|
оба закоммититься → для источника не остаётся ни одного узла разом. Фикс:
|
||||||
|
`pg_advisory_xact_lock` в начале транзакции сериализует ВСЕ mark_banned-вызовы между
|
||||||
|
собой (xact-scoped — снимается сам на commit/rollback, leak невозможен). Один
|
||||||
|
глобальный ключ вместо per-source — сериализует и непересекающиеся баны тоже, но
|
||||||
|
частота вызовов низкая (несколько банов в час, не hot-path) — цена оправдана
|
||||||
|
простотой против per-row `SELECT ... FOR UPDATE` по кандидатам (выше риск deadlock
|
||||||
|
между параллельными mark_banned, лочащими пересекающиеся строки в разном порядке).
|
||||||
|
ponytail: global advisory lock, не per-source — переходи на составной ключ
|
||||||
(напр. hashtext(source)) если частота банов когда-нибудь станет hot-path.
|
(напр. hashtext(source)) если частота банов когда-нибудь станет hot-path.
|
||||||
|
|
||||||
Идемпотентно: узел уже `enabled=false` (ранее забанен/выключен вручную) — no-op,
|
Идемпотентно: повторный бан той же пары не создаёт дубль (PK (proxy_id, source)) —
|
||||||
INFO-лог, disabled_reason НЕ перезаписывается другим source (WHERE enabled в UPDATE).
|
продлевает срок по правилу эскалации. Несуществующий proxy_id — no-op + WARNING.
|
||||||
|
|
||||||
Best-effort по контракту вызывающих (`BrowserFetcher.report_ban`, `curl_proxy_url`) —
|
Best-effort по контракту вызывающих (`BrowserFetcher.report_ban`, `curl_proxy_url`) —
|
||||||
сюда попадают уже обёрнутыми в try/except, но сам мark_banned ошибки БД не глотает
|
сюда попадают уже обёрнутыми в try/except, но сам mark_banned ошибки БД не глотает
|
||||||
(падает как обычно) — caller решает, ловить или нет.
|
(падает как обычно) — caller решает, ловить или нет.
|
||||||
"""
|
"""
|
||||||
# Сериализует check+update ниже с другими конкурентными mark_banned (см. докстринг
|
# Сериализует check+insert ниже с другими конкурентными mark_banned (см. докстринг
|
||||||
# "КОНКУРЕНТНОСТЬ"). Держится до db.commit()/rollback() этой транзакции.
|
# "КОНКУРЕНТНОСТЬ"). Держится до db.commit()/rollback() этой транзакции.
|
||||||
db.execute(
|
db.execute(
|
||||||
text("SELECT pg_advisory_xact_lock(CAST(:key AS bigint))"),
|
text("SELECT pg_advisory_xact_lock(CAST(:key AS bigint))"),
|
||||||
{"key": _MARK_BANNED_ADVISORY_LOCK_KEY},
|
{"key": _MARK_BANNED_ADVISORY_LOCK_KEY},
|
||||||
)
|
)
|
||||||
|
# INSERT ... SELECT ... WHERE EXISTS: guard'ы в WHERE источника строк — не прошли,
|
||||||
|
# значит строк на вставку нет, конфликта нет, эскалации нет (0 rows → ветка логов ниже).
|
||||||
|
# LEAST(ban_count, 16) в показателе — страховка от переполнения double при абсурдном
|
||||||
|
# ban_count (потолок SOURCE_BAN_MAX_HOURS всё равно срежет результат гораздо раньше).
|
||||||
row = (
|
row = (
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_proxies
|
INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason)
|
||||||
SET enabled = false,
|
SELECT CAST(:proxy_id AS bigint),
|
||||||
disabled_reason = CAST(:reason AS text),
|
CAST(:source AS text),
|
||||||
updated_at = now()
|
now() + make_interval(hours => CAST(:base_hours AS integer)),
|
||||||
|
CAST(:reason AS text)
|
||||||
|
WHERE EXISTS (
|
||||||
|
SELECT 1 FROM scrape_proxies
|
||||||
WHERE id = CAST(:proxy_id AS bigint)
|
WHERE id = CAST(:proxy_id AS bigint)
|
||||||
AND enabled
|
)
|
||||||
AND EXISTS (
|
AND EXISTS (
|
||||||
SELECT 1
|
SELECT 1
|
||||||
FROM scrape_proxies sp
|
FROM scrape_proxies sp
|
||||||
WHERE sp.id <> CAST(:proxy_id AS bigint)
|
WHERE sp.id <> CAST(:proxy_id AS bigint)
|
||||||
AND sp.enabled
|
AND sp.enabled
|
||||||
AND sp.consecutive_fails < CAST(:max_fails AS integer)
|
AND sp.consecutive_fails < CAST(:max_fails AS integer)
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM scrape_proxy_source_bans b
|
||||||
|
WHERE b.proxy_id = sp.id
|
||||||
|
AND b.source = CAST(:source AS text)
|
||||||
|
AND b.banned_until > now()
|
||||||
|
)
|
||||||
AND (
|
AND (
|
||||||
sp.provider_affinity IN (:source, 'any')
|
sp.provider_affinity IN (:source, 'any')
|
||||||
|
-- other.id <> sp.id (а не NOT IN (sp.id, :proxy_id), как в
|
||||||
|
-- п.1): банимый узел остаётся enabled и по-прежнему обслуживает
|
||||||
|
-- СВОЮ affinity — значит он и есть валидный backup для неё.
|
||||||
|
-- NOT EXISTS b2 — тот же критерий пригодности, что в acquire()
|
||||||
|
-- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом
|
||||||
|
-- не считается (иначе защита сочла бы affinity живой, когда она
|
||||||
|
-- уже нет).
|
||||||
OR EXISTS (
|
OR EXISTS (
|
||||||
SELECT 1
|
SELECT 1
|
||||||
FROM scrape_proxies other
|
FROM scrape_proxies other
|
||||||
WHERE other.provider_affinity = sp.provider_affinity
|
WHERE other.provider_affinity = sp.provider_affinity
|
||||||
AND other.enabled
|
AND other.enabled
|
||||||
AND other.id NOT IN (sp.id, CAST(:proxy_id AS bigint))
|
AND other.id <> sp.id
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM scrape_proxy_source_bans b2
|
||||||
|
WHERE b2.proxy_id = other.id
|
||||||
|
AND b2.source = other.provider_affinity
|
||||||
|
AND b2.banned_until > now()
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
RETURNING id
|
)
|
||||||
|
ON CONFLICT (proxy_id, source) DO UPDATE
|
||||||
|
SET ban_count = scrape_proxy_source_bans.ban_count + 1,
|
||||||
|
banned_until = now() + make_interval(hours => CAST(
|
||||||
|
LEAST(
|
||||||
|
CAST(:base_hours AS integer)
|
||||||
|
* power(2, LEAST(scrape_proxy_source_bans.ban_count, 16)),
|
||||||
|
CAST(:max_hours AS integer)
|
||||||
|
) AS integer)),
|
||||||
|
reason = CAST(:reason AS text),
|
||||||
|
updated_at = now()
|
||||||
|
RETURNING ban_count, banned_until
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
{
|
{
|
||||||
"proxy_id": proxy_id,
|
"proxy_id": proxy_id,
|
||||||
"source": source,
|
"source": source,
|
||||||
"reason": f"banned:{source}",
|
"reason": f"banned:{source}",
|
||||||
|
"base_hours": SOURCE_BAN_BASE_HOURS,
|
||||||
|
"max_hours": SOURCE_BAN_MAX_HOURS,
|
||||||
"max_fails": MAX_CONSECUTIVE_FAILS,
|
"max_fails": MAX_CONSECUTIVE_FAILS,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
@ -502,16 +630,18 @@ def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
|
||||||
db.commit()
|
db.commit()
|
||||||
if row is not None:
|
if row is not None:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"proxy_pool: proxy id=%d BANNED by source=%s — disabled "
|
"proxy_pool: proxy id=%d BANNED by source=%s — узел снят с выдачи ТОЛЬКО для "
|
||||||
"(disabled_reason='banned:%s'), ipify-проба больше НЕ воскресит (#2610)",
|
"этого источника до %s (ban_count=%s); для остальных источников остаётся в "
|
||||||
|
"строю (#2600 п.2)",
|
||||||
proxy_id,
|
proxy_id,
|
||||||
source,
|
source,
|
||||||
source,
|
row["banned_until"],
|
||||||
|
row["ban_count"],
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
# 0 rows: либо уже disabled (идемпотентно, no-op), либо защита последнего узла
|
# 0 rows: либо узла нет, либо защита последнего узла отменила запись бана — читаем
|
||||||
# сработала — читаем текущее состояние ТОЛЬКО для точного лога (не влияет на решение).
|
# текущее состояние ТОЛЬКО для точного лога (на решение уже не влияет).
|
||||||
current = (
|
current = (
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
|
|
@ -525,23 +655,63 @@ def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
|
||||||
)
|
)
|
||||||
if current is None:
|
if current is None:
|
||||||
logger.warning("proxy_pool: mark_banned id=%d not found — no-op", proxy_id)
|
logger.warning("proxy_pool: mark_banned id=%d not found — no-op", proxy_id)
|
||||||
elif not current["enabled"]:
|
|
||||||
logger.info(
|
|
||||||
"proxy_pool: mark_banned id=%d source=%s — already disabled (reason=%r), no-op",
|
|
||||||
proxy_id,
|
|
||||||
source,
|
|
||||||
current["disabled_reason"],
|
|
||||||
)
|
|
||||||
else:
|
else:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"proxy_pool: proxy id=%d BANNED by source=%s but NOT disabled — last live node "
|
"proxy_pool: proxy id=%d — бан не записан: это последний узел, достижимый для "
|
||||||
"reachable for this source (protection, mirrors acquire() fallback rule, #2600). "
|
"source=%s; нужны новые прокси (см. #2638). Узел продолжит выдаваться этому "
|
||||||
"Ban logged only — operator should investigate/add capacity.",
|
"источнику (голодание хуже, чем работа через забаненный узел).",
|
||||||
proxy_id,
|
proxy_id,
|
||||||
source,
|
source,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def clear_source_bans(db: Session, proxy_id: int, *, source: str | None = None, reason: str) -> int:
|
||||||
|
"""Снять баны узла по источникам (#2600 п.2). Returns число снятых строк.
|
||||||
|
|
||||||
|
ЗАЧЕМ ОТДЕЛЬНАЯ РУЧКА: до п.2 ложный бан лечился оператором через
|
||||||
|
`PATCH /proxies/{id} enabled=true` — включение обнуляло `disabled_reason`, и узел
|
||||||
|
возвращался в строй. Теперь бан живёт в отдельной таблице и сам по себе истекает
|
||||||
|
только по таймеру, вплоть до 72 часов при эскалации. Без этой функции ложное
|
||||||
|
срабатывание детектора капчи (#2642) парковало бы узел на часы, а снять это можно
|
||||||
|
было бы только руками в SQL.
|
||||||
|
|
||||||
|
ГДЕ ВЫЗЫВАЕТСЯ:
|
||||||
|
- `admin.patch_proxy` при ручном включении узла — «оператор включил» означает
|
||||||
|
чистый лист, ровно как обнуление disabled_reason рядом (#2610);
|
||||||
|
- после УСПЕШНОЙ ротации exit-IP (`proxy_rotation.rotate_proxy`) — площадка
|
||||||
|
банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса,
|
||||||
|
держа узел вне выдачи уже без причины.
|
||||||
|
|
||||||
|
source=None — снять все баны узла; конкретный source — только его. DELETE, а не
|
||||||
|
`banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации),
|
||||||
|
а здесь мы как раз объявляем историю недействительной — новый бан начнётся с базовых
|
||||||
|
SOURCE_BAN_BASE_HOURS.
|
||||||
|
|
||||||
|
`reason` идёт только в лог (человекочитаемый повод — «manual enable», «ip rotated»).
|
||||||
|
"""
|
||||||
|
rows = db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
DELETE FROM scrape_proxy_source_bans
|
||||||
|
WHERE proxy_id = CAST(:proxy_id AS bigint)
|
||||||
|
AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text))
|
||||||
|
RETURNING source
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"proxy_id": proxy_id, "source": source},
|
||||||
|
).fetchall()
|
||||||
|
db.commit()
|
||||||
|
if rows:
|
||||||
|
logger.info(
|
||||||
|
"proxy_pool: cleared %d source ban(s) for proxy id=%d (%s) — reason=%s",
|
||||||
|
len(rows),
|
||||||
|
proxy_id,
|
||||||
|
[r.source for r in rows],
|
||||||
|
reason,
|
||||||
|
)
|
||||||
|
return len(rows)
|
||||||
|
|
||||||
|
|
||||||
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:
|
||||||
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
|
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
|
||||||
|
|
||||||
|
|
@ -635,9 +805,12 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
следующего disable/enable цикла), но mark_health их не воскрешает — revived не растёт,
|
следующего disable/enable цикла), но mark_health их не воскрешает — revived не растёт,
|
||||||
WARNING пишет сам mark_health.
|
WARNING пишет сам mark_health.
|
||||||
|
|
||||||
|
В конце — purge бан-строк (#2600 п.2), истёкших дольше SOURCE_BAN_PURGE_DAYS назад
|
||||||
|
(см. комментарий у самого DELETE: отложенность — это и есть сброс ban_count).
|
||||||
|
|
||||||
Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на
|
Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на
|
||||||
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
|
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
|
||||||
{reaped, checked, ok, failed, revived}.
|
{reaped, checked, ok, failed, revived, bans_purged}.
|
||||||
"""
|
"""
|
||||||
reaped = reap_stale_leases(db)
|
reaped = reap_stale_leases(db)
|
||||||
|
|
||||||
|
|
@ -687,13 +860,35 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
else:
|
else:
|
||||||
failed += 1
|
failed += 1
|
||||||
|
|
||||||
|
# Purge ДАВНО истёкших бан-строк (#2600 п.2). Порог — banned_until + SOURCE_BAN_PURGE_DAYS,
|
||||||
|
# НЕ просто `banned_until < now()`: строка после истечения бана ещё ничего не блокирует
|
||||||
|
# (acquire фильтрует по banned_until > now()), но хранит ban_count — память об эскалации.
|
||||||
|
# Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с
|
||||||
|
# 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса:
|
||||||
|
# неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать».
|
||||||
|
purged = len(
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
DELETE FROM scrape_proxy_source_bans
|
||||||
|
WHERE banned_until < now() - make_interval(days => CAST(:days AS integer))
|
||||||
|
RETURNING proxy_id
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"days": SOURCE_BAN_PURGE_DAYS},
|
||||||
|
).fetchall()
|
||||||
|
)
|
||||||
|
db.commit()
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d revived=%d",
|
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d revived=%d "
|
||||||
|
"bans_purged=%d",
|
||||||
reaped,
|
reaped,
|
||||||
checked,
|
checked,
|
||||||
ok_count,
|
ok_count,
|
||||||
failed,
|
failed,
|
||||||
revived,
|
revived,
|
||||||
|
purged,
|
||||||
)
|
)
|
||||||
return {
|
return {
|
||||||
"reaped": reaped,
|
"reaped": reaped,
|
||||||
|
|
@ -701,4 +896,5 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
"ok": ok_count,
|
"ok": ok_count,
|
||||||
"failed": failed,
|
"failed": failed,
|
||||||
"revived": revived,
|
"revived": revived,
|
||||||
|
"bans_purged": purged,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -72,6 +72,7 @@ from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.core.config import settings
|
from app.core.config import settings
|
||||||
|
from app.services.proxy_pool import clear_source_bans
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -363,6 +364,10 @@ async def rotate_proxy(db: Session, proxy_id: int) -> RotationResult:
|
||||||
new_ip = _extract_new_ip(resp)
|
new_ip = _extract_new_ip(resp)
|
||||||
logger.info("proxy_rotation: rotated proxy_id=%d status=%d new_ip=%s", proxy_id, status, new_ip)
|
logger.info("proxy_rotation: rotated proxy_id=%d status=%d new_ip=%s", proxy_id, status, new_ip)
|
||||||
_record_attempt(db, proxy_id, success=True, http_status=status, note=None)
|
_record_attempt(db, proxy_id, success=True, http_status=status, note=None)
|
||||||
|
# Площадки банили СТАРЫЙ exit-IP, а строка бана привязана к proxy_id (#2600 п.2) —
|
||||||
|
# после смены адреса она держала бы узел вне выдачи уже без причины, вплоть до 72ч
|
||||||
|
# при эскалации. Ротация прошла → история банов этого узла недействительна.
|
||||||
|
clear_source_bans(db, proxy_id, reason=f"exit ip rotated (status={status})")
|
||||||
return RotationResult(
|
return RotationResult(
|
||||||
ok=True,
|
ok=True,
|
||||||
reason=None,
|
reason=None,
|
||||||
|
|
|
||||||
105
tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql
Normal file
105
tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql
Normal file
|
|
@ -0,0 +1,105 @@
|
||||||
|
-- 210_scrape_proxy_source_bans.sql
|
||||||
|
-- Здоровье прокси по ПАРЕ «узел × источник» (#2600 п.2).
|
||||||
|
--
|
||||||
|
-- WHY:
|
||||||
|
-- До сих пор бан был ГЛОБАЛЬНЫМ: #2600 п.1 (mark_banned) на распознанный бан
|
||||||
|
-- площадкой выключал узел целиком — enabled=false, disabled_reason='banned:<source>'.
|
||||||
|
-- Реальность другая: Авито банит IP, а Яндекс через тот же IP ходит чисто. Один
|
||||||
|
-- забаненный источник выкидывал живой узел из пула для ВСЕХ источников, пул худел
|
||||||
|
-- в разы быстрее, чем его успевают пополнять (#2638).
|
||||||
|
--
|
||||||
|
-- WHAT:
|
||||||
|
-- scrape_proxy_source_bans — по строке на пару (proxy_id, source). Пока
|
||||||
|
-- banned_until > now(), acquire(source) этот узел НЕ выдаёт; для любого ДРУГОГО
|
||||||
|
-- источника узел остаётся первосортным. Узел больше не выключается глобально —
|
||||||
|
-- enabled/disabled_reason остаются исключительно за оператором (#2610) и за
|
||||||
|
-- авто-disable'ом по серии транспортных сбоев (mark_health).
|
||||||
|
--
|
||||||
|
-- ban_count — счётчик повторных банов той же пары: срок эскалирует
|
||||||
|
-- 6ч → 12ч → 24ч → 48ч → 72ч (потолок), см. SOURCE_BAN_BASE_HOURS/
|
||||||
|
-- SOURCE_BAN_MAX_HOURS в app/services/proxy_pool.py. Истёкшие строки НЕ
|
||||||
|
-- удаляются сразу — purge в run_proxy_healthcheck сносит их только через 7 суток
|
||||||
|
-- после истечения (SOURCE_BAN_PURGE_DAYS), и это же механизм сброса ban_count:
|
||||||
|
-- узел, неделю чистый после снятия бана, начинает эскалацию с нуля.
|
||||||
|
--
|
||||||
|
-- КОНВЕРСИЯ СТАРЫХ ГЛОБАЛЬНЫХ БАНОВ (обязательная часть миграции):
|
||||||
|
-- После #2600 п.1 на проде могли остаться узлы enabled=false с
|
||||||
|
-- disabled_reason LIKE 'banned:%'. Новый код такой семантики больше НЕ пишет и
|
||||||
|
-- ничего её не снимает, а mark_health(ok=True) не воскрешает узлы с non-NULL
|
||||||
|
-- disabled_reason (#2610) — узел завис бы выключенным навсегда, до ручного PATCH.
|
||||||
|
-- Поэтому здесь каждый такой узел конвертируется в per-source бан на 6 часов
|
||||||
|
-- (тот же SOURCE_BAN_BASE_HOURS) и возвращается в строй: enabled=true,
|
||||||
|
-- disabled_reason=NULL. Матчинг по ТОЧНОМУ списку 'banned:<источник>', а не по
|
||||||
|
-- LIKE — ручные тексты оператора (в т.ч. начинающиеся с 'banned:', этот формат
|
||||||
|
-- подсказан комментарием 209-й) НЕ трогаются, это его решение.
|
||||||
|
--
|
||||||
|
-- IDEMPOTENCY / SAFETY:
|
||||||
|
-- - Весь файл в одной транзакции BEGIN/COMMIT.
|
||||||
|
-- - CREATE TABLE / INDEX IF NOT EXISTS, INSERT ... ON CONFLICT DO NOTHING →
|
||||||
|
-- повторный прогон no-op (auto-apply strict на деплое это требует).
|
||||||
|
-- - Конверсионный UPDATE после повторного прогона не находит строк (первый
|
||||||
|
-- прогон уже снял disabled_reason) — тоже no-op.
|
||||||
|
-- - ON DELETE CASCADE: удаление прокси уносит его баны, «висячих» строк нет.
|
||||||
|
--
|
||||||
|
-- Dependencies: 157_scrape_proxies.sql, 209_scrape_proxies_disabled_reason.sql
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS scrape_proxy_source_bans (
|
||||||
|
proxy_id bigint NOT NULL REFERENCES scrape_proxies(id) ON DELETE CASCADE,
|
||||||
|
source text NOT NULL,
|
||||||
|
banned_until timestamptz NOT NULL,
|
||||||
|
reason text,
|
||||||
|
ban_count integer NOT NULL DEFAULT 1,
|
||||||
|
banned_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
PRIMARY KEY (proxy_id, source)
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE scrape_proxy_source_bans IS
|
||||||
|
'Баны прокси по паре (узел, источник), #2600 п.2. Активна строка с '
|
||||||
|
'banned_until > now() — acquire(source) такой узел не выдаёт, для других '
|
||||||
|
'источников узел остаётся доступным. Глобальное выключение узла (enabled=false) '
|
||||||
|
'сюда НЕ относится — это ручное действие оператора или авто-disable по серии '
|
||||||
|
'транспортных сбоев.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN scrape_proxy_source_bans.ban_count IS
|
||||||
|
'Сколько раз эта пара банилась. Срок ТЕКУЩЕГО бана (banned_until - banned_at) = '
|
||||||
|
'base * 2^(ban_count-1), потолок SOURCE_BAN_MAX_HOURS: ban_count=1 → 6ч, 2 → 12ч, '
|
||||||
|
'3 → 24ч и т.д. Сбрасывается удалением строки — либо purge''ем через '
|
||||||
|
'SOURCE_BAN_PURGE_DAYS после истечения, либо proxy_pool.clear_source_bans '
|
||||||
|
'(ручное включение узла оператором / успешная ротация exit-IP).';
|
||||||
|
|
||||||
|
-- Горячий путь — NOT EXISTS-фильтр в acquire(): (proxy_id, source) уже покрыт PK,
|
||||||
|
-- этот индекс закрывает purge/листинг активных банов по времени.
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_scrape_proxy_source_bans_until
|
||||||
|
ON scrape_proxy_source_bans (banned_until);
|
||||||
|
|
||||||
|
-- ── конверсия старых глобальных банов (#2600 п.1 → п.2) ─────────────────────
|
||||||
|
--
|
||||||
|
-- ТОЧНЫЙ список значений, а не LIKE 'banned:%': 209-я миграция сама предлагает этот
|
||||||
|
-- формат в комментарии, поэтому оператор мог написать руками что-то вроде
|
||||||
|
-- 'banned:avito вручную'. LIKE тогда дал бы source='avito вручную' (бан-строка, которая
|
||||||
|
-- ни с чем не сматчится) и МОЛЧА отменил бы ручное выключение. Домен ниже — тот же, что
|
||||||
|
-- у scrape_proxies.provider_affinity (на практике mark_banned п.1 писал только
|
||||||
|
-- avito/cian/yandex/domclick — это значения BrowserFetcher._source).
|
||||||
|
INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason)
|
||||||
|
SELECT id,
|
||||||
|
substring(disabled_reason from 8), -- отрезает префикс 'banned:' (7 символов)
|
||||||
|
now() + interval '6 hours',
|
||||||
|
'migrated from disabled_reason (210)'
|
||||||
|
FROM scrape_proxies
|
||||||
|
WHERE NOT enabled
|
||||||
|
AND disabled_reason IN ('banned:avito', 'banned:cian', 'banned:yandex',
|
||||||
|
'banned:domclick', 'banned:generic', 'banned:any')
|
||||||
|
ON CONFLICT (proxy_id, source) DO NOTHING;
|
||||||
|
|
||||||
|
UPDATE scrape_proxies
|
||||||
|
SET enabled = true,
|
||||||
|
disabled_reason = NULL,
|
||||||
|
updated_at = now()
|
||||||
|
WHERE NOT enabled
|
||||||
|
AND disabled_reason IN ('banned:avito', 'banned:cian', 'banned:yandex',
|
||||||
|
'banned:domclick', 'banned:generic', 'banned:any');
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -26,6 +26,19 @@ reap_stale_leases проверяются по фактическому изме
|
||||||
* ручно-выключенный узел (disabled_reason НЕ NULL) НЕ воскресает даже при
|
* ручно-выключенный узел (disabled_reason НЕ NULL) НЕ воскресает даже при
|
||||||
ok=True, и это логируется (WARNING);
|
ok=True, и это логируется (WARNING);
|
||||||
* run_proxy_healthcheck не считает ручно-выключенный узел в revived.
|
* run_proxy_healthcheck не считает ручно-выключенный узел в revived.
|
||||||
|
- бан по паре «узел × источник» (#2600 п.2, scrape_proxy_source_bans):
|
||||||
|
* mark_banned пишет строку бана и НЕ выключает узел глобально;
|
||||||
|
* повторный бан той же пары эскалирует срок (ban_count растёт);
|
||||||
|
* acquire не выдаёт узел с активным баном по ЭТОМУ source, но выдаёт по ДРУГОМУ
|
||||||
|
(суть задачи: Авито забанил — Яндекс продолжает ходить через тот же узел);
|
||||||
|
* истёкший бан снова не мешает выдаче;
|
||||||
|
* защита последнего узла: бан НЕ записывается, если у acquire(source) не
|
||||||
|
останется кандидатов;
|
||||||
|
* run_proxy_healthcheck сносит бан-строки, истёкшие дольше SOURCE_BAN_PURGE_DAYS,
|
||||||
|
и НЕ трогает истёкшие недавно (в них живёт ban_count для эскалации);
|
||||||
|
* backup выделенной affinity засчитывается, только если он сам не забанен своим
|
||||||
|
источником (иначе fallback уводил бы последний рабочий узел, deep-review);
|
||||||
|
* clear_source_bans снимает баны узла (все / один source) и обнуляет эскалацию.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -44,6 +57,9 @@ from app.services.proxy_pool import (
|
||||||
DISABLE_THRESHOLD,
|
DISABLE_THRESHOLD,
|
||||||
DISABLED_RECHECK_MINUTES,
|
DISABLED_RECHECK_MINUTES,
|
||||||
MAX_CONSECUTIVE_FAILS,
|
MAX_CONSECUTIVE_FAILS,
|
||||||
|
SOURCE_BAN_BASE_HOURS,
|
||||||
|
SOURCE_BAN_MAX_HOURS,
|
||||||
|
SOURCE_BAN_PURGE_DAYS,
|
||||||
STALE_LEASE_MINUTES,
|
STALE_LEASE_MINUTES,
|
||||||
acquire,
|
acquire,
|
||||||
mark_health,
|
mark_health,
|
||||||
|
|
@ -88,10 +104,13 @@ class _FakeResult:
|
||||||
|
|
||||||
|
|
||||||
class FakeSession:
|
class FakeSession:
|
||||||
"""Эмуляция Session поверх in-memory списка scrape_proxies-строк."""
|
"""Эмуляция Session поверх in-memory списков scrape_proxies + scrape_proxy_source_bans."""
|
||||||
|
|
||||||
def __init__(self, rows: list[dict[str, Any]]):
|
def __init__(self, rows: list[dict[str, Any]], bans: list[dict[str, Any]] | None = None):
|
||||||
self.rows = rows
|
self.rows = rows
|
||||||
|
# #2600 п.2: строки scrape_proxy_source_bans (proxy_id, source, banned_until,
|
||||||
|
# ban_count) — бан теперь свойство ПАРЫ «узел × источник», не узла.
|
||||||
|
self.bans: list[dict[str, Any]] = bans or []
|
||||||
# deep-review fix 2 (#2600): вызовы pg_advisory_xact_lock — для теста
|
# deep-review fix 2 (#2600): вызовы pg_advisory_xact_lock — для теста
|
||||||
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
|
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
|
||||||
# проверить — фейксессия однопоточна).
|
# проверить — фейксессия однопоточна).
|
||||||
|
|
@ -101,14 +120,30 @@ class FakeSession:
|
||||||
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
||||||
return next((r for r in self.rows if r["id"] == pid), None)
|
return next((r for r in self.rows if r["id"] == pid), None)
|
||||||
|
|
||||||
|
def _ban(self, pid: int, source: str) -> dict[str, Any] | None:
|
||||||
|
return next((b for b in self.bans if b["proxy_id"] == pid and b["source"] == source), None)
|
||||||
|
|
||||||
|
def _has_active_ban(self, pid: int, source: str) -> bool:
|
||||||
|
ban = self._ban(pid, source)
|
||||||
|
return ban is not None and ban["banned_until"] > datetime.now(UTC)
|
||||||
|
|
||||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||||
sql = str(stmt)
|
sql = str(stmt)
|
||||||
p = params or {}
|
p = params or {}
|
||||||
|
|
||||||
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
|
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
|
||||||
max_fails = p["max_fails"]
|
max_fails = p["max_fails"]
|
||||||
if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any'
|
|
||||||
provider = p["provider"]
|
provider = p["provider"]
|
||||||
|
# #2600 п.2: узел с активным баном по ЭТОМУ source не выдаётся. Как и с
|
||||||
|
# protects_last_node ниже — фильтруем ТОЛЬКО если сам SQL реально содержит
|
||||||
|
# NOT EXISTS по scrape_proxy_source_bans, иначе мок реализовывал бы логику
|
||||||
|
# независимо от проверяемого кода и не отличил бы старый запрос от нового.
|
||||||
|
filters_bans = "scrape_proxy_source_bans" in sql
|
||||||
|
|
||||||
|
def _not_banned(row: dict[str, Any]) -> bool:
|
||||||
|
return not filters_bans or not self._has_active_ban(row["id"], provider)
|
||||||
|
|
||||||
|
if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any'
|
||||||
cands = [
|
cands = [
|
||||||
r
|
r
|
||||||
for r in self.rows
|
for r in self.rows
|
||||||
|
|
@ -116,6 +151,7 @@ class FakeSession:
|
||||||
and r["consecutive_fails"] < max_fails
|
and r["consecutive_fails"] < max_fails
|
||||||
and r["provider_affinity"] in (provider, "any")
|
and r["provider_affinity"] in (provider, "any")
|
||||||
and r["leased_by"] is None
|
and r["leased_by"] is None
|
||||||
|
and _not_banned(r)
|
||||||
]
|
]
|
||||||
else: # fallback: любая affinity, но не последний узел выделенной affinity
|
else: # fallback: любая affinity, но не последний узел выделенной affinity
|
||||||
# (domclick и т.п. — #2600 review). ВАЖНО: применяем эту фильтрацию,
|
# (domclick и т.п. — #2600 review). ВАЖНО: применяем эту фильтрацию,
|
||||||
|
|
@ -124,6 +160,10 @@ class FakeSession:
|
||||||
# кода и не смог бы отличить старый (незащищённый) fallback-запрос от
|
# кода и не смог бы отличить старый (незащищённый) fallback-запрос от
|
||||||
# нового. Тот же класс бага, что был с "enabled" в mark_health-моке.
|
# нового. Тот же класс бага, что был с "enabled" в mark_health-моке.
|
||||||
protects_last_node = "EXISTS" in sql
|
protects_last_node = "EXISTS" in sql
|
||||||
|
# backup засчитывается, только если он ПРИГОДЕН для своей affinity —
|
||||||
|
# тоже гейтим по подстроке (b2-подзапрос), иначе мок «чинил» бы
|
||||||
|
# незащищённый SQL сам.
|
||||||
|
backup_must_be_usable = "b2.banned_until > now()" in sql
|
||||||
|
|
||||||
def _has_backup(row: dict[str, Any]) -> bool:
|
def _has_backup(row: dict[str, Any]) -> bool:
|
||||||
if row["provider_affinity"] == "any":
|
if row["provider_affinity"] == "any":
|
||||||
|
|
@ -132,6 +172,10 @@ class FakeSession:
|
||||||
other["provider_affinity"] == row["provider_affinity"]
|
other["provider_affinity"] == row["provider_affinity"]
|
||||||
and other["enabled"]
|
and other["enabled"]
|
||||||
and other["id"] != row["id"]
|
and other["id"] != row["id"]
|
||||||
|
and not (
|
||||||
|
backup_must_be_usable
|
||||||
|
and self._has_active_ban(other["id"], other["provider_affinity"])
|
||||||
|
)
|
||||||
for other in self.rows
|
for other in self.rows
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -141,6 +185,7 @@ class FakeSession:
|
||||||
if r["enabled"]
|
if r["enabled"]
|
||||||
and r["consecutive_fails"] < max_fails
|
and r["consecutive_fails"] < max_fails
|
||||||
and r["leased_by"] is None
|
and r["leased_by"] is None
|
||||||
|
and _not_banned(r)
|
||||||
and (not protects_last_node or _has_backup(r))
|
and (not protects_last_node or _has_backup(r))
|
||||||
]
|
]
|
||||||
# ORDER BY last_ok_at NULLS LAST, id
|
# ORDER BY last_ok_at NULLS LAST, id
|
||||||
|
|
@ -233,34 +278,81 @@ class FakeSession:
|
||||||
self.advisory_lock_calls.append(p["key"])
|
self.advisory_lock_calls.append(p["key"])
|
||||||
return _FakeResult([])
|
return _FakeResult([])
|
||||||
|
|
||||||
if "SET enabled = false" in sql: # mark_banned UPDATE (#2600 п.1)
|
if "INSERT INTO scrape_proxy_source_bans" in sql: # mark_banned UPSERT (#2600 п.2)
|
||||||
proxy_id = p["proxy_id"]
|
proxy_id = p["proxy_id"]
|
||||||
source = p["source"]
|
source = p["source"]
|
||||||
max_fails = p["max_fails"]
|
max_fails = p["max_fails"]
|
||||||
row = self._by_id(proxy_id)
|
if self._by_id(proxy_id) is None:
|
||||||
if row is None or not row["enabled"]:
|
return _FakeResult([]) # узла нет — no-op
|
||||||
return _FakeResult([]) # already disabled / not found — no-op
|
|
||||||
|
# Оба ban-предиката гейтим по подстрокам самого SQL (как в acquire-ветке):
|
||||||
|
# иначе мок реализовывал бы защиту сам и тест оставался бы зелёным даже
|
||||||
|
# после удаления NOT EXISTS из боевого запроса.
|
||||||
|
filters_bans = "b.banned_until > now()" in sql
|
||||||
|
backup_must_be_usable = "b2.banned_until > now()" in sql
|
||||||
|
|
||||||
def _is_candidate(sp: dict[str, Any]) -> bool:
|
def _is_candidate(sp: dict[str, Any]) -> bool:
|
||||||
if not (sp["enabled"] and sp["consecutive_fails"] < max_fails):
|
if not (sp["enabled"] and sp["consecutive_fails"] < max_fails):
|
||||||
return False
|
return False
|
||||||
|
if filters_bans and self._has_active_ban(sp["id"], source):
|
||||||
|
return False # уже забанен этим же источником — не кандидат
|
||||||
if sp["provider_affinity"] in (source, "any"):
|
if sp["provider_affinity"] in (source, "any"):
|
||||||
return True
|
return True
|
||||||
# fallback-safe: другой enabled узел ТОЙ ЖЕ affinity (кроме sp/proxy_id).
|
# fallback-safe: другой ПРИГОДНЫЙ узел ТОЙ ЖЕ affinity (банимый узел
|
||||||
|
# остаётся enabled и тоже считается — бан теперь per-source; а вот
|
||||||
|
# забаненный своим же источником backup'ом не считается).
|
||||||
return any(
|
return any(
|
||||||
other["provider_affinity"] == sp["provider_affinity"]
|
other["provider_affinity"] == sp["provider_affinity"]
|
||||||
and other["enabled"]
|
and other["enabled"]
|
||||||
and other["id"] not in (sp["id"], proxy_id)
|
and other["id"] != sp["id"]
|
||||||
|
and not (
|
||||||
|
backup_must_be_usable
|
||||||
|
and self._has_active_ban(other["id"], other["provider_affinity"])
|
||||||
|
)
|
||||||
for other in self.rows
|
for other in self.rows
|
||||||
)
|
)
|
||||||
|
|
||||||
still_available = any(r["id"] != proxy_id and _is_candidate(r) for r in self.rows)
|
still_available = any(r["id"] != proxy_id and _is_candidate(r) for r in self.rows)
|
||||||
if not still_available:
|
if not still_available:
|
||||||
return _FakeResult([]) # protected — последний живой узел, не выключаем
|
return _FakeResult([]) # protected — последний узел для source, бан не пишем
|
||||||
|
|
||||||
row["enabled"] = False
|
now = datetime.now(UTC)
|
||||||
row["disabled_reason"] = p["reason"]
|
ban = self._ban(proxy_id, source)
|
||||||
return _FakeResult([{"id": proxy_id}])
|
if ban is None:
|
||||||
|
ban = {
|
||||||
|
"proxy_id": proxy_id,
|
||||||
|
"source": source,
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": now + timedelta(hours=p["base_hours"]),
|
||||||
|
"reason": p["reason"],
|
||||||
|
}
|
||||||
|
self.bans.append(ban)
|
||||||
|
else:
|
||||||
|
# эскалация: срок = base * 2^(новый ban_count - 1), потолок max_hours
|
||||||
|
ban["ban_count"] += 1
|
||||||
|
hours = min(p["base_hours"] * 2 ** (ban["ban_count"] - 1), p["max_hours"])
|
||||||
|
ban["banned_until"] = now + timedelta(hours=hours)
|
||||||
|
ban["reason"] = p["reason"]
|
||||||
|
return _FakeResult(
|
||||||
|
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
|
||||||
|
)
|
||||||
|
|
||||||
|
if "DELETE FROM scrape_proxy_source_bans" in sql and "proxy_id = CAST" in sql:
|
||||||
|
# clear_source_bans: снять баны узла (все либо один source), #2600 п.2
|
||||||
|
cleared = [
|
||||||
|
b
|
||||||
|
for b in self.bans
|
||||||
|
if b["proxy_id"] == p["proxy_id"]
|
||||||
|
and (p["source"] is None or b["source"] == p["source"])
|
||||||
|
]
|
||||||
|
self.bans = [b for b in self.bans if b not in cleared]
|
||||||
|
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||||
|
|
||||||
|
if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2)
|
||||||
|
cutoff = datetime.now(UTC) - timedelta(days=p["days"])
|
||||||
|
purged = [{"proxy_id": b["proxy_id"]} for b in self.bans if b["banned_until"] < cutoff]
|
||||||
|
self.bans = [b for b in self.bans if b["banned_until"] >= cutoff]
|
||||||
|
return _FakeResult(purged)
|
||||||
|
|
||||||
if "SELECT enabled, disabled_reason FROM scrape_proxies" in sql: # mark_banned diag read
|
if "SELECT enabled, disabled_reason FROM scrape_proxies" in sql: # mark_banned diag read
|
||||||
row = self._by_id(p["id"])
|
row = self._by_id(p["id"])
|
||||||
|
|
@ -415,6 +507,33 @@ def test_acquire_fallback_allows_when_dedicated_affinity_has_backup() -> None:
|
||||||
assert db._by_id(2)["leased_by"] is None # у domclick остался живой запасной узел
|
assert db._by_id(2)["leased_by"] is None # у domclick остался живой запасной узел
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_fallback_backup_must_be_usable_for_its_own_source() -> None:
|
||||||
|
"""deep-review #2600 п.2: «backup» — это ПРИГОДНЫЙ узел, а не просто enabled.
|
||||||
|
|
||||||
|
Два узла domclick: node1 забанен САМИМ domclick'ом (законно — node2 тогда был жив) и
|
||||||
|
сейчас занят чужим прогоном, node2 свободен. Для domclick node2 — последний рабочий.
|
||||||
|
Раньше EXISTS видел node1 как backup (он ведь enabled) и разрешал fallback увести
|
||||||
|
node2 под avito — domclick оставался бы без прокси вообще при двух включённых узлах.
|
||||||
|
"""
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="domclick", leased_by=99), _proxy(2, affinity="domclick")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "domclick",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:domclick",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||||||
|
assert db._by_id(2)["leased_by"] is None # последний рабочий узел domclick не тронут
|
||||||
|
# сам domclick при этом обслуживается: node2 свободен и не забанен
|
||||||
|
lease = acquire(db, "domclick", run_id=2) # type: ignore[arg-type]
|
||||||
|
assert lease is not None and lease.id == 2
|
||||||
|
|
||||||
|
|
||||||
# ── release ──────────────────────────────────────────────────────────────────
|
# ── release ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -704,68 +823,96 @@ async def test_healthcheck_checks_disabled_proxy_never_checked_before(
|
||||||
assert db._by_id(1)["enabled"] is False
|
assert db._by_id(1)["enabled"] is False
|
||||||
|
|
||||||
|
|
||||||
# ── mark_banned (#2600 п.1 — довести сигнал бана до пула) ──────────────────────
|
# ── mark_banned (#2600 п.2 — бан по паре «узел × источник») ────────────────────
|
||||||
#
|
#
|
||||||
# Red/green контракт issue: (a) распознанный бан → узел выключен с
|
# Red/green контракт issue: (a) распознанный бан → строка в scrape_proxy_source_bans,
|
||||||
# disabled_reason='banned:<source>'; (b) ipify-проба его не воскрешает (уже
|
# узел НЕ выключен глобально (в п.1 было enabled=false — площадка забанила IP, а не
|
||||||
# покрыто disabled_reason-веткой mark_health выше, #2610 — здесь только
|
# сломала прокси; узел обязан остаться живым для остальных источников); (b) повторный
|
||||||
# убеждаемся, что mark_banned проставляет ТОТ ЖЕ non-NULL disabled_reason);
|
# бан той же пары эскалирует срок; (c) последний достижимый для source узел НЕ банится
|
||||||
# (c) последний живой узел НЕ выключается (только лог); (d) сетевой сбой
|
# (только лог); (d) сетевой сбой по-прежнему идёт через mark_health(ok=False), НЕ через
|
||||||
# по-прежнему идёт через mark_health(ok=False), НЕ через mark_banned (проверяется
|
# mark_banned (проверяется на уровне browser_fetcher/curl_proxy_url тестов — здесь
|
||||||
# на уровне browser_fetcher/curl_proxy_url тестов — здесь mark_banned сам по себе
|
# mark_banned сам по себе не участвует в различении причин, это забота caller'а).
|
||||||
# не участвует в различении причин, это забота вызывающего кода).
|
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_disables_with_reason() -> None:
|
def test_mark_banned_writes_source_ban_row() -> None:
|
||||||
"""Забаненный узел выключается, disabled_reason='banned:<source>' (второй здоровый
|
"""Бан пишется в scrape_proxy_source_bans: source, ban_count=1, срок = base-часы."""
|
||||||
узел affinity='any' есть — защита последнего узла не срабатывает)."""
|
assert mark_banned is not None
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
ban = db._ban(1, "avito")
|
||||||
|
assert ban is not None
|
||||||
|
assert ban["ban_count"] == 1
|
||||||
|
assert ban["reason"] == "banned:avito"
|
||||||
|
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS)
|
||||||
|
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_does_not_disable_node_globally() -> None:
|
||||||
|
"""ГЛАВНОЕ отличие от #2600 п.1: узел остаётся enabled и без disabled_reason —
|
||||||
|
глобальное выключение теперь только за оператором/авто-disable'ом (#2610)."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
row = db._by_id(1)
|
row = db._by_id(1)
|
||||||
assert row["enabled"] is False
|
assert row["enabled"] is True
|
||||||
assert row["disabled_reason"] == "banned:avito"
|
assert row["disabled_reason"] is None
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_disabled_reason_blocks_auto_revive() -> None:
|
def test_mark_banned_repeat_escalates_ban_count_and_duration() -> None:
|
||||||
"""disabled_reason non-NULL после mark_banned → mark_health(ok=True) НЕ
|
"""Повторный бан той же пары: ban_count растёт, срок удваивается (base * 2^(N-1))."""
|
||||||
воскрешает узел (та же ветка #2610, что и ручное выключение)."""
|
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
first_until = db._ban(1, "avito")["banned_until"]
|
||||||
|
|
||||||
mark_health(db, 1, ok=True) # type: ignore[arg-type]
|
|
||||||
|
|
||||||
row = db._by_id(1)
|
|
||||||
assert row["enabled"] is False # НЕ воскрешён голым ipify-успехом
|
|
||||||
assert row["disabled_reason"] == "banned:avito"
|
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_idempotent_already_disabled() -> None:
|
|
||||||
"""Уже выключенный узел (ручной либо предыдущий бан) — no-op, причина не
|
|
||||||
перезаписывается ДРУГИМ source."""
|
|
||||||
assert mark_banned is not None
|
|
||||||
db = FakeSession(
|
|
||||||
[
|
|
||||||
_proxy(1, affinity="avito", enabled=False, disabled_reason="banned:cian"),
|
|
||||||
_proxy(2, affinity="any"),
|
|
||||||
]
|
|
||||||
)
|
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
row = db._by_id(1)
|
|
||||||
assert row["enabled"] is False
|
ban = db._ban(1, "avito")
|
||||||
assert row["disabled_reason"] == "banned:cian" # не перезаписан
|
assert ban["ban_count"] == 2
|
||||||
|
assert len(db.bans) == 1 # PK (proxy_id, source) — дубля нет
|
||||||
|
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS * 2)
|
||||||
|
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||||||
|
assert ban["banned_until"] > first_until
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_escalation_capped_at_max_hours() -> None:
|
||||||
|
"""Эскалация упирается в SOURCE_BAN_MAX_HOURS, а не растёт до бесконечности."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||||||
|
for _ in range(8):
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
ban = db._ban(1, "avito")
|
||||||
|
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS)
|
||||||
|
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_different_sources_are_independent_rows() -> None:
|
||||||
|
"""Бан Авито и бан Циана на одном узле — две независимые строки, не перезапись."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
mark_banned(db, 1, source="cian") # type: ignore[arg-type]
|
||||||
|
assert {b["source"] for b in db.bans} == {"avito", "cian"}
|
||||||
|
assert all(b["ban_count"] == 1 for b in db.bans)
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_unknown_id_is_noop() -> None:
|
||||||
|
"""Несуществующий proxy_id — бан не пишется, падать не должно (best-effort caller)."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||||||
|
mark_banned(db, 999, source="avito") # type: ignore[arg-type]
|
||||||
|
assert db.bans == []
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_protects_last_live_node_own_affinity() -> None:
|
def test_mark_banned_protects_last_live_node_own_affinity() -> None:
|
||||||
"""Единственный узел avito, других (свободных/'any') нет вообще — НЕ выключается,
|
"""Единственный узел avito, других (свободных/'any') нет вообще — бан НЕ пишется,
|
||||||
только лог (issue #2600: бан не должен обрушить единственный источник целиком)."""
|
только лог (issue #2600: бан не должен обрушить единственный источник целиком —
|
||||||
|
ходить через забаненный узел лучше, чем не ходить вообще)."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito")])
|
db = FakeSession([_proxy(1, affinity="avito")])
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
row = db._by_id(1)
|
assert db.bans == []
|
||||||
assert row["enabled"] is True # НЕ тронут
|
assert db._by_id(1)["enabled"] is True # и глобально не тронут
|
||||||
assert row["disabled_reason"] is None
|
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_protects_last_live_node_logs_warning(
|
def test_mark_banned_protects_last_live_node_logs_warning(
|
||||||
|
|
@ -775,17 +922,38 @@ def test_mark_banned_protects_last_live_node_logs_warning(
|
||||||
db = FakeSession([_proxy(1, affinity="avito")])
|
db = FakeSession([_proxy(1, affinity="avito")])
|
||||||
with caplog.at_level("WARNING"):
|
with caplog.at_level("WARNING"):
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
assert any("last live node" in rec.message for rec in caplog.records)
|
assert any("бан не записан" in rec.message for rec in caplog.records)
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_disables_when_any_affinity_backup_exists() -> None:
|
def test_mark_banned_protection_counts_active_bans_of_other_nodes() -> None:
|
||||||
|
"""Второй узел формально жив, но уже забанен ЭТИМ ЖЕ источником — кандидатом для
|
||||||
|
source он не является, значит бан первого узла оставил бы avito без прокси вообще.
|
||||||
|
Защита обязана сработать (иначе оба узла разом выпадут для avito)."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="avito"), _proxy(2, affinity="any")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 2,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
assert db._ban(1, "avito") is None # бан первого узла не записан
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_writes_when_any_affinity_backup_exists() -> None:
|
||||||
"""Забанен единственный avito-специфичный узел, но есть 'any' — 'any' закрывает
|
"""Забанен единственный avito-специфичный узел, но есть 'any' — 'any' закрывает
|
||||||
availability для avito (та же семантика, что acquire()'s primary IN (provider,
|
availability для avito (та же семантика, что acquire()'s primary IN (provider,
|
||||||
'any')) → выключаем безопасно."""
|
'any')) → бан записывается безопасно."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
assert db._by_id(1)["enabled"] is False
|
assert db._ban(1, "avito") is not None
|
||||||
|
|
||||||
|
|
||||||
# ── mark_banned: affinity-aware last-node (orchestrator follow-up, свежий прод-факт) ──
|
# ── mark_banned: affinity-aware last-node (orchestrator follow-up, свежий прод-факт) ──
|
||||||
|
|
@ -799,20 +967,19 @@ def test_mark_banned_disables_when_any_affinity_backup_exists() -> None:
|
||||||
def test_mark_banned_dedicated_affinity_alone_does_not_count_as_backup_for_other_source() -> None:
|
def test_mark_banned_dedicated_affinity_alone_does_not_count_as_backup_for_other_source() -> None:
|
||||||
"""avito банится; в пуле остаётся только один domclick-узел (affinity выделенная,
|
"""avito банится; в пуле остаётся только один domclick-узел (affinity выделенная,
|
||||||
БЕЗ backup) — для avito это НЕ доступный узел (acquire('avito') не взял бы его через
|
БЕЗ backup) — для avito это НЕ доступный узел (acquire('avito') не взял бы его через
|
||||||
fallback, #2609 protects last node of domclick). Защита должна сработать — avito-узел
|
fallback, #2609 protects last node of domclick). Защита должна сработать — бан
|
||||||
НЕ выключается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
|
НЕ записывается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="domclick")])
|
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="domclick")])
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
row = db._by_id(1)
|
assert db.bans == [] # domclick-узел НЕ считается доступной заменой
|
||||||
assert row["enabled"] is True # domclick-узел НЕ считается доступной заменой
|
|
||||||
assert db._by_id(2)["enabled"] is True # и сам не тронут
|
assert db._by_id(2)["enabled"] is True # и сам не тронут
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None:
|
def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None:
|
||||||
"""Та же ситуация, но у domclick есть ВТОРОЙ узел (backup) — тогда fallback может
|
"""Та же ситуация, но у domclick есть ВТОРОЙ узел (backup) — тогда fallback может
|
||||||
забрать ОДИН из них под avito (acquire()'s EXISTS-правило #2609), доступность для
|
забрать ОДИН из них под avito (acquire()'s EXISTS-правило #2609), доступность для
|
||||||
avito сохраняется через fallback → banned avito-узел безопасно выключается."""
|
avito сохраняется через fallback → бан avito-узла записывается безопасно."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession(
|
db = FakeSession(
|
||||||
[
|
[
|
||||||
|
|
@ -822,7 +989,37 @@ def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
assert db._by_id(1)["enabled"] is False
|
assert db._ban(1, "avito") is not None
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_banned_backup_of_dedicated_affinity_must_be_usable() -> None:
|
||||||
|
"""Зеркало acquire-правила в защите (deep-review #2600 п.2).
|
||||||
|
|
||||||
|
Узел 3 (domclick) забанен САМИМ domclick'ом и вдобавок в карантине по fails, т.е.
|
||||||
|
сам заменой для avito быть не может. Узел 2 — последний РАБОЧИЙ узел domclick,
|
||||||
|
fallback не имеет права его забрать. Значит замены для avito нет вообще → бан
|
||||||
|
avito-узла НЕ пишется. Со старым правилом («backup = любой enabled той же affinity»)
|
||||||
|
узел 3 засчитался бы бэкапом, узел 2 стал бы «доступным» и бан бы записался.
|
||||||
|
"""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession(
|
||||||
|
[
|
||||||
|
_proxy(1, affinity="avito"),
|
||||||
|
_proxy(2, affinity="domclick"),
|
||||||
|
_proxy(3, affinity="domclick", fails=MAX_CONSECUTIVE_FAILS),
|
||||||
|
],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 3,
|
||||||
|
"source": "domclick",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:domclick",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
assert db._ban(1, "avito") is None
|
||||||
|
|
||||||
|
|
||||||
def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
|
def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
|
||||||
|
|
@ -833,7 +1030,7 @@ def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
|
||||||
[_proxy(1, affinity="avito"), _proxy(2, affinity="any", fails=MAX_CONSECUTIVE_FAILS)]
|
[_proxy(1, affinity="avito"), _proxy(2, affinity="any", fails=MAX_CONSECUTIVE_FAILS)]
|
||||||
)
|
)
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
assert db._by_id(1)["enabled"] is True # карантинный узел не спасает
|
assert db.bans == [] # карантинный узел не спасает
|
||||||
|
|
||||||
|
|
||||||
# ── mark_banned: TOCTOU-защита (deep-review fix 2, #2600) ───────────────────────
|
# ── mark_banned: TOCTOU-защита (deep-review fix 2, #2600) ───────────────────────
|
||||||
|
|
@ -857,10 +1054,213 @@ def test_mark_banned_takes_advisory_xact_lock_with_fixed_key() -> None:
|
||||||
|
|
||||||
def test_mark_banned_takes_advisory_lock_even_when_protected() -> None:
|
def test_mark_banned_takes_advisory_lock_even_when_protected() -> None:
|
||||||
"""Лок берётся ПЕРЕД проверкой доступности — даже когда защита последнего узла
|
"""Лок берётся ПЕРЕД проверкой доступности — даже когда защита последнего узла
|
||||||
в итоге отменяет disable, лок всё равно взят (сериализация check+decide, не
|
в итоге отменяет запись бана, лок всё равно взят (сериализация check+decide, не
|
||||||
только update)."""
|
только сам INSERT)."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession([_proxy(1, affinity="avito")]) # единственный узел — protected
|
db = FakeSession([_proxy(1, affinity="avito")]) # единственный узел — protected
|
||||||
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
assert db.advisory_lock_calls # лок взят, хотя disable не произошёл
|
assert db.advisory_lock_calls # лок взят, хотя бан не записан
|
||||||
assert db._by_id(1)["enabled"] is True
|
assert db.bans == []
|
||||||
|
|
||||||
|
|
||||||
|
# ── acquire × бан по источнику (#2600 п.2 — суть задачи) ───────────────────────
|
||||||
|
#
|
||||||
|
# Прод-факт, из-за которого всё затевалось: Авито банит IP, Яндекс через тот же IP
|
||||||
|
# ходит чисто. Глобальный enabled=false выкидывал живой узел из пула для ВСЕХ
|
||||||
|
# источников; per-source бан обязан выключать выдачу ровно одному.
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_skips_node_banned_for_this_source() -> None:
|
||||||
|
"""Единственный узел забанен ЭТИМ источником → acquire(source) возвращает None."""
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
assert acquire(db, "avito", run_id=1) is None # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_still_issues_node_banned_by_other_source() -> None:
|
||||||
|
"""ГЛАВНЫЙ тест задачи: узел забанен Авито — Яндексу он выдаётся как ни в чём не
|
||||||
|
бывало (бан — свойство пары, а не узла)."""
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
lease = acquire(db, "yandex", run_id=2) # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert lease.id == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_ignores_expired_ban() -> None:
|
||||||
|
"""Истёкшая бан-строка (banned_until в прошлом) выдаче не мешает — она ещё лежит
|
||||||
|
только ради ban_count (purge снесёт её позже, SOURCE_BAN_PURGE_DAYS)."""
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="avito")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 2,
|
||||||
|
"banned_until": datetime.now(UTC) - timedelta(hours=1),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
assert lease is not None and lease.id == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_acquire_fallback_also_respects_source_ban() -> None:
|
||||||
|
"""Fallback-заход (чужая affinity) тоже отсекает забаненные для source узлы: два
|
||||||
|
cian-узла, один забанен avito → avito достаётся ВТОРОЙ, а не забаненный."""
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="cian"), _proxy(2, affinity="cian")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 1,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
)
|
||||||
|
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
assert lease is not None
|
||||||
|
assert lease.id == 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_ban_then_acquire_end_to_end() -> None:
|
||||||
|
"""Сквозной сценарий: mark_banned('avito') на узле 1 → avito получает узел 2,
|
||||||
|
yandex по-прежнему может получить узел 1."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession([_proxy(1, affinity="any"), _proxy(2, affinity="any")])
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
|
||||||
|
avito_lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
assert avito_lease is not None and avito_lease.id == 2
|
||||||
|
|
||||||
|
release(db, 2) # type: ignore[arg-type]
|
||||||
|
yandex_lease = acquire(db, "yandex", run_id=2) # type: ignore[arg-type]
|
||||||
|
assert yandex_lease is not None and yandex_lease.id == 1 # забанен только для avito
|
||||||
|
|
||||||
|
|
||||||
|
# ── purge истёкших бан-строк в run_proxy_healthcheck (#2600 п.2) ───────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_healthcheck_purges_long_expired_bans_only(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Сносим строки, истёкшие дольше SOURCE_BAN_PURGE_DAYS назад; истёкшую вчера
|
||||||
|
оставляем — в ней живёт ban_count (память об эскалации, см. комментарий у DELETE)."""
|
||||||
|
now = datetime.now(UTC)
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any")],
|
||||||
|
bans=[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"ban_count": 3,
|
||||||
|
"banned_until": now - timedelta(days=SOURCE_BAN_PURGE_DAYS + 1),
|
||||||
|
"reason": "banned:avito",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "cian",
|
||||||
|
"ban_count": 2,
|
||||||
|
"banned_until": now - timedelta(days=1),
|
||||||
|
"reason": "banned:cian",
|
||||||
|
},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
|
||||||
|
async def _fake_probe(url: str) -> tuple[bool, str | None, int | None, str | None]:
|
||||||
|
return True, "1.2.3.4", 10, None
|
||||||
|
|
||||||
|
monkeypatch.setattr(proxy_pool, "_probe_proxy", _fake_probe)
|
||||||
|
|
||||||
|
counters = await proxy_pool.run_proxy_healthcheck(db) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert counters["bans_purged"] == 1
|
||||||
|
assert [b["source"] for b in db.bans] == ["cian"]
|
||||||
|
|
||||||
|
|
||||||
|
# ── clear_source_bans (#2600 п.2 — рычаг оператора против ложного бана) ────────
|
||||||
|
#
|
||||||
|
# До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан
|
||||||
|
# в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой
|
||||||
|
# ручки ложное срабатывание детектора капчи (#2642) снималось бы только руками в SQL.
|
||||||
|
|
||||||
|
|
||||||
|
def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"proxy_id": pid,
|
||||||
|
"source": source,
|
||||||
|
"ban_count": ban_count,
|
||||||
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS),
|
||||||
|
"reason": f"banned:{source}",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_source_bans_removes_all_bans_of_node() -> None:
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||||||
|
bans=[_active_ban(1, "avito"), _active_ban(1, "cian"), _active_ban(2, "avito")],
|
||||||
|
)
|
||||||
|
cleared = proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
|
||||||
|
assert cleared == 2
|
||||||
|
assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел
|
||||||
|
# узел снова выдаётся источнику, который его банил
|
||||||
|
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
|
assert lease is not None and lease.id == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_source_bans_single_source_keeps_others() -> None:
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any")],
|
||||||
|
bans=[_active_ban(1, "avito"), _active_ban(1, "cian")],
|
||||||
|
)
|
||||||
|
cleared = proxy_pool.clear_source_bans( # type: ignore[arg-type]
|
||||||
|
db, 1, source="avito", reason="ip rotated"
|
||||||
|
)
|
||||||
|
assert cleared == 1
|
||||||
|
assert [b["source"] for b in db.bans] == ["cian"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
|
||||||
|
db = FakeSession([_proxy(1, affinity="any")])
|
||||||
|
assert proxy_pool.clear_source_bans(db, 1, reason="manual enable") == 0 # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_source_bans_resets_escalation() -> None:
|
||||||
|
"""DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан
|
||||||
|
начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию."""
|
||||||
|
assert mark_banned is not None
|
||||||
|
db = FakeSession(
|
||||||
|
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||||||
|
bans=[_active_ban(1, "avito", ban_count=4)],
|
||||||
|
)
|
||||||
|
proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
|
||||||
|
|
||||||
|
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
|
||||||
|
|
||||||
|
ban = db._ban(1, "avito")
|
||||||
|
assert ban["ban_count"] == 1
|
||||||
|
expected = datetime.now(UTC) + timedelta(hours=SOURCE_BAN_BASE_HOURS)
|
||||||
|
assert abs((ban["banned_until"] - expected).total_seconds()) < 60
|
||||||
|
|
|
||||||
|
|
@ -48,6 +48,10 @@ class _FakeResult:
|
||||||
def fetchone(self) -> dict[str, Any] | None:
|
def fetchone(self) -> dict[str, Any] | None:
|
||||||
return self._rows[0] if self._rows else None
|
return self._rows[0] if self._rows else None
|
||||||
|
|
||||||
|
def fetchall(self) -> list[Any]:
|
||||||
|
# RETURNING source у clear_source_bans — код читает r.source (attribute access)
|
||||||
|
return [type("Row", (), r)() for r in self._rows]
|
||||||
|
|
||||||
|
|
||||||
class FakeSession:
|
class FakeSession:
|
||||||
"""Эмуляция Session: одна строка scrape_proxies + append-only
|
"""Эмуляция Session: одна строка scrape_proxies + append-only
|
||||||
|
|
@ -58,9 +62,13 @@ class FakeSession:
|
||||||
self,
|
self,
|
||||||
proxy_row: dict[str, Any] | None,
|
proxy_row: dict[str, Any] | None,
|
||||||
rotations: list[dict[str, Any]] | None = None,
|
rotations: list[dict[str, Any]] | None = None,
|
||||||
|
source_bans: list[dict[str, Any]] | None = None,
|
||||||
):
|
):
|
||||||
self.proxy_row = proxy_row
|
self.proxy_row = proxy_row
|
||||||
self.rotations: list[dict[str, Any]] = rotations or []
|
self.rotations: list[dict[str, Any]] = rotations or []
|
||||||
|
# #2600 п.2: успешная ротация снимает баны узла (IP сменился — бан старого
|
||||||
|
# адреса недействителен), см. proxy_pool.clear_source_bans.
|
||||||
|
self.source_bans: list[dict[str, Any]] = source_bans or []
|
||||||
self.commits = 0
|
self.commits = 0
|
||||||
|
|
||||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||||
|
|
@ -96,6 +104,11 @@ class FakeSession:
|
||||||
)
|
)
|
||||||
return _FakeResult([])
|
return _FakeResult([])
|
||||||
|
|
||||||
|
if "DELETE FROM scrape_proxy_source_bans" in sql: # clear_source_bans (#2600 п.2)
|
||||||
|
cleared = [b for b in self.source_bans if b["proxy_id"] == p["proxy_id"]]
|
||||||
|
self.source_bans = [b for b in self.source_bans if b not in cleared]
|
||||||
|
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||||
|
|
||||||
raise AssertionError(f"unhandled SQL: {sql}")
|
raise AssertionError(f"unhandled SQL: {sql}")
|
||||||
|
|
||||||
def commit(self) -> None:
|
def commit(self) -> None:
|
||||||
|
|
@ -328,6 +341,43 @@ async def test_successful_rotation_writes_history_row(monkeypatch: pytest.Monkey
|
||||||
assert db.commits >= 1
|
assert db.commits >= 1
|
||||||
|
|
||||||
|
|
||||||
|
async def test_successful_rotation_clears_source_bans(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""(#2600 п.2) Сменился exit-IP → баны площадок на СТАРОМ адресе недействительны.
|
||||||
|
|
||||||
|
Строка бана привязана к proxy_id, а не к IP — без снятия узел остался бы вне выдачи
|
||||||
|
источнику до 72 часов уже без причины.
|
||||||
|
"""
|
||||||
|
monkeypatch.setattr(proxy_rotation.settings, "asocks_api_token", SECRET_TOKEN)
|
||||||
|
fake_client, _calls = _fake_async_client(response=(200, {"ip": "9.9.9.9"}), exception=None)
|
||||||
|
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
||||||
|
|
||||||
|
db = FakeSession(
|
||||||
|
_proxy_row(),
|
||||||
|
source_bans=[
|
||||||
|
{"proxy_id": 1, "source": "avito"},
|
||||||
|
{"proxy_id": 1, "source": "cian"},
|
||||||
|
{"proxy_id": 2, "source": "avito"}, # чужой узел — не трогаем
|
||||||
|
],
|
||||||
|
)
|
||||||
|
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result.ok is True
|
||||||
|
assert db.source_bans == [{"proxy_id": 2, "source": "avito"}]
|
||||||
|
|
||||||
|
|
||||||
|
async def test_failed_rotation_keeps_source_bans(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Провайдер ответил ошибкой — IP НЕ сменился, баны обязаны остаться."""
|
||||||
|
monkeypatch.setattr(proxy_rotation.settings, "asocks_api_token", SECRET_TOKEN)
|
||||||
|
fake_client, _calls = _fake_async_client(response=(500, None), exception=None)
|
||||||
|
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
||||||
|
|
||||||
|
db = FakeSession(_proxy_row(), source_bans=[{"proxy_id": 1, "source": "avito"}])
|
||||||
|
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result.ok is False
|
||||||
|
assert db.source_bans == [{"proxy_id": 1, "source": "avito"}]
|
||||||
|
|
||||||
|
|
||||||
# ── 401 → loud failure ───────────────────────────────────────────────────────
|
# ── 401 → loud failure ───────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,9 @@
|
||||||
|
|
||||||
Покрытие (db мокается, NO live network/DB):
|
Покрытие (db мокается, NO live network/DB):
|
||||||
- POST /api/v1/admin/proxies/bulk — UPSERT-счётчики, валидация affinity/kind
|
- POST /api/v1/admin/proxies/bulk — UPSERT-счётчики, валидация affinity/kind
|
||||||
- GET /api/v1/admin/proxies — маскировка пароля, фильтры, disabled_reason в ответе
|
- GET /api/v1/admin/proxies — маскировка пароля, фильтры, disabled_reason в ответе,
|
||||||
|
активные баны по источникам (#2600 п.2 — узел бывает enabled=true и при этом не
|
||||||
|
выдаётся конкретному источнику)
|
||||||
- PATCH /api/v1/admin/proxies/{id} — enable/disable, 404
|
- PATCH /api/v1/admin/proxies/{id} — enable/disable, 404
|
||||||
- #2610: PATCH enabled=false ставит disabled_reason (ручное выключение отличимо от
|
- #2610: PATCH enabled=false ставит disabled_reason (ручное выключение отличимо от
|
||||||
авто); PATCH enabled=true сбрасывает disabled_reason в NULL (снова авто-восстанавливаем)
|
авто); PATCH enabled=true сбрасывает disabled_reason в NULL (снова авто-восстанавливаем)
|
||||||
|
|
@ -49,6 +51,28 @@ def _scalar_result(value: object) -> MagicMock:
|
||||||
return res
|
return res
|
||||||
|
|
||||||
|
|
||||||
|
def _cleared_bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock:
|
||||||
|
"""Ответ на DELETE ... RETURNING source (proxy_pool.clear_source_bans, #2600 п.2).
|
||||||
|
|
||||||
|
PATCH enabled=true снимает баны узла по источникам — «ручное включение = чистый
|
||||||
|
лист», как и обнуление disabled_reason рядом.
|
||||||
|
"""
|
||||||
|
res = MagicMock()
|
||||||
|
res.fetchall.return_value = [type("Row", (), r)() for r in (rows or [])]
|
||||||
|
return res
|
||||||
|
|
||||||
|
|
||||||
|
def _bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock:
|
||||||
|
"""Ответ на ВТОРОЙ execute в /proxies-ручках — активные баны по источникам (#2600 п.2).
|
||||||
|
|
||||||
|
list_proxies/patch_proxy после основного запроса дочитывают scrape_proxy_source_bans
|
||||||
|
(_fetch_source_bans), поэтому мок обязан отдавать два разных результата по порядку.
|
||||||
|
"""
|
||||||
|
res = MagicMock()
|
||||||
|
res.mappings.return_value.all.return_value = rows or []
|
||||||
|
return res
|
||||||
|
|
||||||
|
|
||||||
# ── _mask_proxy_url unit ─────────────────────────────────────────────────────
|
# ── _mask_proxy_url unit ─────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -156,7 +180,7 @@ def _proxy_db_row(**over: Any) -> dict[str, Any]:
|
||||||
def test_list_masks_password(client: TestClient, db: MagicMock) -> None:
|
def test_list_masks_password(client: TestClient, db: MagicMock) -> None:
|
||||||
result = MagicMock()
|
result = MagicMock()
|
||||||
result.mappings.return_value.all.return_value = [_proxy_db_row()]
|
result.mappings.return_value.all.return_value = [_proxy_db_row()]
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
r = client.get("/api/v1/admin/proxies")
|
r = client.get("/api/v1/admin/proxies")
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
|
|
@ -174,7 +198,7 @@ def test_list_exposes_disabled_reason(client: TestClient, db: MagicMock) -> None
|
||||||
_proxy_db_row(id=1, enabled=False, disabled_reason=None),
|
_proxy_db_row(id=1, enabled=False, disabled_reason=None),
|
||||||
_proxy_db_row(id=2, enabled=False, disabled_reason="забанен Авито"),
|
_proxy_db_row(id=2, enabled=False, disabled_reason="забанен Авито"),
|
||||||
]
|
]
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
r = client.get("/api/v1/admin/proxies")
|
r = client.get("/api/v1/admin/proxies")
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
|
|
@ -183,6 +207,37 @@ def test_list_exposes_disabled_reason(client: TestClient, db: MagicMock) -> None
|
||||||
assert rows[2]["disabled_reason"] == "забанен Авито" # выключен вручную
|
assert rows[2]["disabled_reason"] == "забанен Авито" # выключен вручную
|
||||||
|
|
||||||
|
|
||||||
|
def test_list_exposes_active_source_bans(client: TestClient, db: MagicMock) -> None:
|
||||||
|
"""(#2600 п.2) Узел enabled=true, но забанен Авито — оператор должен видеть, почему
|
||||||
|
он не выдаётся конкретному источнику; для остальных источников узел в строю."""
|
||||||
|
result = MagicMock()
|
||||||
|
result.mappings.return_value.all.return_value = [
|
||||||
|
_proxy_db_row(id=1, enabled=True),
|
||||||
|
_proxy_db_row(id=2, enabled=True),
|
||||||
|
]
|
||||||
|
db.execute.side_effect = [
|
||||||
|
result,
|
||||||
|
_bans_result(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"proxy_id": 1,
|
||||||
|
"source": "avito",
|
||||||
|
"banned_until": datetime(2026, 8, 5, 12, tzinfo=UTC),
|
||||||
|
"ban_count": 2,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
),
|
||||||
|
]
|
||||||
|
|
||||||
|
r = client.get("/api/v1/admin/proxies")
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
rows = {row["id"]: row for row in r.json()}
|
||||||
|
assert rows[1]["source_bans"] == [
|
||||||
|
{"source": "avito", "banned_until": "2026-08-05T12:00:00+00:00", "ban_count": 2}
|
||||||
|
]
|
||||||
|
assert rows[2]["source_bans"] == [] # чистый узел — пустой список, а не отсутствие поля
|
||||||
|
|
||||||
|
|
||||||
def test_list_passes_filters(client: TestClient, db: MagicMock) -> None:
|
def test_list_passes_filters(client: TestClient, db: MagicMock) -> None:
|
||||||
result = MagicMock()
|
result = MagicMock()
|
||||||
result.mappings.return_value.all.return_value = []
|
result.mappings.return_value.all.return_value = []
|
||||||
|
|
@ -201,7 +256,7 @@ def test_list_passes_filters(client: TestClient, db: MagicMock) -> None:
|
||||||
def test_patch_disable(client: TestClient, db: MagicMock) -> None:
|
def test_patch_disable(client: TestClient, db: MagicMock) -> None:
|
||||||
result = MagicMock()
|
result = MagicMock()
|
||||||
result.mappings.return_value.fetchone.return_value = _proxy_db_row(enabled=False)
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(enabled=False)
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
|
|
@ -227,13 +282,14 @@ def test_patch_disable_sets_disabled_reason_default(client: TestClient, db: Magi
|
||||||
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
||||||
enabled=False, disabled_reason="manually disabled via admin API"
|
enabled=False, disabled_reason="manually disabled via admin API"
|
||||||
)
|
)
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
assert r.json()["disabled_reason"] == "manually disabled via admin API"
|
assert r.json()["disabled_reason"] == "manually disabled via admin API"
|
||||||
# дефолтная причина реально передана в SQL как fallback-параметр
|
# дефолтная причина реально передана в SQL как fallback-параметр (первый execute —
|
||||||
params = db.execute.call_args.args[1]
|
# сам UPDATE; второй, #2600 п.2, дочитывает активные баны по источникам)
|
||||||
|
params = db.execute.call_args_list[0].args[1]
|
||||||
assert params["default_reason"]
|
assert params["default_reason"]
|
||||||
assert params["reason"] is None
|
assert params["reason"] is None
|
||||||
|
|
||||||
|
|
@ -244,12 +300,12 @@ def test_patch_disable_sets_disabled_reason_custom(client: TestClient, db: Magic
|
||||||
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
||||||
enabled=False, disabled_reason="забанен Авито"
|
enabled=False, disabled_reason="забанен Авито"
|
||||||
)
|
)
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False, "reason": "забанен Авито"})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False, "reason": "забанен Авито"})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
assert r.json()["disabled_reason"] == "забанен Авито"
|
assert r.json()["disabled_reason"] == "забанен Авито"
|
||||||
params = db.execute.call_args.args[1]
|
params = db.execute.call_args_list[0].args[1]
|
||||||
assert params["reason"] == "забанен Авито"
|
assert params["reason"] == "забанен Авито"
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -260,15 +316,46 @@ def test_patch_enable_clears_disabled_reason(client: TestClient, db: MagicMock)
|
||||||
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
||||||
enabled=True, disabled_reason=None
|
enabled=True, disabled_reason=None
|
||||||
)
|
)
|
||||||
db.execute.return_value = result
|
db.execute.side_effect = [result, _cleared_bans_result(), _bans_result()]
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
assert r.json()["disabled_reason"] is None
|
assert r.json()["disabled_reason"] is None
|
||||||
params = db.execute.call_args.args[1]
|
params = db.execute.call_args_list[0].args[1]
|
||||||
assert params["enabled"] is True
|
assert params["enabled"] is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_patch_enable_clears_source_bans(client: TestClient, db: MagicMock) -> None:
|
||||||
|
"""(#2600 п.2) Ручное включение = чистый лист: снимаются и per-source баны, иначе у
|
||||||
|
оператора нет способа отменить ложный бан (детектор капчи, #2642) — узел был бы
|
||||||
|
enabled=true и всё равно невыдаваемым источнику до 72 часов."""
|
||||||
|
result = MagicMock()
|
||||||
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
|
||||||
|
enabled=True, disabled_reason=None
|
||||||
|
)
|
||||||
|
db.execute.side_effect = [result, _cleared_bans_result([{"source": "avito"}]), _bans_result()]
|
||||||
|
|
||||||
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
delete_sql = str(db.execute.call_args_list[1].args[0])
|
||||||
|
assert "DELETE FROM scrape_proxy_source_bans" in delete_sql
|
||||||
|
assert db.execute.call_args_list[1].args[1]["proxy_id"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_patch_disable_keeps_source_bans(client: TestClient, db: MagicMock) -> None:
|
||||||
|
"""Выключение узла бан-строки НЕ снимает — снятие это «оператор говорит, что узел
|
||||||
|
в порядке», а выключение утверждает обратное."""
|
||||||
|
result = MagicMock()
|
||||||
|
result.mappings.return_value.fetchone.return_value = _proxy_db_row(enabled=False)
|
||||||
|
db.execute.side_effect = [result, _bans_result()]
|
||||||
|
|
||||||
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
assert not any(
|
||||||
|
"DELETE FROM scrape_proxy_source_bans" in str(c.args[0]) for c in db.execute.call_args_list
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
# ── POST /proxies/bulk — не глушит ручной disable (#2610) ──────────────────
|
# ── POST /proxies/bulk — не глушит ручной disable (#2610) ──────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -323,7 +323,7 @@ class BrowserFetcher:
|
||||||
)
|
)
|
||||||
|
|
||||||
def report_ban(self, reason: str) -> None:
|
def report_ban(self, reason: str) -> None:
|
||||||
"""Пометить ТЕКУЩИЙ lease забаненным площадкой (#2600 п.1).
|
"""Пометить ТЕКУЩИЙ lease забаненным площадкой (#2600 п.1, п.2).
|
||||||
|
|
||||||
Вызывать из точки детекта бана (заглушка HTTP 200 / капча / QRATOR-маркер),
|
Вызывать из точки детекта бана (заглушка HTTP 200 / капча / QRATOR-маркер),
|
||||||
ПОКА lease ещё держится (до `__aexit__`/`_release_lease`) — `fetch()` уже
|
ПОКА lease ещё держится (до `__aexit__`/`_release_lease`) — `fetch()` уже
|
||||||
|
|
@ -336,8 +336,11 @@ class BrowserFetcher:
|
||||||
подключён — best-effort, как touch/mark_health/release: проблема пула не должна
|
подключён — best-effort, как touch/mark_health/release: проблема пула не должна
|
||||||
ронять сбор. Lease НЕ освобождается и НЕ ротируется здесь — вызывающий код обычно
|
ронять сбор. Lease НЕ освобождается и НЕ ротируется здесь — вызывающий код обычно
|
||||||
сразу поднимает исключение и завершает сессию (release произойдёт как обычно в
|
сразу поднимает исключение и завершает сессию (release произойдёт как обычно в
|
||||||
`__aexit__`); пометка узла (`enabled=false`) переживает release — `acquire()`
|
`__aexit__`); бан переживает release — с #2600 п.2 это строка в
|
||||||
фильтрует по `enabled`, свежий lease его больше не возьмёт.
|
`scrape_proxy_source_bans` для пары (узел, `self._source`), и `acquire(source)`
|
||||||
|
её фильтрует, так что свежий lease ЭТОГО источника узел больше не возьмёт. Узел
|
||||||
|
при этом остаётся `enabled` и продолжает работать на другие источники: площадка
|
||||||
|
забанила IP, а не сломала прокси.
|
||||||
"""
|
"""
|
||||||
if self._lease is None or self._proxy_provider is None:
|
if self._lease is None or self._proxy_provider is None:
|
||||||
return
|
return
|
||||||
|
|
|
||||||
|
|
@ -211,14 +211,16 @@ class ProxyProvider(Protocol):
|
||||||
...
|
...
|
||||||
|
|
||||||
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
|
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
|
||||||
"""Пометить lease забаненным площадкой `source` (#2600 п.1).
|
"""Пометить lease забаненным площадкой `source` (#2600 п.1, п.2).
|
||||||
|
|
||||||
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
|
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
|
||||||
авто-disable'ит только после DISABLE_THRESHOLD подряд неудач (транзиентный
|
авто-disable'ит только после DISABLE_THRESHOLD подряд неудач (транзиентный
|
||||||
сбой должен пережить пару неудач). Здесь сигнал УЖЕ надёжно распознан (валидная
|
сбой должен пережить пару неудач). Здесь сигнал УЖЕ надёжно распознан (валидная
|
||||||
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) — узел выключается сразу
|
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) — узел немедленно снимается
|
||||||
(`disabled_reason='banned:<source>'`), кроме случая когда это последний живой
|
с выдачи ЭТОМУ источнику (строка в `scrape_proxy_source_bans`, срок эскалирует
|
||||||
узел для `source` (см. `app.services.proxy_pool.mark_banned` — там же защита).
|
на повторных банах), для остальных источников остаётся в строю: площадка банит
|
||||||
|
IP, а не ломает прокси. Исключение — последний узел, достижимый для `source`:
|
||||||
|
бан не записывается (см. `app.services.proxy_pool.mark_banned`, там же защита).
|
||||||
|
|
||||||
Вызывать из точки детекта бана, ПОКА lease ещё держится (до release/__aexit__) —
|
Вызывать из точки детекта бана, ПОКА lease ещё держится (до release/__aexit__) —
|
||||||
иначе id узла, который был использован, потерян. Best-effort — caller
|
иначе id узла, который был использован, потерян. Best-effort — caller
|
||||||
|
|
|
||||||
|
|
@ -19,8 +19,9 @@ release ВСЕГДА в finally — lease не должен течь, даже
|
||||||
url:`, — `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
|
url:`, — `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
|
||||||
`DomClickBlockedError` и т.п. — см. `proxy_errors.ProxyBanError`), это НЕ просто
|
`DomClickBlockedError` и т.п. — см. `proxy_errors.ProxyBanError`), это НЕ просто
|
||||||
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
|
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
|
||||||
`mark_banned` — узел выключается сразу (`disabled_reason='banned:<provider>'`), кроме
|
`mark_banned` — узел сразу снимается с выдачи ЭТОМУ провайдеру (per-source бан, #2600 п.2;
|
||||||
случая когда это последний живой узел (защита в `app.services.proxy_pool.mark_banned`).
|
для остальных источников остаётся в строю), кроме случая когда это последний узел,
|
||||||
|
достижимый для провайдера (защита в `app.services.proxy_pool.mark_banned`).
|
||||||
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
|
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
|
||||||
ИЗНУТРИ блока, получает сигнал бесплатно — этот модуль намеренно НЕ импортирует
|
ИЗНУТРИ блока, получает сигнал бесплатно — этот модуль намеренно НЕ импортирует
|
||||||
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные
|
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue