From b89788ee99cdcf3ad17e8f1a332179fb6dc8d4ec Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 13:05:20 +0300 Subject: [PATCH] =?UTF-8?q?feat(tradein/proxy):=20=D0=BF=D1=80=D0=BE=D0=B3?= =?UTF-8?q?=D0=BE=D0=BD=20=D0=B7=D0=BD=D0=B0=D0=B5=D1=82=20=D1=81=D0=B2?= =?UTF-8?q?=D0=BE=D0=B9=20=D1=83=D0=B7=D0=B5=D0=BB,=20=D0=B0=20=D1=81?= =?UTF-8?q?=D0=BD=D1=8F=D1=82=D1=8B=D0=B9=20=D0=B1=D0=B0=D0=BD=20=D0=BF?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D1=81=D1=82=D0=B0=D1=91=D1=82=20=D1=81=D1=82?= =?UTF-8?q?=D0=B8=D1=80=D0=B0=D1=82=D1=8C=20=D0=B8=D1=81=D1=82=D0=BE=D1=80?= =?UTF-8?q?=D0=B8=D1=8E=20(#3404)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Выбор оператора мобильного прокси опирался на две ненадёжные опоры. Первая: `scrape_runs` не знала, через какой узел шёл прогон — колонка `proxy_id` была только у банов и ротаций. «Какой узел собрал 5 карточек из 21» не выяснялось ни одним запросом. Вторая: `clear_source_bans` делала DELETE, а зовётся она после КАЖДОЙ успешной ротации exit-IP. У #540723 (МегаФон) 23 успешные ротации и ноль строк банов, у #540722 (Tele2) ротаций почти не было и 7 банов. «7 против 0» читалось как «Tele2 хуже», хотя в той же мере это «у МегаФона историю стёрли 23 раза». Теперь: - `scrape_runs.proxy_id` — последний выданный прогону узел; полная цепочка (если узел менялся mid-run) копится в `counters.proxy_ids`. Пишет `proxy_pool.attribute_run_proxy` из единственной точки — сразу после выдачи лиза в `acquire()`, поэтому curl-путь, браузерный sticky lease и ре-acquire при ротации покрыты одинаково. `run_id` доходит до адаптера через ContextVar (`scraper_kit.orchestration.run_context`): протокол `ProxyProvider.acquire` его не несёт, а `RealProxyProvider` живёт одним объектом на весь планировщик. Best-effort: `lock_timeout` 2с и проглоченное исключение — диагностика не вправе ронять выдачу прокси или ждать на блокировке строки прогона. - `clear_source_bans` гасит строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`) вместо удаления. Эскалация сохраняется 1:1: формула в `mark_banned` берёт ПРЕДЫДУЩИЙ `ban_count` показателем степени, при нуле это ровно `SOURCE_BAN_BASE_HOURS` — как после DELETE. Строка доживает до штатного purge по `SOURCE_BAN_PURGE_DAYS`. Для всех читателей `scrape_proxy_source_bans` погашенная строка неотличима от отсутствующей: acquire, оба guard-подзапроса `mark_banned`, `proxy_egress` (ранжирование по `ban_count` даёт 0, как у узла без истории), admin `_active_ban` — все гейтятся по `banned_until > now()`. Ничего не бэкфиллится: связать прошедшие прогоны с узлами нечем (`leased_by` исторически = NON_RUN_LEASE_MARKER), врать восстановленным значением нельзя. Миграция 287. Тесты: 9 новых на обе части (главный — эскалация после гашения даёт базовые 6ч, а не удвоенные) + 14 существующих переведены с DELETE-семантики на гашение, включая проверку, что секрет ротации не утекает в новое `cleared_reason`. Полный прогон бэкенда: 5600 passed, 37 skipped. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_011WHFxVPWoBnSZihkdH1Uou --- tradein-mvp/backend/app/core/config.py | 3 +- .../backend/app/services/proxy_pool.py | 137 +++++- .../backend/app/services/scraper_adapters.py | 9 +- .../data/sql/287_proxy_run_attribution.sql | 107 +++++ .../backend/tests/services/test_proxy_pool.py | 86 +++- .../tests/services/test_proxy_rotation.py | 53 ++- .../tests/test_2800_per_source_probe.py | 9 +- .../tests/test_3404_proxy_run_attribution.py | 403 ++++++++++++++++++ .../backend/tests/test_admin_proxies.py | 14 +- .../scraper_kit/orchestration/run_context.py | 35 ++ .../src/scraper_kit/orchestration/runs.py | 22 +- 11 files changed, 832 insertions(+), 46 deletions(-) create mode 100644 tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql create mode 100644 tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py create mode 100644 tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/run_context.py diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 44109609..748a70d9 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -1187,7 +1187,8 @@ class Settings(BaseSettings): # #3283g: ротация exit-IP НА САМ БАН площадки, а не только по счётчику попыток. # Бан привязан к IP (замерено вживую: rotate_proxy() лечит забаненный узел за - # секунды, clear_source_bans снимает запись из scrape_proxy_source_bans), но + # секунды, clear_source_bans гасит бан в scrape_proxy_source_bans — строка живёт + # до purge, #3404), но # #3251/#3212 запрещают сбрасывать browser-context на КАЖДЫЙ блок -- сброс без # смены IP выбрасывает пройденный QRATOR-PoW и запускает самоподдерживающийся # каскад блоков на том же адресе. rotate_on_ban МЕНЯЕТ IP вместе со сбросом, diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index 9f52e565..dcdd4bb1 100644 --- a/tradein-mvp/backend/app/services/proxy_pool.py +++ b/tradein-mvp/backend/app/services/proxy_pool.py @@ -54,7 +54,10 @@ Self-healing (#2600): WARNING (пул надо пополнять, #2638). - Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился - IP, поэтому смена адреса делает строку недействительной. + IP, поэтому смена адреса делает строку недействительной. С #3404 «снятие» гасит + строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`), + а не удаляет её — строка живёт до штатного purge (SOURCE_BAN_PURGE_DAYS), но для + выдачи и для эскалации следующего бана это неотличимо от прежнего DELETE. Ручное выключение vs авто-выключение (#2610): - scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины @@ -143,6 +146,7 @@ __all__ = [ "STALE_LEASE_MINUTES", "ProxyLease", "acquire", + "attribute_run_proxy", "clear_source_bans", "mark_banned", "mark_browser_health", @@ -441,6 +445,11 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe {"run_id": lease_marker, "id": proxy_id}, ) db.commit() + if run_id is not None and run_id != NON_RUN_LEASE_MARKER: + # #3404: одна точка, покрывающая ВСЕ пути выдачи (curl — acquire на каждый + # вызов, браузер — sticky lease на весь прогон, ре-acquire при ротации узла + # mid-run) — см. attribute_run_proxy docstring. + attribute_run_proxy(db, run_id, proxy_id) if fallback_used: logger.warning( "proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity " @@ -474,6 +483,78 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe ) +def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None: + """Записать узел, через который идёт прогон run_id, в scrape_runs (#3404). + + Единственный писатель — `acquire()` сразу после выдачи lease'а: покрывает и + curl-путь (acquire на каждый вызов), и браузерный sticky lease (один acquire на + весь прогон), и ре-acquire при ротации узла mid-run (`browser_fetcher`, + `_LEASE_ROTATE_AFTER_FAILS`) — то есть смена узла ЗА прогон фиксируется сама, + без отдельного вызова с чьей-либо стороны. + + `scrape_runs.proxy_id` — ПОСЛЕДНИЙ использованный узел (перезаписывается при + каждой новой выдаче); полная цепочка узлов, если она менялась, — в + `counters.proxy_ids` (список id, без дублей). Пишем через `||`-мерж + `counters` (тот же контракт, что у `runs.update_heartbeat`/`mark_done`) — + чужие ключи (чекпоинт, метка interrupted) не затираются. + + Идемпотентно: повторная выдача ТОГО ЖЕ узла не дублирует его в `proxy_ids` + (`@>`-проверка перед append). Best-effort: любой сбой (например, run_id уже + не существует — гонка с финализацией) логируется WARNING и проглатывается — + атрибуция прогону не должна ронять выдачу прокси, это диагностика, а не + часть контракта lease'а. 0 rows (run_id не найден) — DEBUG, не ошибка: сама + выдача при этом уже произошла и коммитнута предыдущим db.commit() в acquire(). + """ + try: + # Строку прогона параллельно обновляет heartbeat/финализатор из ДРУГОЙ сессии + # (короткие транзакции, каждая со своим commit). Пересечение маловероятно, но + # ждать на блокировке в пути выдачи прокси нельзя — диагностика не должна + # тормозить сбор. Не дождались за 2с — уходим в except ниже (WARNING, lease цел). + db.execute(text("SET LOCAL lock_timeout = '2s'")) + row = db.execute( + text( + """ + UPDATE scrape_runs + SET proxy_id = CAST(:proxy_id AS bigint), + counters = COALESCE(counters, '{}'::jsonb) || jsonb_build_object( + 'proxy_ids', + CASE + WHEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb) + @> to_jsonb(CAST(:proxy_id AS bigint)) + THEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb) + ELSE COALESCE(counters -> 'proxy_ids', '[]'::jsonb) + || jsonb_build_array(CAST(:proxy_id AS bigint)) + END + ) + WHERE id = CAST(:run_id AS bigint) + RETURNING id + """ + ), + {"proxy_id": proxy_id, "run_id": run_id}, + ).first() + db.commit() + if row is None: + logger.debug( + "proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already " + "finalized?)", + run_id, + ) + except Exception: + # Best-effort (см. docstring) — атрибуция диагностическая, не часть + # контракта lease'а: lease уже выдан и не должен теряться из-за неё. + logger.warning( + "proxy_pool: attribute_run_proxy failed run_id=%d proxy_id=%d — lease " + "issued regardless", + run_id, + proxy_id, + exc_info=True, + ) + try: + db.rollback() + except Exception: + pass + + def release(db: Session, proxy_id: int) -> None: """Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен).""" db.execute( @@ -935,15 +1016,21 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None = ) ) ON CONFLICT (proxy_id, source) DO UPDATE - SET ban_count = scrape_proxy_source_bans.ban_count + 1, - banned_until = now() + make_interval(hours => CAST( + 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() + reason = CAST(:reason AS text), + -- #3404: строка могла быть погашена clear_source_bans (banned_until + -- в прошлом попадает в WHERE ниже) — новый бан затирает её метки + -- гашения, иначе на СНОВА забаненной паре висели бы cleared_at/ + -- cleared_reason от предыдущего, уже неактуального гашения. + cleared_at = NULL, + cleared_reason = NULL, + updated_at = now() -- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет), -- свою же (обычная эскалация) или перебиваем боевым сбором — он сильнее -- пробы. Иначе 0 rows и ветка "deferred" ниже (дефект #2803). @@ -1058,12 +1145,25 @@ def clear_source_bans( банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса, держа узел вне выдачи уже без причины. - source=None — снять все баны узла; конкретный source — только его. DELETE, а не - `banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации), - а здесь мы как раз объявляем историю недействительной — новый бан начнётся с базовых - SOURCE_BAN_BASE_HOURS. + source=None — снять все баны узла; конкретный source — только его. Гасим строку + (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`), а НЕ + удаляем (#3404, было DELETE): строка доживает до штатного purge'а + (`run_proxy_healthcheck`, SOURCE_BAN_PURGE_DAYS), но перестаёт блокировать + выдачу немедленно и перестаёт нести историю эскалации — `ban_count = 0` даёт + следующему бану той же пары ровно те же SOURCE_BAN_BASE_HOURS, что и раньше + после DELETE (формула `mark_banned` берёт ПРЕДЫДУЩИЙ ban_count показателем + степени: 0 → база, без множителя). Причина держать строку — трассируемость + (видно, что бан БЫЛ и когда/кем снят), а не поведение: для читателей ниже + погашенная строка неотличима от отсутствующей (см. риски в шапке PR #3404). - `reason` идёт только в лог (человекочитаемый повод — «manual enable», «ip rotated»). + Идемпотентно и в другую сторону: повторный вызов на уже погашенной строке + (последний предикат в WHERE) её не трогает — 0 rows, `banned_until` НЕ + сдвигается вперёд. Без этого условия повторный `PATCH enabled=true` двигал бы + `banned_until` на каждый вызов и отодвигал бы purge на неопределённый срок. + + `reason` идёт в лог (человекочитаемый повод — «manual enable», «ip rotated») и + теперь ЕЩЁ в колонку `cleared_reason` — постоянный след того, кто и почему + погасил бан. `only_reason` — ФИЛЬТР по колонке reason, т.е. «снимать только строки, которые написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое @@ -1075,14 +1175,22 @@ def clear_source_bans( rows = db.execute( text( """ - DELETE FROM scrape_proxy_source_bans + UPDATE scrape_proxy_source_bans + SET banned_until = now(), + ban_count = 0, + cleared_at = now(), + cleared_reason = CAST(:reason AS text), + updated_at = now() WHERE proxy_id = CAST(:proxy_id AS bigint) AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text)) AND (CAST(:only_reason AS text) IS NULL OR reason = CAST(:only_reason AS text)) + -- Уже погашенная строка (гейт покоя, см. докстринг) — не трогаем: без + -- него повторный вызов сдвигал бы banned_until вперёд и отодвигал purge. + AND NOT (cleared_at IS NOT NULL AND ban_count = 0 AND banned_until <= now()) RETURNING source """ ), - {"proxy_id": proxy_id, "source": source, "only_reason": only_reason}, + {"proxy_id": proxy_id, "source": source, "only_reason": only_reason, "reason": reason}, ).fetchall() db.commit() if rows: @@ -1465,6 +1573,11 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]: # Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с # 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса: # неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать». + # #3404: под этот же порог теперь попадают и ПОГАШЕННЫЕ clear_source_bans строки — + # для них banned_until == момент гашения (== cleared_at), т.е. таймер до purge + # отсчитывается от гашения, а не от исходного истечения бана. ban_count у них уже + # 0 к моменту гашения, так что покидающий purge их не «сбрасывает» повторно — + # он просто убирает уже неактуальный след из таблицы. purged = len( db.execute( text( diff --git a/tradein-mvp/backend/app/services/scraper_adapters.py b/tradein-mvp/backend/app/services/scraper_adapters.py index 44bb627c..e712d0d6 100644 --- a/tradein-mvp/backend/app/services/scraper_adapters.py +++ b/tradein-mvp/backend/app/services/scraper_adapters.py @@ -221,10 +221,17 @@ class RealProxyProvider: def acquire(self, provider: str) -> ProxyLease | None: from scraper_kit.contracts import ProxyLease as _KitProxyLease + from scraper_kit.orchestration.run_context import current_run_id + # #3404: протокол ProxyProvider.acquire(provider) не несёт run_id (этот + # адаптер — один объект на весь scheduler_main.py), поэтому берём его из + # ContextVar, который выставляет runs.create_run. None — вызов вне прогона + # (health-check, эстиматор) — proxy_pool.acquire в этом случае лизит под + # NON_RUN_LEASE_MARKER, как и раньше, атрибуцию в scrape_runs не пишет. + run_id = current_run_id.get() db = _SessionLocal() try: - lease = _proxy_pool.acquire(db, provider) + lease = _proxy_pool.acquire(db, provider, run_id=run_id) finally: db.close() if lease is None: diff --git a/tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql b/tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql new file mode 100644 index 00000000..20c72e8c --- /dev/null +++ b/tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql @@ -0,0 +1,107 @@ +-- 287_proxy_run_attribution.sql +-- scrape_runs.proxy_id — узел прогона (#3404 A) + soft-clear банов по источнику (#3404 B). +-- +-- Dependencies: 015_scrape_runs.sql (scrape_runs), 157_scrape_proxies.sql (scrape_proxies), +-- 210_scrape_proxy_source_bans.sql (scrape_proxy_source_bans). +-- Apply after: 286_offer_price_history_decimal_slips.sql +-- +-- ЧАСТЬ A — WHY: +-- До сих пор ни одна строка scrape_runs не знала, через какой узел пула шёл прогон: +-- ProxyProvider.acquire(provider) run_id не принимает, а leased_by у боевого пути — +-- NON_RUN_LEASE_MARKER. Разбор исхода прогона по узлу (кто плодит баны/провалы) +-- был возможен только вручную, по времени. Код-часть (app/services/proxy_pool.py: +-- attribute_run_proxy) пишет сюда после каждой выдачи лиза; здесь только схема. +-- +-- ЧАСТЬ A — WHAT: +-- proxy_id — узел, через который шёл прогон. Если за прогон узел МЕНЯЛСЯ (ротация +-- при повторных провалах в browser_fetcher, либо curl-путь берёт лиз на каждый вызов +-- в providers/_proxy.py), здесь остаётся ПОСЛЕДНИЙ; полная цепочка узлов копится в +-- scrape_runs.counters->'proxy_ids' (jsonb-массив, пишет тот же attribute_run_proxy). +-- ON DELETE SET NULL, а не CASCADE — узел из пула может быть выведен/удалён оператором, +-- история прогонов (аналитика, отчёты) не должна пропадать вместе с ним. +-- Индекс (proxy_id, started_at DESC) — под разрез «исход прогона по узлу за период» +-- (WHERE proxy_id = ... ORDER BY started_at DESC); partial по proxy_id IS NOT NULL не +-- делаем, потому что колонка сортировки (started_at) в самом индексе — Postgres и так +-- не будет использовать индекс без него для прогонов без узла. +-- +-- НИЧЕГО НЕ БЭКФИЛЛИТСЯ: связать уже прошедшие прогоны с конкретным узлом задним +-- числом нечем — leased_by исторически = NON_RUN_LEASE_MARKER, а лог выдачи лизов +-- не хранит run_id. Для всех строк scrape_runs, созданных ДО этой миграции, +-- proxy_id остаётся NULL навсегда — это не «прогон без прокси», а «прогон, для +-- которого атрибуция не собиралась». Врать восстановленным/угаданным значением +-- нельзя, поэтому backfill-UPDATE здесь сознательно отсутствует. +-- +-- ЧАСТЬ B — WHY: +-- clear_source_bans() (proxy_pool.py) сейчас делает DELETE строки +-- scrape_proxy_source_bans. Это стирает историю эскалации (ban_count) и не оставляет +-- следа, что бан был снят ДОСРОЧНО (оператором/успешной ротацией exit-IP), в отличие +-- от бана, который просто истёк сам. Код-часть переводит функцию на UPDATE +-- (гашение: banned_until=now(), ban_count=0, cleared_at/cleared_reason проставляются), +-- строка доживает до штатного purge в run_proxy_healthcheck (SOURCE_BAN_PURGE_DAYS). +-- Эскалация при повторном бане той же пары сохраняется 1:1: формула в mark_banned +-- берёт СТАРЫЙ ban_count как показатель степени (base * 2^ban_count), при +-- ban_count=0 после гашения это ровно SOURCE_BAN_BASE_HOURS=6ч — байт-в-байт как +-- свежий INSERT после DELETE. +-- +-- ЧАСТЬ B — WHAT: +-- cleared_at — момент досрочного снятия бана (не путать с истечением banned_until +-- само по себе: NULL значит «бан снят не был / истёк сам», не-NULL — снят +-- оператором или ротацией exit-IP до истечения срока или сразу после). +-- cleared_reason — свободный текст причины снятия (тот же 'reason', что передаётся +-- в clear_source_bans). +-- +-- ИДЕМПОТЕНТНОСТЬ: +-- ADD COLUMN IF NOT EXISTS × 3, CREATE INDEX IF NOT EXISTS, COMMENT ON COLUMN +-- (безусловны, но идемпотентны сами по себе — просто перезаписывают тот же текст). +-- Backfill-DML в файле нет вовсе, поэтому повторный прогон — чистый no-op. + +BEGIN; + +SET LOCAL lock_timeout = '5s'; + +-- ── Часть A: scrape_runs.proxy_id ──────────────────────────────────────────── + +ALTER TABLE scrape_runs + ADD COLUMN IF NOT EXISTS proxy_id bigint REFERENCES scrape_proxies(id) ON DELETE SET NULL; + +CREATE INDEX IF NOT EXISTS idx_scrape_runs_proxy_id_started_at + ON scrape_runs (proxy_id, started_at DESC); + +COMMENT ON COLUMN scrape_runs.proxy_id IS + 'Узел пула (scrape_proxies.id), через который шёл прогон. NULL = прогон без ' + 'прокси (эстиматор, admin-инициированные вызовы с NON_RUN_LEASE_MARKER) либо ' + 'прогон ДО применения миграции 287 (backfill не делался — связать нечем). ' + 'Если узел менялся mid-run (ротация после серии провалов в browser_fetcher, ' + 'либо curl-путь берёт лиз заново на каждый вызов) — здесь ПОСЛЕДНИЙ выданный ' + 'узел, полная цепочка — counters->''proxy_ids'' (jsonb-массив id, в порядке ' + 'первой выдачи). ON DELETE SET NULL: удаление узла из пула не должно уносить ' + 'историю прогонов.'; + +-- ── Часть B: soft-clear в scrape_proxy_source_bans ─────────────────────────── + +ALTER TABLE scrape_proxy_source_bans + ADD COLUMN IF NOT EXISTS cleared_at timestamptz, + ADD COLUMN IF NOT EXISTS cleared_reason text; + +COMMENT ON COLUMN scrape_proxy_source_bans.cleared_at IS + 'Момент досрочного снятия бана (proxy_pool.clear_source_bans, #3404) — ' + 'оператором или успешной ротацией exit-IP. NULL = бан не снимался вручную ' + '(либо ещё активен, либо истёк сам по banned_until). Строка при гашении НЕ ' + 'удаляется — доживает до штатного purge (SOURCE_BAN_PURGE_DAYS), таймер ' + 'которого для погашенных строк отсчитывается от banned_until = момент гашения.'; + +COMMENT ON COLUMN scrape_proxy_source_bans.cleared_reason IS + 'Причина досрочного снятия бана (тот же текст, что передан в ' + 'clear_source_bans(reason=...)). NULL, если строка не гасилась вручную.'; + +COMMENT ON COLUMN scrape_proxy_source_bans.ban_count IS + 'Сколько раз эта пара банилась. Срок ТЕКУЩЕГО бана (banned_until - banned_at) = ' + 'base * 2^(ban_count-1), потолок SOURCE_BAN_MAX_HOURS: ban_count=1 → 6ч, 2 → 12ч, ' + '3 → 24ч и т.д. Сбрасывается либо purge''ем через SOURCE_BAN_PURGE_DAYS после ' + 'истечения, либо proxy_pool.clear_source_bans (#3404: досрочное ГАШЕНИЕ строки —' + ' banned_until=now(), ban_count=0, cleared_at/cleared_reason проставляются; ' + 'строка НЕ удаляется, живёт до purge). Оба пути одинаково обнуляют ban_count, ' + 'поэтому следующий бан той же пары в обоих случаях стартует заново с ' + 'SOURCE_BAN_BASE_HOURS.'; + +COMMIT; diff --git a/tradein-mvp/backend/tests/services/test_proxy_pool.py b/tradein-mvp/backend/tests/services/test_proxy_pool.py index 9a1b9ea0..1a40a8aa 100644 --- a/tradein-mvp/backend/tests/services/test_proxy_pool.py +++ b/tradein-mvp/backend/tests/services/test_proxy_pool.py @@ -438,21 +438,38 @@ class FakeSession: [{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}] ) - if "DELETE FROM scrape_proxy_source_bans" in sql and "proxy_id = CAST" in sql: - # clear_source_bans: снять баны узла (все либо один source), #2600 п.2. - # Фильтр по reason (#2800) гейтим по подстроке боевого SQL — как ban-фильтры - # в acquire-ветке: иначе мок «чинил» бы код, который фильтра не содержит, и - # тест на «успешная проба не гасит чужой бан» остался бы зелёным на сломанном. + if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql: + # clear_source_bans (#3404): гасим строку (banned_until=now(), ban_count=0, + # cleared_at/cleared_reason) вместо DELETE — трассируемость снятия бана, строка + # доживает до штатного purge. Фильтр по reason (#2800) гейтим по подстроке + # боевого SQL — как ban-фильтры в acquire-ветке: иначе мок «чинил» бы код, + # который фильтра не содержит, и тест на «успешная проба не гасит чужой бан» + # остался бы зелёным на сломанном. filters_reason = "reason = CAST(:only_reason AS text)" in sql only_reason = p.get("only_reason") if filters_reason else None - cleared = [ - b - for b in self.bans - if b["proxy_id"] == p["proxy_id"] - and (p["source"] is None or b["source"] == p["source"]) - and (only_reason is None or b.get("reason") == only_reason) - ] - self.bans = [b for b in self.bans if b not in cleared] + now = datetime.now(UTC) + cleared: list[dict[str, Any]] = [] + for b in self.bans: + if b["proxy_id"] != p["proxy_id"]: + continue + if p["source"] is not None and b["source"] != p["source"]: + continue + if only_reason is not None and b.get("reason") != only_reason: + continue + # Гейт покоя (реальный WHERE): уже погашенную строку повторно не трогаем — + # без него повторный вызов сдвигал бы banned_until вперёд. + if ( + b.get("cleared_at") is not None + and b["ban_count"] == 0 + and b["banned_until"] <= now + ): + continue + b["banned_until"] = now + b["ban_count"] = 0 + b["cleared_at"] = now + b["cleared_reason"] = p["reason"] + b["updated_at"] = now + cleared.append(b) return _FakeResult([{"source": b["source"]} for b in cleared]) if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2) @@ -1419,6 +1436,10 @@ async def test_healthcheck_purges_long_expired_bans_only( # До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан # в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой # ручки ложное срабатывание детектора капчи (#2642) снималось бы только руками в SQL. +# +# #3404: снятие гасит строку (banned_until=now(), ban_count=0, cleared_at/cleared_reason), +# а не удаляет её — строка живёт для трассируемости до штатного purge, но для выдачи и +# для эскалации следующего бана неотличима от прежнего DELETE. def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]: @@ -1428,17 +1449,37 @@ def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]: "ban_count": ban_count, "banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS), "reason": f"banned:{source}", + "cleared_at": None, + "cleared_reason": None, } -def test_clear_source_bans_removes_all_bans_of_node() -> None: +def test_clear_source_bans_gates_all_bans_of_node() -> None: db = FakeSession( [_proxy(1, affinity="any"), _proxy(2, affinity="any")], bans=[_active_ban(1, "avito"), _active_ban(1, "cian"), _active_ban(2, "avito")], ) cleared = proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type] assert cleared == 2 - assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел + # #3404: строки НЕ удаляются — все три остаются в таблице (трассируемость). + assert {(b["proxy_id"], b["source"]) for b in db.bans} == { + (1, "avito"), + (1, "cian"), + (2, "avito"), + } + now = datetime.now(UTC) + node1_bans = [b for b in db.bans if b["proxy_id"] == 1] + assert len(node1_bans) == 2 + for b in node1_bans: + assert b["banned_until"] <= now + assert b["ban_count"] == 0 + assert b["cleared_at"] is not None + assert b["cleared_reason"] == "manual enable" + # чужой бан (proxy_id=2) не тронут — остаётся активным + other_ban = next(b for b in db.bans if b["proxy_id"] == 2) + assert other_ban["banned_until"] > now + assert other_ban["ban_count"] == 1 + assert other_ban["cleared_at"] is None # узел снова выдаётся источнику, который его банил lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type] assert lease is not None and lease.id == 1 @@ -1453,7 +1494,16 @@ def test_clear_source_bans_single_source_keeps_others() -> None: db, 1, source="avito", reason="ip rotated" ) assert cleared == 1 - assert [b["source"] for b in db.bans] == ["cian"] + # обе строки остаются (#3404), но только "avito" погашена + avito_ban = db._ban(1, "avito") + cian_ban = db._ban(1, "cian") + now = datetime.now(UTC) + assert avito_ban["banned_until"] <= now + assert avito_ban["ban_count"] == 0 + assert avito_ban["cleared_reason"] == "ip rotated" + assert cian_ban["banned_until"] > now + assert cian_ban["ban_count"] == 1 + assert cian_ban["cleared_at"] is None def test_clear_source_bans_noop_when_nothing_to_clear() -> None: @@ -1462,8 +1512,8 @@ def test_clear_source_bans_noop_when_nothing_to_clear() -> None: def test_clear_source_bans_resets_escalation() -> None: - """DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан - начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию.""" + """Гашение (banned_until=now(), ban_count=0), а не удаление строки: следующий бан той + же пары начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию (#3404).""" assert mark_banned is not None db = FakeSession( [_proxy(1, affinity="any"), _proxy(2, affinity="any")], diff --git a/tradein-mvp/backend/tests/services/test_proxy_rotation.py b/tradein-mvp/backend/tests/services/test_proxy_rotation.py index e10c73df..759a8f78 100644 --- a/tradein-mvp/backend/tests/services/test_proxy_rotation.py +++ b/tradein-mvp/backend/tests/services/test_proxy_rotation.py @@ -105,9 +105,19 @@ class FakeSession: ) return _FakeResult([]) - if "DELETE FROM scrape_proxy_source_bans" in sql: # clear_source_bans (#2600 п.2) - cleared = [b for b in self.source_bans if b["proxy_id"] == p["proxy_id"]] - self.source_bans = [b for b in self.source_bans if b not in cleared] + if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql: + # clear_source_bans (#3404): гасит строку (cleared_reason=reason), а НЕ + # удаляет — строка остаётся для трассируемости. Фейк мутирует найденные + # записи в месте, не убирая их из self.source_bans. + target_source = p.get("source") + cleared = [ + b + for b in self.source_bans + if b["proxy_id"] == p["proxy_id"] + and (target_source is None or b["source"] == target_source) + ] + for b in cleared: + b["cleared_reason"] = p["reason"] return _FakeResult([{"source": b["source"]} for b in cleared]) if "UPDATE scrape_proxies" in sql and "exit_ip" in sql: @@ -404,7 +414,13 @@ async def test_successful_rotation_clears_source_bans(monkeypatch: pytest.Monkey result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type] assert result.ok is True - assert db.source_bans == [{"proxy_id": 2, "source": "avito"}] + # #3404: строки не удаляются — гасятся (cleared_reason проставлен), остаются для + # трассируемости. Чужой узел (proxy_id=2) не тронут вообще. + proxy1_bans = [b for b in db.source_bans if b["proxy_id"] == 1] + assert len(proxy1_bans) == 2 + assert all(b.get("cleared_reason") == "exit ip rotated (status=200)" for b in proxy1_bans) + other_ban = next(b for b in db.source_bans if b["proxy_id"] == 2) + assert other_ban.get("cleared_reason") is None async def test_failed_rotation_keeps_source_bans(monkeypatch: pytest.MonkeyPatch) -> None: @@ -471,10 +487,13 @@ async def test_token_never_appears_in_reason_success(monkeypatch: pytest.MonkeyP fake_client, _ = _fake_async_client(response=(200, {"ip": "1.1.1.1"}), exception=None) monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client) - db = FakeSession(_proxy_row()) + db = FakeSession(_proxy_row(), source_bans=[{"proxy_id": 1, "source": "avito"}]) result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type] assert SECRET_TOKEN not in (result.reason or "") assert SECRET_TOKEN not in (result.new_ip or "") + # #3404: успешная ротация гасит бан и пишет cleared_reason — секрет не должен + # попасть и туда (та же гигиена, что для reason/note/логов). + assert SECRET_TOKEN not in (db.source_bans[0].get("cleared_reason") or "") async def test_token_never_appears_in_reason_on_401(monkeypatch: pytest.MonkeyPatch) -> None: @@ -574,7 +593,10 @@ async def test_token_never_appears_in_log_messages_or_sentry_text( else: monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", _no_http_allowed()) - db = FakeSession(_proxy_row(rotate_url=rotate_url)) + db = FakeSession( + _proxy_row(rotate_url=rotate_url), + source_bans=[{"proxy_id": 1, "source": "avito"}], + ) with caplog.at_level(logging.DEBUG): caplog.clear() await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type] @@ -583,6 +605,12 @@ async def test_token_never_appears_in_log_messages_or_sentry_text( assert SECRET_TOKEN not in record.getMessage(), ( f"scenario={name}: token leaked into log message args" ) + if name == "success": + # #3404: успешная ротация гасит бан и пишет cleared_reason — секрет не + # должен попасть и туда. + assert SECRET_TOKEN not in (db.source_bans[0].get("cleared_reason") or ""), ( + f"scenario={name}: token leaked into cleared_reason" + ) assert sentry_texts, "expected at least one Sentry capture (401 scenario)" assert all(SECRET_TOKEN not in text for text in sentry_texts) @@ -646,7 +674,9 @@ async def test_mobileproxy_success_clears_source_bans(monkeypatch: pytest.Monkey result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type] assert result.ok is True - assert db.source_bans == [] + # #3404: гашение, не удаление — строка остаётся, cleared_reason помечает снятие. + assert len(db.source_bans) == 1 + assert db.source_bans[0]["cleared_reason"] == "exit ip rotated (mobileproxy status=200)" async def test_mobileproxy_non_ok_status_is_failure_and_consumes_quota( @@ -737,7 +767,10 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs( fake_client, _ = _fake_async_client_get(response=response, exception=exception) monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client) - db = FakeSession(_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL)) + db = FakeSession( + _proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL), + source_bans=[{"proxy_id": 1, "source": "avito"}], + ) with caplog.at_level(logging.DEBUG): caplog.clear() result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type] @@ -749,6 +782,10 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs( assert MOBILEPROXY_KEY not in record.getMessage(), ( f"scenario={name}: proxy_key leaked into log message args" ) + if name == "success": + # #3404: успешная ротация гасит бан и пишет cleared_reason — ключ не + # должен попасть и туда. + assert MOBILEPROXY_KEY not in (db.source_bans[0].get("cleared_reason") or ""), name async def test_mobileproxy_unknown_query_shape_still_masks_url_in_refusal( diff --git a/tradein-mvp/backend/tests/test_2800_per_source_probe.py b/tradein-mvp/backend/tests/test_2800_per_source_probe.py index 7fa59d45..fcc0472e 100644 --- a/tradein-mvp/backend/tests/test_2800_per_source_probe.py +++ b/tradein-mvp/backend/tests/test_2800_per_source_probe.py @@ -302,7 +302,14 @@ async def test_probe_clears_only_its_own_ban(monkeypatch: pytest.MonkeyPatch) -> avito = db._ban(1, "avito") assert avito is not None, "чужой бан проба снимать не имеет права" assert (avito["reason"], avito["ban_count"]) == ("banned:avito", 1), "и не переписывать" - assert db._ban(1, "cian") is None, "свой вердикт проба обязана снять" + cian = db._ban(1, "cian") + # #3404: снятие теперь ГАСИТ строку, а не удаляет — вердикт пробы перестаёт + # блокировать выдачу (banned_until в прошлом, ban_count обнулён), но остаётся + # виден в истории как снятый досрочно. + assert cian is not None, "погашенная строка живёт до purge (#3404)" + assert cian["ban_count"] == 0, "свой вердикт проба обязана снять" + assert cian["banned_until"] <= datetime.now(UTC), "и снять немедленно" + assert cian["cleared_at"] is not None, "снятие должно быть отмечено" assert counters["pair_cleared"] == 1 diff --git a/tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py b/tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py new file mode 100644 index 00000000..7ff728c3 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py @@ -0,0 +1,403 @@ +"""#3404: атрибуция прогона к узлу (scrape_runs.proxy_id) + гашение бана вместо DELETE. + +Offline-тесты (без live БД) в стиле test_proxy_pool.py / test_3390_single_runs_module.py: +stateful fake-сессия эмулирует таблицы scrape_proxies / scrape_proxy_source_bans / +scrape_runs, интерпретируя SQL по ключевым фрагментам, так что acquire/mark_banned/ +clear_source_bans/attribute_run_proxy проверяются по фактическому изменению состояния, +а не по замоканному возврату. + +Покрытие: + 1. clear_source_bans гасит строку (banned_until<=now, ban_count=0, cleared_at + заполнен), а НЕ удаляет — строка остаётся в таблице. + 2. Эскалация 1:1: бан → clear → бан той же пары снова стартует с ровно + SOURCE_BAN_BASE_HOURS, а не удвоенного срока (главный тест issue). + 3. Погашенная строка не мешает acquire(source) — узел выдаётся как обычно. + 4. attribute_run_proxy доводит proxy_id прогона до scrape_runs через acquire(); + путь без run_id (health-checker, run_id=None) оставляет scrape_runs нетронутым. + 5. Смена узла за прогон (повторный acquire тем же run_id на другой узел) + отражается в counters['proxy_ids'] без дублей. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from datetime import UTC, datetime, timedelta +from typing import Any + +from app.services.proxy_pool import ( + SOURCE_BAN_BASE_HOURS, + acquire, + attribute_run_proxy, + clear_source_bans, + mark_banned, +) + +# ── stateful fake session (scrape_proxies + scrape_proxy_source_bans + scrape_runs) ── + + +class _FakeResult: + def __init__(self, rows: list[Any]) -> None: + self._rows = rows + + def mappings(self) -> _FakeResult: + return self + + def fetchone(self) -> Any: + return self._rows[0] if self._rows else None + + def first(self) -> Any: + return self._rows[0] if self._rows else None + + def fetchall(self) -> list[Any]: + return [type("Row", (), r)() if isinstance(r, dict) else r for r in self._rows] + + def all(self) -> list[Any]: + return list(self._rows) + + +def _proxy( + pid: int, + *, + affinity: str = "avito", + enabled: bool = True, + leased_by: int | None = None, +) -> dict[str, Any]: + return { + "id": pid, + "url": f"http://proxy{pid}", + "kind": "residential", + "rotate_url": None, + "browser_unfit_since": None, + "enabled": enabled, + "consecutive_fails": 0, + "provider_affinity": affinity, + "leased_by": leased_by, + "leased_at": None, + "expires_at": None, + "last_ok_at": None, + } + + +class RunAttributionDb: + """Мини-Postgres: scrape_proxies + scrape_proxy_source_bans + scrape_runs. + + Гейты по SQL-подстрокам скопированы из проверяемого кода (proxy_pool.py) — как в + test_proxy_pool.py/test_3390: значение проверяется по факту исполнения реального + statement'а, а не зашито ожиданием теста. + """ + + def __init__(self, proxies: list[dict[str, Any]]) -> None: + self.proxies = proxies + self.bans: list[dict[str, Any]] = [] + self.runs: dict[int, dict[str, Any]] = {} + self.set_local_statements: list[str] = [] + + def add_run(self, run_id: int) -> None: + self.runs[run_id] = {"proxy_id": None, "counters": {}} + + def _by_id(self, pid: int) -> dict[str, Any] | None: + return next((r for r in self.proxies 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: + b = self._ban(pid, source) + return b is not None and b["banned_until"] > datetime.now(UTC) + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = " ".join(str(stmt).split()) + p = params or {} + + if "pg_advisory_xact_lock" in sql: + return _FakeResult([]) + + if "expires_at <= now()" in sql and "FOR UPDATE SKIP LOCKED" not in sql: + return _FakeResult([]) # без просроченных узлов в этих тестах + + if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary/fallback) + provider = p["provider"] + max_fails = p["max_fails"] + primary = "provider_affinity IN" in sql + cands = [ + r + for r in self.proxies + if r["enabled"] + and r["consecutive_fails"] < max_fails + and r["leased_by"] is None + and not self._has_active_ban(r["id"], provider) + and (r["provider_affinity"] in (provider, "any") if primary else True) + ] + cands.sort( + key=lambda r: ( + r["browser_unfit_since"] is not None, + r["last_ok_at"] is None, + r["last_ok_at"] or datetime.min.replace(tzinfo=UTC), + r["id"], + ) + ) + return _FakeResult(cands[:1]) + + if "SET leased_by = CAST(:run_id" in sql: # acquire lease UPDATE + row = self._by_id(p["id"]) + if row is not None: + row["leased_by"] = p["run_id"] + row["leased_at"] = datetime.now(UTC) + return _FakeResult([]) + + if sql.startswith("SET LOCAL"): + # attribute_run_proxy ставит lock_timeout перед UPDATE scrape_runs: + # ждать на блокировке строки прогона в пути выдачи прокси нельзя. + # Для фейка это no-op, но проглатывать молча нечестно — гейт ниже + # ловит любой ДРУГОЙ незнакомый SQL. + self.set_local_statements.append(sql) + return _FakeResult([]) + + if "SET leased_by = NULL" in sql: # release + row = self._by_id(p["id"]) + if row is not None: + row["leased_by"] = None + row["leased_at"] = None + return _FakeResult([]) + + if "INSERT INTO scrape_proxy_source_bans" in sql: # mark_banned upsert + proxy_id, source = p["proxy_id"], p["source"] + if self._by_id(proxy_id) is None: + return _FakeResult([]) + max_fails = p["max_fails"] + still_available = any( + r["id"] != proxy_id + and r["enabled"] + and r["consecutive_fails"] < max_fails + and not self._has_active_ban(r["id"], source) + for r in self.proxies + ) + if not still_available: + return _FakeResult([]) # protected — последний узел + + now = datetime.now(UTC) + ban = self._ban(proxy_id, source) + reason = p["reason"] + if ban is None: + ban = { + "proxy_id": proxy_id, + "source": source, + "ban_count": 1, + "banned_until": now + timedelta(hours=p["base_hours"]), + "reason": reason, + "cleared_at": None, + "cleared_reason": None, + } + self.bans.append(ban) + else: + is_active_foreign = ( + ban["banned_until"] > now + and ban.get("reason") != reason + and reason != p.get("live_reason") + ) + if is_active_foreign: + return _FakeResult([]) # deferred — чужой активный владелец + 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"] = reason + ban["cleared_at"] = None + ban["cleared_reason"] = None + return _FakeResult( + [{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}] + ) + + if "UPDATE scrape_proxy_source_bans" in sql and "cleared_at" in sql: # clear_source_bans + proxy_id = p["proxy_id"] + source = p.get("source") + only_reason = p.get("only_reason") + reason = p["reason"] + now = datetime.now(UTC) + cleared = [] + for b in self.bans: + if b["proxy_id"] != proxy_id: + continue + if source is not None and b["source"] != source: + continue + if only_reason is not None and b.get("reason") != only_reason: + continue + if ( + b.get("cleared_at") is not None + and b["ban_count"] == 0 + and b["banned_until"] <= now + ): + continue # гейт покоя — уже погашена, no-op + b["banned_until"] = now + b["ban_count"] = 0 + b["cleared_at"] = now + b["cleared_reason"] = reason + cleared.append({"source": b["source"]}) + return _FakeResult(cleared) + + if "UPDATE scrape_runs" in sql: # attribute_run_proxy (#3404) + run = self.runs.get(p["run_id"]) + if run is None: + return _FakeResult([]) + run["proxy_id"] = p["proxy_id"] + proxy_ids: list[int] = run["counters"].get("proxy_ids", []) + if p["proxy_id"] not in proxy_ids: + proxy_ids = [*proxy_ids, p["proxy_id"]] + run["counters"]["proxy_ids"] = proxy_ids + return _FakeResult([{"id": p["run_id"]}]) + + raise AssertionError(f"unhandled SQL in fake session: {sql[:120]}") + + def commit(self) -> None: + pass + + def rollback(self) -> None: + pass + + +# ── 1. clear_source_bans гасит, а не удаляет ──────────────────────────────── + + +def test_clear_source_bans_gates_not_deletes() -> None: + db = RunAttributionDb([_proxy(1), _proxy(2)]) + assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type] + assert len(db.bans) == 1 + + cleared = clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type] + + assert cleared == 1 + assert len(db.bans) == 1, "строка должна остаться на месте, не удалиться" + ban = db.bans[0] + assert ban["banned_until"] <= datetime.now(UTC) + assert ban["ban_count"] == 0 + assert ban["cleared_at"] is not None + assert ban["cleared_reason"] == "manual enable" + + +# ── 2. Эскалация сохранилась 1:1 (главный тест issue) ─────────────────────── + + +def test_escalation_resets_to_base_after_clear() -> None: + db = RunAttributionDb([_proxy(1), _proxy(2)]) + + assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type] + first_ban = db.bans[0] + assert first_ban["ban_count"] == 1 + first_span = first_ban["banned_until"] - datetime.now(UTC) + assert timedelta(hours=SOURCE_BAN_BASE_HOURS - 1) < first_span <= timedelta( + hours=SOURCE_BAN_BASE_HOURS + ) + + clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type] + assert db.bans[0]["ban_count"] == 0 + + assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type] + second_ban = db.bans[0] + + assert second_ban["ban_count"] == 1, "после clear эскалация обязана стартовать заново" + second_span = second_ban["banned_until"] - datetime.now(UTC) + assert timedelta(hours=SOURCE_BAN_BASE_HOURS - 1) < second_span <= timedelta( + hours=SOURCE_BAN_BASE_HOURS + ), f"срок {second_span} обязан быть БАЗОВЫМ ({SOURCE_BAN_BASE_HOURS}ч), а не удвоенным" + + +# ── 3. Погашенная строка не мешает выдаче узла ────────────────────────────── + + +def test_acquire_ignores_cleared_ban() -> None: + db = RunAttributionDb([_proxy(1), _proxy(2)]) + mark_banned(db, 1, source="avito", reason="banned:avito") # type: ignore[arg-type] + clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type] + # узел 2 занят чужим прогоном — иначе acquire мог бы честно выдать его вместо + # узла 1, и тест перестал бы проверять именно "погашенный бан не блокирует". + db.proxies[1]["leased_by"] = 999 + + lease = acquire(db, "avito") # type: ignore[arg-type] + + assert lease is not None, "погашенный бан не должен блокировать acquire" + assert lease.id == 1 + + +def test_acquire_still_blocks_active_ban_after_clear_gate_check() -> None: + """Контроль ложноположительного теста выше: АКТИВНЫЙ (не погашенный) бан acquire + по-прежнему блокирует — это не сломано этим же изменением.""" + db = RunAttributionDb([_proxy(1), _proxy(2)]) + mark_banned(db, 1, source="avito", reason="banned:avito") # type: ignore[arg-type] + + lease = acquire(db, "avito") # type: ignore[arg-type] + + assert lease is not None + assert lease.id == 2, "узел 1 активно забанен по avito — выдан должен быть узел 2" + + +# ── 4. proxy_id прогона доезжает до scrape_runs ───────────────────────────── + + +def test_acquire_attributes_run_proxy() -> None: + db = RunAttributionDb([_proxy(1)]) + db.add_run(42) + + lease = acquire(db, "avito", run_id=42) # type: ignore[arg-type] + + assert lease is not None + assert db.runs[42]["proxy_id"] == lease.id == 1 + assert db.runs[42]["counters"]["proxy_ids"] == [1] + # Атрибуция обязана ограничить ожидание блокировки: строку прогона параллельно + # пишет heartbeat из другой сессии, а путь выдачи прокси ждать не может. + assert any("lock_timeout" in stmt for stmt in db.set_local_statements) + + +def test_acquire_without_run_id_leaves_scrape_runs_untouched() -> None: + """health-checker и прочие не-run вызовы (run_id=None) не трогают scrape_runs.""" + db = RunAttributionDb([_proxy(1)]) + db.add_run(42) + + lease = acquire(db, "avito") # type: ignore[arg-type] + + assert lease is not None + assert db.runs[42]["proxy_id"] is None + assert db.runs[42]["counters"] == {} + + +# ── 5. Смена узла за прогон отражается в counters ─────────────────────────── + + +def test_mid_run_proxy_rotation_recorded_in_counters() -> None: + """Ре-acquire тем же run_id на другой узел (старый lease ещё держится, как при + ротации в browser_fetcher — `_LEASE_ROTATE_AFTER_FAILS`) — оба узла в proxy_ids, + proxy_id несёт ПОСЛЕДНИЙ (сверено с attribute_run_proxy docstring, а не с ТЗ).""" + db = RunAttributionDb([_proxy(1), _proxy(2)]) + db.add_run(7) + + first = acquire(db, "avito", run_id=7) # type: ignore[arg-type] + assert first is not None + # старый lease НЕ освобождён (leased_by=7) — второй acquire тем же run_id + # неизбежно возьмёт другой свободный узел, детерминированно. + assert db._by_id(first.id)["leased_by"] == 7 + + second = acquire(db, "avito", run_id=7) # type: ignore[arg-type] + assert second is not None + assert second.id != first.id, "тестовая обвязка обязана взять ДРУГОЙ узел" + + assert db.runs[7]["proxy_id"] == second.id + assert set(db.runs[7]["counters"]["proxy_ids"]) == {first.id, second.id} + + +def test_attribute_run_proxy_idempotent_same_proxy() -> None: + """Повторная атрибуция ТЕМ ЖЕ узлом не дублирует id в proxy_ids.""" + db = RunAttributionDb([_proxy(1)]) + db.add_run(9) + + attribute_run_proxy(db, 9, 1) # type: ignore[arg-type] + attribute_run_proxy(db, 9, 1) # type: ignore[arg-type] + + assert db.runs[9]["counters"]["proxy_ids"] == [1] + + +def test_attribute_run_proxy_best_effort_on_missing_run() -> None: + """run_id не найден (гонка с финализацией) — best-effort no-op, исключение не летит.""" + db = RunAttributionDb([_proxy(1)]) + # run 999 никогда не создавался в этой fake-БД + attribute_run_proxy(db, 999, 1) # type: ignore[arg-type] # не должно бросить diff --git a/tradein-mvp/backend/tests/test_admin_proxies.py b/tradein-mvp/backend/tests/test_admin_proxies.py index 8726e69c..51dd39e0 100644 --- a/tradein-mvp/backend/tests/test_admin_proxies.py +++ b/tradein-mvp/backend/tests/test_admin_proxies.py @@ -52,7 +52,8 @@ def _scalar_result(value: object) -> MagicMock: def _cleared_bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock: - """Ответ на DELETE ... RETURNING source (proxy_pool.clear_source_bans, #2600 п.2). + """Ответ на UPDATE ... RETURNING source (proxy_pool.clear_source_bans, #3404: гасит + строку — banned_until=now()/ban_count=0/cleared_at/cleared_reason, — а не DELETE). PATCH enabled=true снимает баны узла по источникам — «ручное включение = чистый лист», как и обнуление disabled_reason рядом. @@ -337,8 +338,11 @@ def test_patch_enable_clears_source_bans(client: TestClient, db: MagicMock) -> N r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True}) assert r.status_code == 200, r.text - delete_sql = str(db.execute.call_args_list[1].args[0]) - assert "DELETE FROM scrape_proxy_source_bans" in delete_sql + clear_sql = str(db.execute.call_args_list[1].args[0]) + # #3404: гашение (UPDATE ... SET cleared_reason=...), а не DELETE — строка остаётся + # для трассируемости, но перестаёт блокировать выдачу. + assert "UPDATE scrape_proxy_source_bans" in clear_sql + assert "cleared_reason" in clear_sql assert db.execute.call_args_list[1].args[1]["proxy_id"] == 1 @@ -351,8 +355,10 @@ def test_patch_disable_keeps_source_bans(client: TestClient, db: MagicMock) -> N r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False}) assert r.status_code == 200, r.text + # #3404: clear_source_bans теперь UPDATE, а не DELETE — ищем по cleared_reason, + # уникальному для этого запроса маркеру. assert not any( - "DELETE FROM scrape_proxy_source_bans" in str(c.args[0]) for c in db.execute.call_args_list + "cleared_reason" in str(c.args[0]) for c in db.execute.call_args_list ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/run_context.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/run_context.py new file mode 100644 index 00000000..6ddf8049 --- /dev/null +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/run_context.py @@ -0,0 +1,35 @@ +"""ContextVar с id текущего прогона (#3404) — канал для мест, куда `run_id` не доходит +аргументом. + +ЗАЧЕМ ЭТО СУЩЕСТВУЕТ +-------------------- +`RealProxyProvider` (`app.services.scraper_adapters`) — ОДИН объект на весь жизненный +цикл `scheduler_main.py`, реализующий kit-протокол `ProxyProvider.acquire(provider)` +(`scraper_kit/contracts.py`) — сигнатура протокола не несёт `run_id`, а конструктору +адаптера неоткуда его взять: он создаётся один раз при старте планировщика, задолго до +того, как какой-либо конкретный прогон начнётся. Пробрасывать `run_id` через цепочку +вызовов от `create_run` до `proxy_pool.acquire` означало бы менять сигнатуры всех +провайдеров пайплайна (curl-путь `providers/_proxy.py`, `browser_fetcher`) ради одного +параметра, который в 99% вызовов (health-check, эстиматор) не нужен вовсе. + +`ContextVar`, а не глобальная переменная: `asyncio`-таски и синхронные Celery-контексты +не делят один поток предсказуемо, а `ContextVar` копируется в `asyncio.to_thread`/ +дочерние таски и не путается между КОНКУРЕНТНЫМИ прогонами в одном процессе — тот же +довод, что и у `app.core.public_request.is_public_request`. + +Что это НЕ гарантирует: если прогон и вызов `acquire()` идут в разных тредах/тасках, +созданных ДО `current_run_id.set()`, значение не унаследуется — атрибуция для такого +прогона молча останется NULL (не ошибка, см. `proxy_pool.attribute_run_proxy`). +""" + +from __future__ import annotations + +from contextvars import ContextVar + +#: id прогона (`scrape_runs.id`), который сейчас выполняется в этом контексте. +#: None — контекст не привязан ни к какому прогону (health-check, эстиматор, admin). +#: `runs.py` управляет им напрямую через `.set()`/`.set(None)` (не через +#: context-manager: жизненный цикл прогона не вложен в один `with`-блок — он +#: пересекает несколько функций/финализаторов, см. `create_run`/`mark_done`/ +#: `mark_failed`). +current_run_id: ContextVar[int | None] = ContextVar("scraper_kit_current_run_id", default=None) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index e2d8d009..4b38558a 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -45,6 +45,8 @@ from typing import Any from sqlalchemy import text from sqlalchemy.orm import Session +from scraper_kit.orchestration.run_context import current_run_id + try: # sentry опционален — kit не тянет его в зависимостях import sentry_sdk except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует @@ -621,6 +623,12 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …); отдельной колонки run_type больше нет — она 3244 прогона подряд молчала дефолтом 'city_sweep' и подписывала им, например, proxy_healthcheck (#2674). + + После коммита выставляет `current_run_id` (#3404) — канал, которым `run_id` + доходит до `RealProxyProvider.acquire`, живущего одним объектом на весь + планировщик (см. `run_context.py`). Финализаторы (`mark_done`/`mark_failed`/ + `mark_banned`/`mark_cancelled`) сбрасывают его обратно в None, чтобы следующий + прогон в том же треде не унаследовал чужой id. Returns run_id (bigint). """ row = db.execute( @@ -637,7 +645,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: ).fetchone() db.commit() assert row is not None, "scrape_runs INSERT returned no id" - return int(row.id) + run_id = int(row.id) + current_run_id.set(run_id) + return run_id def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = None) -> int: @@ -881,6 +891,9 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: if row is None: logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id) db.commit() + # #3404: прогон завершён — сбрасываем канал атрибуции прокси, чтобы следующий + # прогон в том же треде (или health-check между ними) не унаследовал этот run_id. + current_run_id.set(None) # #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая # для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort, # после коммита — статус уже персистирован в БД. @@ -923,6 +936,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) if row is None: logger.warning("mark_failed no-op: run_id=%d not in 'running' state", run_id) db.commit() + # #3404: см. mark_done — сброс канала атрибуции прокси. + current_run_id.set(None) # Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает). _alert_on_run_id(db, run_id) @@ -988,6 +1003,8 @@ def mark_banned( if row is None: logger.warning("mark_banned no-op: run_id=%d not in 'running' state", run_id) db.commit() + # #3404: см. mark_done — сброс канала атрибуции прокси. + current_run_id.set(None) # Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает). _alert_on_run_id(db, run_id) @@ -1173,6 +1190,9 @@ def mark_cancelled(db: Session, run_id: int) -> bool: {"run_id": run_id}, ).fetchone() db.commit() + if result is not None: + # #3404: см. mark_done — сброс канала атрибуции прокси. + current_run_id.set(None) return result is not None