From 964867a9433f1e6612a54100e85b0a189aceedf8 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 5 Aug 2026 17:14:12 +0500 Subject: [PATCH] =?UTF-8?q?feat(tradein/proxy):=20=D0=B7=D0=B4=D0=BE=D1=80?= =?UTF-8?q?=D0=BE=D0=B2=D1=8C=D0=B5=20=D0=BF=D1=80=D0=BE=D0=BA=D1=81=D0=B8?= =?UTF-8?q?=20=D0=BF=D0=BE=20=D0=BF=D0=B0=D1=80=D0=B5=20=C2=AB=D1=83=D0=B7?= =?UTF-8?q?=D0=B5=D0=BB=20=C3=97=20=D0=B8=D1=81=D1=82=D0=BE=D1=87=D0=BD?= =?UTF-8?q?=D0=B8=D0=BA=C2=BB=20(#2600=20=D0=BF.2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Бан площадкой был глобальным: п.1 на распознанный бан выключал узел целиком (enabled=false, disabled_reason='banned:'). Реальность другая — Авито банит 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 --- tradein-mvp/backend/app/api/v1/admin.py | 55 +++ .../backend/app/services/proxy_pool.py | 233 ++++++++--- .../data/sql/210_scrape_proxy_source_bans.sql | 98 +++++ .../backend/tests/services/test_proxy_pool.py | 385 ++++++++++++++---- .../backend/tests/test_admin_proxies.py | 67 ++- .../scraper-kit/src/scraper_kit/contracts.py | 10 +- .../src/scraper_kit/providers/_proxy.py | 5 +- 7 files changed, 706 insertions(+), 147 deletions(-) create mode 100644 tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql diff --git a/tradein-mvp/backend/app/api/v1/admin.py b/tradein-mvp/backend/app/api/v1/admin.py index 1e5add94..9eedea04 100644 --- a/tradein-mvp/backend/app/api/v1/admin.py +++ b/tradein-mvp/backend/app/api/v1/admin.py @@ -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"]), []), ) diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index e614a623..fb297169 100644 --- a/tradein-mvp/backend/app/services/proxy_pool.py +++ b/tradein-mvp/backend/app/services/proxy_pool.py @@ -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:'`. Тот же - 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:'` — так было в п.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, } diff --git a/tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql b/tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql new file mode 100644 index 00000000..65b756e8 --- /dev/null +++ b/tradein-mvp/backend/data/sql/210_scrape_proxy_source_bans.sql @@ -0,0 +1,98 @@ +-- 210_scrape_proxy_source_bans.sql +-- Здоровье прокси по ПАРЕ «узел × источник» (#2600 п.2). +-- +-- WHY: +-- До сих пор бан был ГЛОБАЛЬНЫМ: #2600 п.1 (mark_banned) на распознанный бан +-- площадкой выключал узел целиком — enabled=false, disabled_reason='banned:'. +-- Реальность другая: Авито банит 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; diff --git a/tradein-mvp/backend/tests/services/test_proxy_pool.py b/tradein-mvp/backend/tests/services/test_proxy_pool.py index 3fd7597d..1ee3a85b 100644 --- a/tradein-mvp/backend/tests/services/test_proxy_pool.py +++ b/tradein-mvp/backend/tests/services/test_proxy_pool.py @@ -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:'; (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:' (второй здоровый - узел 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"] diff --git a/tradein-mvp/backend/tests/test_admin_proxies.py b/tradein-mvp/backend/tests/test_admin_proxies.py index f4619380..789953c1 100644 --- a/tradein-mvp/backend/tests/test_admin_proxies.py +++ b/tradein-mvp/backend/tests/test_admin_proxies.py @@ -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 diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py index 8a8883b4..56324055 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py @@ -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` (см. `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 diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py index 83912475..baba2926 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py @@ -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:'`), кроме -случая когда это последний живой узел (защита в `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-прокси-слой не должен знать про конкретные