feat(tradein/proxy): здоровье прокси по паре «узел × источник» (#2600 п.2)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / 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 2m43s

Бан площадкой был глобальным: п.1 на распознанный бан выключал узел целиком
(enabled=false, disabled_reason='banned:<source>'). Реальность другая — Авито
банит IP, а Яндекс через тот же IP ходит чисто, поэтому один забаненный источник
выкидывал живой узел из пула для всех и худил пул быстрее, чем его пополняют
(#2638). Плюс такое состояние не самолечилось: ipify площадку не эмулирует, бан
не видит, а non-NULL disabled_reason блокирует авто-воскрешение (#2610) — нужен
был ручной PATCH.

Теперь бан — свойство ПАРЫ (proxy_id, source) в scrape_proxy_source_bans:
acquire(source) не выдаёт узел только этому источнику, для остальных узел
первосортный; снимается сам по времени. Срок эскалирует 6ч → 12 → 24 → 48 → 72
(потолок) на повторных банах той же пары; ban_count сбрасывается purge'ем
истёкших строк через 7 суток — поэтому purge намеренно отложенный, а не по
banned_until < now(). Защита последнего узла сохранена, но считается по
источнику: если после бана у acquire(source) не останется кандидатов — бан не
пишется, WARNING зовёт пополнять пул.

Миграция 210 конвертирует прод-остатки п.1 (enabled=false + disabled_reason
LIKE 'banned:%') в 6-часовые per-source баны и возвращает узлы в строй — иначе
они висели бы выключенными вечно.

Оператору активные баны видны в GET/PATCH /admin/proxies (source_bans) — без
этого «узел включён, но не выдаётся» необъяснимо.

Refs #2600
This commit is contained in:
bot-backend 2026-08-05 17:14:12 +05:00
parent d362b16d7c
commit 964867a943
7 changed files with 706 additions and 147 deletions

View file

@ -2638,6 +2638,52 @@ class ProxyBulkResponse(BaseModel):
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):
id: int
label: str | None
@ -2659,6 +2705,8 @@ class ProxyRow(BaseModel):
expires_at: str | None
created_at: str | None
updated_at: str | None
# #2600 п.2: активные баны площадками. Пустой список = узел выдаётся всем источникам.
source_bans: list[ProxySourceBan] = Field(default_factory=list)
@router.post("/proxies/bulk", response_model=ProxyBulkResponse)
@ -2744,6 +2792,9 @@ def list_proxies(
"""Список прокси со статусами. Пароли в url/rotate_url маскируются.
Фильтры: provider (=provider_affinity), enabled. Без фильтров все.
source_bans активные баны узла площадками (#2600 п.2): узел может быть
enabled=true и при этом не выдаваться конкретному источнику.
"""
clauses: list[str] = []
params: dict[str, Any] = {}
@ -2778,6 +2829,8 @@ def list_proxies(
def _iso(v: Any) -> str | None:
return v.isoformat() if v is not None else None
bans = _fetch_source_bans(db, [int(r["id"]) for r in rows])
return [
ProxyRow(
id=r["id"],
@ -2800,6 +2853,7 @@ def list_proxies(
expires_at=_iso(r["expires_at"]),
created_at=_iso(r["created_at"]),
updated_at=_iso(r["updated_at"]),
source_bans=bans.get(int(r["id"]), []),
)
for r in rows
]
@ -2900,6 +2954,7 @@ def patch_proxy(
expires_at=_iso(row["expires_at"]),
created_at=_iso(row["created_at"]),
updated_at=_iso(row["updated_at"]),
source_bans=_fetch_source_bans(db, [int(row["id"])]).get(int(row["id"]), []),
)

View file

@ -34,6 +34,25 @@ Self-healing (#2600):
enabled-узел выделенной affinity (пример domclick, один узел на всё, см. acquire
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).
Ручное выключение vs авто-выключение (#2610):
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
enabled=false: пул выключил сам после серии сбоев (disabled_reason IS NULL)
@ -75,6 +94,9 @@ __all__ = [
"DISABLE_THRESHOLD",
"MAX_CONSECUTIVE_FAILS",
"NON_RUN_LEASE_MARKER",
"SOURCE_BAN_BASE_HOURS",
"SOURCE_BAN_MAX_HOURS",
"SOURCE_BAN_PURGE_DAYS",
"STALE_LEASE_MINUTES",
"ProxyLease",
"acquire",
@ -109,6 +131,21 @@ DISABLED_RECHECK_MINUTES = 60
# Маркер lease для не-run вызовов (leased_by NOT NULL = занят, но это не id из scrape_runs).
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
# /scraper/health (_probe_current_ip).
_HEALTH_PROBE_URL = "https://api.ipify.org"
@ -156,6 +193,11 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
affinity='any' ИЛИ у этой affinity есть ДРУГОЙ enabled-узел (EXISTS-подзапрос)
т.е. выдача не обнулит доступность выделенной affinity целиком.
ОБА запроса отсекают узлы с АКТИВНЫМ баном по ЭТОМУ provider'у
(scrape_proxy_source_bans.banned_until > now(), #2600 п.2) — узел, забаненный Авито,
остаётся полноценным кандидатом для Яндекса и остальных источников. Бан по чужому
source на выдачу не влияет вообще.
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
@ -173,6 +215,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
AND consecutive_fails < CAST(:max_fails AS integer)
AND provider_affinity IN (:provider, 'any')
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
FOR UPDATE SKIP LOCKED
LIMIT 1
@ -199,6 +248,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
WHERE sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer)
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 (
sp.provider_affinity = 'any'
OR EXISTS (
@ -214,7 +270,7 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
LIMIT 1
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS},
{"max_fails": MAX_CONSECUTIVE_FAILS, "provider": provider},
)
.mappings()
.fetchone()
@ -407,92 +463,128 @@ def mark_health(
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 и
авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация —
транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно
распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер НЕ
исключение транспорта, НЕ голый network-fail) узел выключается НЕМЕДЛЕННО,
без ожидания порога: `enabled=false`, `disabled_reason='banned:<source>'`. Тот же
non-NULL `disabled_reason`, что и ручное выключение (#2610) — `mark_health(ok=True)`
больше не воскресит узел голым ipify-пробой: сам ipify площадку не эмулирует, бана
не увидит (ровно баг, который #2600 описывает как корень проблемы).
исключение транспорта, НЕ голый network-fail).
ЗАЩИТА ПОСЛЕДНЕГО ЖИВОГО УЗЛА (issue #2600 риск, переиспользует паттерн #2609):
если это последний узел, из-за которого `acquire(source)` вообще способен что-то
вернуть НЕ выключаем, только WARNING-лог. Доступность считается ТЕМ ЖЕ правилом,
что и acquire() (primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не
последняя из своей) EXISTS-проверка ниже, а не наивный `COUNT(*) WHERE enabled`.
Без этого защита не сработала бы, например, если проверять только
"affinity=proxy_id.provider_affinity" (fallback от #2609 позволяет чужой affinity
подменить исчерпанную) источник мог бы остаться без единого узла, даже когда
формально в пуле есть живые строки другой выделенной affinity (domclick).
ЧТО ИМЕННО ДЕЛАЕТСЯ (изменение против #2600 п.1): узел БОЛЬШЕ НЕ выключается
глобально (`enabled=false, disabled_reason='banned:<source>'` так было в п.1).
Пишется строка в `scrape_proxy_source_bans` (миграция 210): пока
`banned_until > now()`, `acquire(source)` этот узел не выдаёт, а для ЛЮБОГО
другого источника он остаётся первосортным. Авито банит IP Яндекс через тот же
IP ходит чисто; глобальное выключение выкидывало живой узел отовсюду и худило пул
в разы быстрее, чем его пополняют (#2638). `enabled`/`disabled_reason` остаются
исключительно за оператором (#2610) и за авто-disable'ом по транспортным сбоям.
КОНКУРЕНТНОСТЬ (deep-review fix 2): один `UPDATE ... WHERE ... AND EXISTS(...)`
НЕ атомарная гарантия поперёк СТРОК. EXISTS читает состояние других строк на момент
своего снапшота (READ COMMITTED), но не лочит их два ПАРАЛЛЕЛЬНЫХ mark_banned для
РАЗНЫХ proxy_id (напр. avito банит A, cian банит B миллисекундами позже) каждый может
увидеть другого как "ещё живого" в своём EXISTS и оба закоммититься пул уходит с 2
живых узлов в 0 разом. Это НЕ самолечится (#2610: disabled_reason блокирует ipify-
воскрешение, нужен ручной PATCH). Фикс: `pg_advisory_xact_lock` в начале транзакции
сериализует ВСЕ mark_banned-вызовы между собой (xact-scoped снимается сам на
commit/rollback, leak невозможен). Один глобальный ключ вместо per-affinity
сериализует и НЕ пересекающиеся по affinity баны тоже (avito vs domclick не гонятся
за одни строки), но частота вызовов низкая (несколько банов в час, не hot-path)
цена оправдана простотой против per-row `SELECT ... FOR UPDATE` по кандидатам
(потребовал бы лочить весь EXISTS-кандидат-сет заранее, выше риск deadlock между
параллельными mark_banned, лочащими пересекающиеся строки в разном порядке).
ponytail: global advisory lock, не per-affinity переходи на составной ключ
ЭСКАЛАЦИЯ: первый бан пары SOURCE_BAN_BASE_HOURS; каждый следующий удваивает
срок (ban_count после инкремента N SOURCE_BAN_BASE_HOURS * 2^(N-1)), но не выше
SOURCE_BAN_MAX_HOURS. Узел, который площадка банит раз за разом, отдыхает от неё
всё дольше, вместо того чтобы жечь прогоны. Сброс ban_count только purge'ем
истёкших строк (run_proxy_healthcheck, SOURCE_BAN_PURGE_DAYS).
ЗАЩИТА ПОСЛЕДНЕГО УЗЛА, ТЕПЕРЬ ПО ИСТОЧНИКУ (issue #2600 риск, паттерн #2609):
если после записи бана у `acquire(source)` не останется НИ ОДНОГО кандидата бан
НЕ пишется, только WARNING. Доступность считается ТЕМ ЖЕ правилом, что и acquire()
(primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не последняя из
своей) ПЛЮС отсутствие активной бан-строки для этого source EXISTS ниже, а не
наивный `COUNT(*) WHERE enabled`. Голодать без прокси хуже, чем ходить через
забаненный: капча хотя бы иногда пропускает, отсутствие узла нет.
КОНКУРЕНТНОСТЬ (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.
Идемпотентно: узел уже `enabled=false` (ранее забанен/выключен вручную) no-op,
INFO-лог, disabled_reason НЕ перезаписывается другим source (WHERE enabled в UPDATE).
Идемпотентно: повторный бан той же пары не создаёт дубль (PK (proxy_id, source))
продлевает срок по правилу эскалации. Несуществующий proxy_id no-op + WARNING.
Best-effort по контракту вызывающих (`BrowserFetcher.report_ban`, `curl_proxy_url`)
сюда попадают уже обёрнутыми в try/except, но сам мark_banned ошибки БД не глотает
сюда попадают уже обёрнутыми в try/except, но сам mark_banned ошибки БД не глотает
(падает как обычно) caller решает, ловить или нет.
"""
# Сериализует check+update ниже с другими конкурентными mark_banned (см. докстринг
# Сериализует check+insert ниже с другими конкурентными mark_banned (см. докстринг
# "КОНКУРЕНТНОСТЬ"). Держится до db.commit()/rollback() этой транзакции.
db.execute(
text("SELECT pg_advisory_xact_lock(CAST(:key AS bigint))"),
{"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 = (
db.execute(
text(
"""
UPDATE scrape_proxies
SET enabled = false,
disabled_reason = CAST(:reason AS text),
updated_at = now()
WHERE id = CAST(:proxy_id AS bigint)
AND enabled
INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason)
SELECT CAST(:proxy_id AS bigint),
CAST(:source AS text),
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)
)
AND EXISTS (
SELECT 1
FROM scrape_proxies sp
WHERE sp.id <> CAST(:proxy_id AS bigint)
AND sp.enabled
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 (
sp.provider_affinity IN (:source, 'any')
-- other.id <> sp.id (а не NOT IN (sp.id, :proxy_id), как в
-- п.1): банимый узел остаётся enabled и по-прежнему обслуживает
-- СВОЮ affinity значит он и есть валидный backup для неё.
OR EXISTS (
SELECT 1
FROM scrape_proxies other
WHERE other.provider_affinity = sp.provider_affinity
AND other.enabled
AND other.id NOT IN (sp.id, CAST(:proxy_id AS bigint))
AND other.id <> sp.id
)
)
)
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,
"source": source,
"reason": f"banned:{source}",
"base_hours": SOURCE_BAN_BASE_HOURS,
"max_hours": SOURCE_BAN_MAX_HOURS,
"max_fails": MAX_CONSECUTIVE_FAILS,
},
)
@ -502,16 +594,18 @@ def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
db.commit()
if row is not None:
logger.warning(
"proxy_pool: proxy id=%d BANNED by source=%s — disabled "
"(disabled_reason='banned:%s'), ipify-проба больше НЕ воскресит (#2610)",
"proxy_pool: proxy id=%d BANNED by source=%s — узел снят с выдачи ТОЛЬКО для "
"этого источника до %s (ban_count=%s); для остальных источников остаётся в "
"строю (#2600 п.2)",
proxy_id,
source,
source,
row["banned_until"],
row["ban_count"],
)
return
# 0 rows: либо уже disabled (идемпотентно, no-op), либо защита последнего узла
# сработала — читаем текущее состояние ТОЛЬКО для точного лога (не влияет на решение).
# 0 rows: либо узла нет, либо защита последнего узла отменила запись бана — читаем
# текущее состояние ТОЛЬКО для точного лога (на решение уже не влияет).
current = (
db.execute(
text(
@ -525,18 +619,11 @@ def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
)
if current is None:
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:
logger.warning(
"proxy_pool: proxy id=%d BANNED by source=%s but NOT disabled — last live node "
"reachable for this source (protection, mirrors acquire() fallback rule, #2600). "
"Ban logged only — operator should investigate/add capacity.",
"proxy_pool: proxy id=%d — бан не записан: это последний узел, достижимый для "
"source=%s; нужны новые прокси (см. #2638). Узел продолжит выдаваться этому "
"источнику (голодание хуже, чем работа через забаненный узел).",
proxy_id,
source,
)
@ -635,9 +722,12 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
следующего disable/enable цикла), но mark_health их не воскрешает revived не растёт,
WARNING пишет сам mark_health.
В конце purge бан-строк (#2600 п.2), истёкших дольше SOURCE_BAN_PURGE_DAYS назад
(см. комментарий у самого DELETE: отложенность это и есть сброс ban_count).
Пробы идут последовательно пул небольшой (десятки узлов), а параллельный залп на
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
{reaped, checked, ok, failed, revived}.
{reaped, checked, ok, failed, revived, bans_purged}.
"""
reaped = reap_stale_leases(db)
@ -687,13 +777,35 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
else:
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(
"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,
checked,
ok_count,
failed,
revived,
purged,
)
return {
"reaped": reaped,
@ -701,4 +813,5 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
"ok": ok_count,
"failed": failed,
"revived": revived,
"bans_purged": purged,
}

View file

@ -0,0 +1,98 @@
-- 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. Ручные выключения (disabled_reason без префикса
-- 'banned:') НЕ трогаются — это решение оператора.
--
-- 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
'Сколько раз эта пара банилась. Срок следующего бана = base * 2^(ban_count-1), '
'потолок SOURCE_BAN_MAX_HOURS. Сбрасывается только удалением строки purge''ем '
'через SOURCE_BAN_PURGE_DAYS после истечения бана.';
-- Горячий путь — 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) ─────────────────────
--
-- substring(... from 8) — отрезает префикс 'banned:' (7 символов). LIKE 'banned:_%'
-- (а не 'banned:%') гарантирует непустой source для NOT NULL-колонки; вырожденное
-- 'banned:' без источника строки бана не получает, но узел ниже всё равно вернётся
-- в строй — висеть вечно выключенным он не должен ни в каком случае.
INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason)
SELECT id,
substring(disabled_reason from 8),
now() + interval '6 hours',
'migrated from disabled_reason (210)'
FROM scrape_proxies
WHERE NOT enabled
AND disabled_reason LIKE 'banned:_%'
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 LIKE 'banned:%';
COMMIT;

View file

@ -26,6 +26,16 @@ reap_stale_leases проверяются по фактическому изме
* ручно-выключенный узел (disabled_reason НЕ NULL) НЕ воскресает даже при
ok=True, и это логируется (WARNING);
* 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 для эскалации).
"""
from __future__ import annotations
@ -44,6 +54,9 @@ from app.services.proxy_pool import (
DISABLE_THRESHOLD,
DISABLED_RECHECK_MINUTES,
MAX_CONSECUTIVE_FAILS,
SOURCE_BAN_BASE_HOURS,
SOURCE_BAN_MAX_HOURS,
SOURCE_BAN_PURGE_DAYS,
STALE_LEASE_MINUTES,
acquire,
mark_health,
@ -88,10 +101,13 @@ class _FakeResult:
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
# #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 — для теста
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
# проверить — фейксессия однопоточна).
@ -101,14 +117,30 @@ class FakeSession:
def _by_id(self, pid: int) -> dict[str, Any] | 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:
sql = str(stmt)
p = params or {}
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary affinity-scoped or fallback)
max_fails = p["max_fails"]
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'
provider = p["provider"]
cands = [
r
for r in self.rows
@ -116,6 +148,7 @@ class FakeSession:
and r["consecutive_fails"] < max_fails
and r["provider_affinity"] in (provider, "any")
and r["leased_by"] is None
and _not_banned(r)
]
else: # fallback: любая affinity, но не последний узел выделенной affinity
# (domclick и т.п. — #2600 review). ВАЖНО: применяем эту фильтрацию,
@ -141,6 +174,7 @@ class FakeSession:
if r["enabled"]
and r["consecutive_fails"] < max_fails
and r["leased_by"] is None
and _not_banned(r)
and (not protects_last_node or _has_backup(r))
]
# ORDER BY last_ok_at NULLS LAST, id
@ -233,34 +267,59 @@ class FakeSession:
self.advisory_lock_calls.append(p["key"])
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"]
source = p["source"]
max_fails = p["max_fails"]
row = self._by_id(proxy_id)
if row is None or not row["enabled"]:
return _FakeResult([]) # already disabled / not found — no-op
if self._by_id(proxy_id) is None:
return _FakeResult([]) # узла нет — no-op
def _is_candidate(sp: dict[str, Any]) -> bool:
if not (sp["enabled"] and sp["consecutive_fails"] < max_fails):
return False
if self._has_active_ban(sp["id"], source):
return False # уже забанен этим же источником — не кандидат
if sp["provider_affinity"] in (source, "any"):
return True
# fallback-safe: другой enabled узел ТОЙ ЖЕ affinity (кроме sp/proxy_id).
# fallback-safe: другой enabled узел ТОЙ ЖЕ affinity (банимый узел
# остаётся enabled и тоже считается — бан теперь per-source).
return any(
other["provider_affinity"] == sp["provider_affinity"]
and other["enabled"]
and other["id"] not in (sp["id"], proxy_id)
and other["id"] != sp["id"]
for other in self.rows
)
still_available = any(r["id"] != proxy_id and _is_candidate(r) for r in self.rows)
if not still_available:
return _FakeResult([]) # protected — последний живой узел, не выключаем
return _FakeResult([]) # protected — последний узел для source, бан не пишем
row["enabled"] = False
row["disabled_reason"] = p["reason"]
return _FakeResult([{"id": proxy_id}])
now = datetime.now(UTC)
ban = self._ban(proxy_id, source)
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: # 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
row = self._by_id(p["id"])
@ -704,68 +763,96 @@ async def test_healthcheck_checks_disabled_proxy_never_checked_before(
assert db._by_id(1)["enabled"] is False
# ── mark_banned (#2600 п.1 — довести сигнал бана до пула) ──────────────────────
# ── mark_banned (#2600 п.2 — бан по паре «узел × источник») ────────────────────
#
# Red/green контракт issue: (a) распознанный бан → узел выключен с
# disabled_reason='banned:<source>'; (b) ipify-проба его не воскрешает (уже
# покрыто disabled_reason-веткой mark_health выше, #2610 — здесь только
# убеждаемся, что mark_banned проставляет ТОТ ЖЕ non-NULL disabled_reason);
# (c) последний живой узел НЕ выключается (только лог); (d) сетевой сбой
# по-прежнему идёт через mark_health(ok=False), НЕ через mark_banned (проверяется
# на уровне browser_fetcher/curl_proxy_url тестов — здесь mark_banned сам по себе
# не участвует в различении причин, это забота вызывающего кода).
# Red/green контракт issue: (a) распознанный бан → строка в scrape_proxy_source_bans,
# узел НЕ выключен глобально (в п.1 было enabled=false — площадка забанила IP, а не
# сломала прокси; узел обязан остаться живым для остальных источников); (b) повторный
# бан той же пары эскалирует срок; (c) последний достижимый для source узел НЕ банится
# (только лог); (d) сетевой сбой по-прежнему идёт через mark_health(ok=False), НЕ через
# mark_banned (проверяется на уровне browser_fetcher/curl_proxy_url тестов — здесь
# mark_banned сам по себе не участвует в различении причин, это забота caller'а).
def test_mark_banned_disables_with_reason() -> None:
"""Забаненный узел выключается, disabled_reason='banned:<source>' (второй здоровый
узел affinity='any' есть защита последнего узла не срабатывает)."""
def test_mark_banned_writes_source_ban_row() -> None:
"""Бан пишется в scrape_proxy_source_bans: source, ban_count=1, срок = base-часы."""
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
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is False
assert row["disabled_reason"] == "banned:avito"
assert row["enabled"] is True
assert row["disabled_reason"] is None
def test_mark_banned_disabled_reason_blocks_auto_revive() -> None:
"""disabled_reason non-NULL после mark_banned → mark_health(ok=True) НЕ
воскрешает узел (та же ветка #2610, что и ручное выключение)."""
def test_mark_banned_repeat_escalates_ban_count_and_duration() -> None:
"""Повторный бан той же пары: ban_count растёт, срок удваивается (base * 2^(N-1))."""
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]
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]
row = db._by_id(1)
assert row["enabled"] is False
assert row["disabled_reason"] == "banned:cian" # не перезаписан
ban = db._ban(1, "avito")
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:
"""Единственный узел avito, других (свободных/'any') нет вообще — НЕ выключается,
только лог (issue #2600: бан не должен обрушить единственный источник целиком)."""
"""Единственный узел avito, других (свободных/'any') нет вообще — бан НЕ пишется,
только лог (issue #2600: бан не должен обрушить единственный источник целиком —
ходить через забаненный узел лучше, чем не ходить вообще)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is True # НЕ тронут
assert row["disabled_reason"] is None
assert db.bans == []
assert db._by_id(1)["enabled"] is True # и глобально не тронут
def test_mark_banned_protects_last_live_node_logs_warning(
@ -775,17 +862,38 @@ def test_mark_banned_protects_last_live_node_logs_warning(
db = FakeSession([_proxy(1, affinity="avito")])
with caplog.at_level("WARNING"):
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' закрывает
availability для avito (та же семантика, что acquire()'s primary IN (provider,
'any')) выключаем безопасно."""
'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]
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, свежий прод-факт) ──
@ -799,20 +907,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:
"""avito банится; в пуле остаётся только один domclick-узел (affinity выделенная,
БЕЗ backup) для avito это НЕ доступный узел (acquire('avito') не взял бы его через
fallback, #2609 protects last node of domclick). Защита должна сработать — avito-узел
НЕ выключается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
fallback, #2609 protects last node of domclick). Защита должна сработать — бан
НЕ записывается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="domclick")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is True # domclick-узел НЕ считается доступной заменой
assert db.bans == [] # domclick-узел НЕ считается доступной заменой
assert db._by_id(2)["enabled"] is True # и сам не тронут
def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None:
"""Та же ситуация, но у domclick есть ВТОРОЙ узел (backup) — тогда fallback может
забрать ОДИН из них под avito (acquire()'s EXISTS-правило #2609), доступность для
avito сохраняется через fallback banned avito-узел безопасно выключается."""
avito сохраняется через fallback бан avito-узла записывается безопасно."""
assert mark_banned is not None
db = FakeSession(
[
@ -822,7 +929,7 @@ def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None
]
)
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_unhealthy_candidate_not_counted_as_backup() -> None:
@ -833,7 +940,7 @@ def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
[_proxy(1, affinity="avito"), _proxy(2, affinity="any", fails=MAX_CONSECUTIVE_FAILS)]
)
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) ───────────────────────
@ -857,10 +964,148 @@ def test_mark_banned_takes_advisory_xact_lock_with_fixed_key() -> None:
def test_mark_banned_takes_advisory_lock_even_when_protected() -> None:
"""Лок берётся ПЕРЕД проверкой доступности — даже когда защита последнего узла
в итоге отменяет disable, лок всё равно взят (сериализация check+decide, не
только update)."""
в итоге отменяет запись бана, лок всё равно взят (сериализация check+decide, не
только сам INSERT)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito")]) # единственный узел — protected
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db.advisory_lock_calls # лок взят, хотя disable не произошёл
assert db._by_id(1)["enabled"] is True
assert db.advisory_lock_calls # лок взят, хотя бан не записан
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"]

View file

@ -2,7 +2,9 @@
Покрытие (db мокается, NO live network/DB):
- 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
- #2610: PATCH enabled=false ставит disabled_reason (ручное выключение отличимо от
авто); PATCH enabled=true сбрасывает disabled_reason в NULL (снова авто-восстанавливаем)
@ -49,6 +51,17 @@ def _scalar_result(value: object) -> MagicMock:
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 ─────────────────────────────────────────────────────
@ -156,7 +169,7 @@ def _proxy_db_row(**over: Any) -> dict[str, Any]:
def test_list_masks_password(client: TestClient, db: MagicMock) -> None:
result = MagicMock()
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")
assert r.status_code == 200, r.text
@ -174,7 +187,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=2, enabled=False, disabled_reason="забанен Авито"),
]
db.execute.return_value = result
db.execute.side_effect = [result, _bans_result()]
r = client.get("/api/v1/admin/proxies")
assert r.status_code == 200, r.text
@ -183,6 +196,37 @@ def test_list_exposes_disabled_reason(client: TestClient, db: MagicMock) -> None
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:
result = MagicMock()
result.mappings.return_value.all.return_value = []
@ -201,7 +245,7 @@ def test_list_passes_filters(client: TestClient, db: MagicMock) -> None:
def test_patch_disable(client: TestClient, db: MagicMock) -> None:
result = MagicMock()
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})
assert r.status_code == 200, r.text
@ -227,13 +271,14 @@ def test_patch_disable_sets_disabled_reason_default(client: TestClient, db: Magi
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
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})
assert r.status_code == 200, r.text
assert r.json()["disabled_reason"] == "manually disabled via admin API"
# дефолтная причина реально передана в SQL как fallback-параметр
params = db.execute.call_args.args[1]
# дефолтная причина реально передана в SQL как fallback-параметр (первый execute —
# сам UPDATE; второй, #2600 п.2, дочитывает активные баны по источникам)
params = db.execute.call_args_list[0].args[1]
assert params["default_reason"]
assert params["reason"] is None
@ -244,12 +289,12 @@ def test_patch_disable_sets_disabled_reason_custom(client: TestClient, db: Magic
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
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": "забанен Авито"})
assert r.status_code == 200, r.text
assert r.json()["disabled_reason"] == "забанен Авито"
params = db.execute.call_args.args[1]
params = db.execute.call_args_list[0].args[1]
assert params["reason"] == "забанен Авито"
@ -260,12 +305,12 @@ def test_patch_enable_clears_disabled_reason(client: TestClient, db: MagicMock)
result.mappings.return_value.fetchone.return_value = _proxy_db_row(
enabled=True, disabled_reason=None
)
db.execute.return_value = result
db.execute.side_effect = [result, _bans_result()]
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
assert r.status_code == 200, r.text
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

View file

@ -211,14 +211,16 @@ class ProxyProvider(Protocol):
...
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
"""Пометить lease забаненным площадкой `source` (#2600 п.1).
"""Пометить lease забаненным площадкой `source` (#2600 п.1, п.2).
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
авто-disable'ит только после DISABLE_THRESHOLD подряд неудач (транзиентный
сбой должен пережить пару неудач). Здесь сигнал УЖЕ надёжно распознан (валидная
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) узел выключается сразу
(`disabled_reason='banned:<source>'`), кроме случая когда это последний живой
узел для `source` (см. `app.services.proxy_pool.mark_banned` там же защита).
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) узел немедленно снимается
с выдачи ЭТОМУ источнику (строка в `scrape_proxy_source_bans`, срок эскалирует
на повторных банах), для остальных источников остаётся в строю: площадка банит
IP, а не ломает прокси. Исключение последний узел, достижимый для `source`:
бан не записывается (см. `app.services.proxy_pool.mark_banned`, там же защита).
Вызывать из точки детекта бана, ПОКА lease ещё держится (до release/__aexit__)
иначе id узла, который был использован, потерян. Best-effort caller

View file

@ -19,8 +19,9 @@ release ВСЕГДА в finally — lease не должен течь, даже
url:`, `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
`DomClickBlockedError` и т.п. см. `proxy_errors.ProxyBanError`), это НЕ просто
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
`mark_banned` узел выключается сразу (`disabled_reason='banned:<provider>'`), кроме
случая когда это последний живой узел (защита в `app.services.proxy_pool.mark_banned`).
`mark_banned` узел сразу снимается с выдачи ЭТОМУ провайдеру (per-source бан, #2600 п.2;
для остальных источников остаётся в строю), кроме случая когда это последний узел,
достижимый для провайдера (защита в `app.services.proxy_pool.mark_banned`).
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
ИЗНУТРИ блока, получает сигнал бесплатно этот модуль намеренно НЕ импортирует
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные