feat(tradein/proxy): прогон знает свой узел, а снятый бан перестаёт стирать историю (#3404) #3405
11 changed files with 832 additions and 46 deletions
|
|
@ -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 вместе со сбросом,
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
107
tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql
Normal file
107
tradein-mvp/backend/data/sql/287_proxy_run_attribution.sql
Normal file
|
|
@ -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;
|
||||
|
|
@ -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")],
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
403
tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py
Normal file
403
tradein-mvp/backend/tests/test_3404_proxy_run_attribution.py
Normal file
|
|
@ -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] # не должно бросить
|
||||
|
|
@ -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
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue