Compare commits
No commits in common. "72e3bc9e2401a2800fbe1d4a8a3bf75fb95fc7b9" and "3d154f4ec04bb227c5fb56f28321358b86e57210" have entirely different histories.
72e3bc9e24
...
3d154f4ec0
11 changed files with 46 additions and 832 deletions
|
|
@ -1187,8 +1187,7 @@ class Settings(BaseSettings):
|
||||||
|
|
||||||
# #3283g: ротация exit-IP НА САМ БАН площадки, а не только по счётчику попыток.
|
# #3283g: ротация exit-IP НА САМ БАН площадки, а не только по счётчику попыток.
|
||||||
# Бан привязан к IP (замерено вживую: rotate_proxy() лечит забаненный узел за
|
# Бан привязан к IP (замерено вживую: rotate_proxy() лечит забаненный узел за
|
||||||
# секунды, clear_source_bans гасит бан в scrape_proxy_source_bans — строка живёт
|
# секунды, clear_source_bans снимает запись из scrape_proxy_source_bans), но
|
||||||
# до purge, #3404), но
|
|
||||||
# #3251/#3212 запрещают сбрасывать browser-context на КАЖДЫЙ блок -- сброс без
|
# #3251/#3212 запрещают сбрасывать browser-context на КАЖДЫЙ блок -- сброс без
|
||||||
# смены IP выбрасывает пройденный QRATOR-PoW и запускает самоподдерживающийся
|
# смены IP выбрасывает пройденный QRATOR-PoW и запускает самоподдерживающийся
|
||||||
# каскад блоков на том же адресе. rotate_on_ban МЕНЯЕТ IP вместе со сбросом,
|
# каскад блоков на том же адресе. rotate_on_ban МЕНЯЕТ IP вместе со сбросом,
|
||||||
|
|
|
||||||
|
|
@ -54,10 +54,7 @@ Self-healing (#2600):
|
||||||
WARNING (пул надо пополнять, #2638).
|
WARNING (пул надо пополнять, #2638).
|
||||||
- Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс
|
- Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс
|
||||||
автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился
|
автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился
|
||||||
IP, поэтому смена адреса делает строку недействительной. С #3404 «снятие» гасит
|
IP, поэтому смена адреса делает строку недействительной.
|
||||||
строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`),
|
|
||||||
а не удаляет её — строка живёт до штатного purge (SOURCE_BAN_PURGE_DAYS), но для
|
|
||||||
выдачи и для эскалации следующего бана это неотличимо от прежнего DELETE.
|
|
||||||
|
|
||||||
Ручное выключение vs авто-выключение (#2610):
|
Ручное выключение vs авто-выключение (#2610):
|
||||||
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
|
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
|
||||||
|
|
@ -146,7 +143,6 @@ __all__ = [
|
||||||
"STALE_LEASE_MINUTES",
|
"STALE_LEASE_MINUTES",
|
||||||
"ProxyLease",
|
"ProxyLease",
|
||||||
"acquire",
|
"acquire",
|
||||||
"attribute_run_proxy",
|
|
||||||
"clear_source_bans",
|
"clear_source_bans",
|
||||||
"mark_banned",
|
"mark_banned",
|
||||||
"mark_browser_health",
|
"mark_browser_health",
|
||||||
|
|
@ -445,11 +441,6 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
|
||||||
{"run_id": lease_marker, "id": proxy_id},
|
{"run_id": lease_marker, "id": proxy_id},
|
||||||
)
|
)
|
||||||
db.commit()
|
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:
|
if fallback_used:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity "
|
"proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity "
|
||||||
|
|
@ -483,78 +474,6 @@ 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:
|
def release(db: Session, proxy_id: int) -> None:
|
||||||
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
|
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
|
||||||
db.execute(
|
db.execute(
|
||||||
|
|
@ -1024,12 +943,6 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None =
|
||||||
CAST(:max_hours AS integer)
|
CAST(:max_hours AS integer)
|
||||||
) AS integer)),
|
) AS integer)),
|
||||||
reason = CAST(:reason AS text),
|
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()
|
updated_at = now()
|
||||||
-- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет),
|
-- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет),
|
||||||
-- свою же (обычная эскалация) или перебиваем боевым сбором — он сильнее
|
-- свою же (обычная эскалация) или перебиваем боевым сбором — он сильнее
|
||||||
|
|
@ -1145,25 +1058,12 @@ def clear_source_bans(
|
||||||
банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса,
|
банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса,
|
||||||
держа узел вне выдачи уже без причины.
|
держа узел вне выдачи уже без причины.
|
||||||
|
|
||||||
source=None — снять все баны узла; конкретный source — только его. Гасим строку
|
source=None — снять все баны узла; конкретный source — только его. DELETE, а не
|
||||||
(`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`), а НЕ
|
`banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации),
|
||||||
удаляем (#3404, было DELETE): строка доживает до штатного purge'а
|
а здесь мы как раз объявляем историю недействительной — новый бан начнётся с базовых
|
||||||
(`run_proxy_healthcheck`, SOURCE_BAN_PURGE_DAYS), но перестаёт блокировать
|
SOURCE_BAN_BASE_HOURS.
|
||||||
выдачу немедленно и перестаёт нести историю эскалации — `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, т.е. «снимать только строки, которые
|
`only_reason` — ФИЛЬТР по колонке reason, т.е. «снимать только строки, которые
|
||||||
написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое
|
написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое
|
||||||
|
|
@ -1175,22 +1075,14 @@ def clear_source_bans(
|
||||||
rows = db.execute(
|
rows = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_proxy_source_bans
|
DELETE FROM 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)
|
WHERE proxy_id = CAST(:proxy_id AS bigint)
|
||||||
AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text))
|
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))
|
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
|
RETURNING source
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason, "reason": reason},
|
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason},
|
||||||
).fetchall()
|
).fetchall()
|
||||||
db.commit()
|
db.commit()
|
||||||
if rows:
|
if rows:
|
||||||
|
|
@ -1573,11 +1465,6 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
|
||||||
# Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с
|
# Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с
|
||||||
# 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса:
|
# 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса:
|
||||||
# неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать».
|
# неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать».
|
||||||
# #3404: под этот же порог теперь попадают и ПОГАШЕННЫЕ clear_source_bans строки —
|
|
||||||
# для них banned_until == момент гашения (== cleared_at), т.е. таймер до purge
|
|
||||||
# отсчитывается от гашения, а не от исходного истечения бана. ban_count у них уже
|
|
||||||
# 0 к моменту гашения, так что покидающий purge их не «сбрасывает» повторно —
|
|
||||||
# он просто убирает уже неактуальный след из таблицы.
|
|
||||||
purged = len(
|
purged = len(
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
|
|
|
||||||
|
|
@ -221,17 +221,10 @@ class RealProxyProvider:
|
||||||
|
|
||||||
def acquire(self, provider: str) -> ProxyLease | None:
|
def acquire(self, provider: str) -> ProxyLease | None:
|
||||||
from scraper_kit.contracts import ProxyLease as _KitProxyLease
|
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()
|
db = _SessionLocal()
|
||||||
try:
|
try:
|
||||||
lease = _proxy_pool.acquire(db, provider, run_id=run_id)
|
lease = _proxy_pool.acquire(db, provider)
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
if lease is None:
|
if lease is None:
|
||||||
|
|
|
||||||
|
|
@ -1,107 +0,0 @@
|
||||||
-- 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,38 +438,21 @@ class FakeSession:
|
||||||
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
|
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
|
||||||
)
|
)
|
||||||
|
|
||||||
if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql:
|
if "DELETE FROM scrape_proxy_source_bans" in sql and "proxy_id = CAST" in sql:
|
||||||
# clear_source_bans (#3404): гасим строку (banned_until=now(), ban_count=0,
|
# clear_source_bans: снять баны узла (все либо один source), #2600 п.2.
|
||||||
# cleared_at/cleared_reason) вместо DELETE — трассируемость снятия бана, строка
|
# Фильтр по reason (#2800) гейтим по подстроке боевого SQL — как ban-фильтры
|
||||||
# доживает до штатного purge. Фильтр по reason (#2800) гейтим по подстроке
|
# в acquire-ветке: иначе мок «чинил» бы код, который фильтра не содержит, и
|
||||||
# боевого SQL — как ban-фильтры в acquire-ветке: иначе мок «чинил» бы код,
|
# тест на «успешная проба не гасит чужой бан» остался бы зелёным на сломанном.
|
||||||
# который фильтра не содержит, и тест на «успешная проба не гасит чужой бан»
|
|
||||||
# остался бы зелёным на сломанном.
|
|
||||||
filters_reason = "reason = CAST(:only_reason AS text)" in sql
|
filters_reason = "reason = CAST(:only_reason AS text)" in sql
|
||||||
only_reason = p.get("only_reason") if filters_reason else None
|
only_reason = p.get("only_reason") if filters_reason else None
|
||||||
now = datetime.now(UTC)
|
cleared = [
|
||||||
cleared: list[dict[str, Any]] = []
|
b
|
||||||
for b in self.bans:
|
for b in self.bans
|
||||||
if b["proxy_id"] != p["proxy_id"]:
|
if b["proxy_id"] == p["proxy_id"]
|
||||||
continue
|
and (p["source"] is None or b["source"] == p["source"])
|
||||||
if p["source"] is not None and b["source"] != p["source"]:
|
and (only_reason is None or b.get("reason") == only_reason)
|
||||||
continue
|
]
|
||||||
if only_reason is not None and b.get("reason") != only_reason:
|
self.bans = [b for b in self.bans if b not in cleared]
|
||||||
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])
|
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||||
|
|
||||||
if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2)
|
if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2)
|
||||||
|
|
@ -1436,10 +1419,6 @@ async def test_healthcheck_purges_long_expired_bans_only(
|
||||||
# До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан
|
# До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан
|
||||||
# в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой
|
# в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой
|
||||||
# ручки ложное срабатывание детектора капчи (#2642) снималось бы только руками в SQL.
|
# ручки ложное срабатывание детектора капчи (#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]:
|
def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
|
||||||
|
|
@ -1449,37 +1428,17 @@ def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
|
||||||
"ban_count": ban_count,
|
"ban_count": ban_count,
|
||||||
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS),
|
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS),
|
||||||
"reason": f"banned:{source}",
|
"reason": f"banned:{source}",
|
||||||
"cleared_at": None,
|
|
||||||
"cleared_reason": None,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
def test_clear_source_bans_gates_all_bans_of_node() -> None:
|
def test_clear_source_bans_removes_all_bans_of_node() -> None:
|
||||||
db = FakeSession(
|
db = FakeSession(
|
||||||
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||||||
bans=[_active_ban(1, "avito"), _active_ban(1, "cian"), _active_ban(2, "avito")],
|
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]
|
cleared = proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
|
||||||
assert cleared == 2
|
assert cleared == 2
|
||||||
# #3404: строки НЕ удаляются — все три остаются в таблице (трассируемость).
|
assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел
|
||||||
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]
|
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||||
assert lease is not None and lease.id == 1
|
assert lease is not None and lease.id == 1
|
||||||
|
|
@ -1494,16 +1453,7 @@ def test_clear_source_bans_single_source_keeps_others() -> None:
|
||||||
db, 1, source="avito", reason="ip rotated"
|
db, 1, source="avito", reason="ip rotated"
|
||||||
)
|
)
|
||||||
assert cleared == 1
|
assert cleared == 1
|
||||||
# обе строки остаются (#3404), но только "avito" погашена
|
assert [b["source"] for b in db.bans] == ["cian"]
|
||||||
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:
|
def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
|
||||||
|
|
@ -1512,8 +1462,8 @@ def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
|
||||||
|
|
||||||
|
|
||||||
def test_clear_source_bans_resets_escalation() -> None:
|
def test_clear_source_bans_resets_escalation() -> None:
|
||||||
"""Гашение (banned_until=now(), ban_count=0), а не удаление строки: следующий бан той
|
"""DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан
|
||||||
же пары начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию (#3404)."""
|
начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию."""
|
||||||
assert mark_banned is not None
|
assert mark_banned is not None
|
||||||
db = FakeSession(
|
db = FakeSession(
|
||||||
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||||||
|
|
|
||||||
|
|
@ -105,19 +105,9 @@ class FakeSession:
|
||||||
)
|
)
|
||||||
return _FakeResult([])
|
return _FakeResult([])
|
||||||
|
|
||||||
if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql:
|
if "DELETE FROM scrape_proxy_source_bans" in sql: # clear_source_bans (#2600 п.2)
|
||||||
# clear_source_bans (#3404): гасит строку (cleared_reason=reason), а НЕ
|
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]
|
||||||
# записи в месте, не убирая их из 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])
|
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||||
|
|
||||||
if "UPDATE scrape_proxies" in sql and "exit_ip" in sql:
|
if "UPDATE scrape_proxies" in sql and "exit_ip" in sql:
|
||||||
|
|
@ -414,13 +404,7 @@ async def test_successful_rotation_clears_source_bans(monkeypatch: pytest.Monkey
|
||||||
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
||||||
assert result.ok is True
|
assert result.ok is True
|
||||||
# #3404: строки не удаляются — гасятся (cleared_reason проставлен), остаются для
|
assert db.source_bans == [{"proxy_id": 2, "source": "avito"}]
|
||||||
# трассируемости. Чужой узел (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:
|
async def test_failed_rotation_keeps_source_bans(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
|
@ -487,13 +471,10 @@ 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)
|
fake_client, _ = _fake_async_client(response=(200, {"ip": "1.1.1.1"}), exception=None)
|
||||||
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
||||||
|
|
||||||
db = FakeSession(_proxy_row(), source_bans=[{"proxy_id": 1, "source": "avito"}])
|
db = FakeSession(_proxy_row())
|
||||||
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
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.reason or "")
|
||||||
assert SECRET_TOKEN not in (result.new_ip 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:
|
async def test_token_never_appears_in_reason_on_401(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
|
@ -593,10 +574,7 @@ async def test_token_never_appears_in_log_messages_or_sentry_text(
|
||||||
else:
|
else:
|
||||||
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", _no_http_allowed())
|
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", _no_http_allowed())
|
||||||
|
|
||||||
db = FakeSession(
|
db = FakeSession(_proxy_row(rotate_url=rotate_url))
|
||||||
_proxy_row(rotate_url=rotate_url),
|
|
||||||
source_bans=[{"proxy_id": 1, "source": "avito"}],
|
|
||||||
)
|
|
||||||
with caplog.at_level(logging.DEBUG):
|
with caplog.at_level(logging.DEBUG):
|
||||||
caplog.clear()
|
caplog.clear()
|
||||||
await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
@ -605,12 +583,6 @@ async def test_token_never_appears_in_log_messages_or_sentry_text(
|
||||||
assert SECRET_TOKEN not in record.getMessage(), (
|
assert SECRET_TOKEN not in record.getMessage(), (
|
||||||
f"scenario={name}: token leaked into log message args"
|
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 sentry_texts, "expected at least one Sentry capture (401 scenario)"
|
||||||
assert all(SECRET_TOKEN not in text for text in sentry_texts)
|
assert all(SECRET_TOKEN not in text for text in sentry_texts)
|
||||||
|
|
@ -674,9 +646,7 @@ async def test_mobileproxy_success_clears_source_bans(monkeypatch: pytest.Monkey
|
||||||
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
||||||
assert result.ok is True
|
assert result.ok is True
|
||||||
# #3404: гашение, не удаление — строка остаётся, cleared_reason помечает снятие.
|
assert db.source_bans == []
|
||||||
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(
|
async def test_mobileproxy_non_ok_status_is_failure_and_consumes_quota(
|
||||||
|
|
@ -767,10 +737,7 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs(
|
||||||
fake_client, _ = _fake_async_client_get(response=response, exception=exception)
|
fake_client, _ = _fake_async_client_get(response=response, exception=exception)
|
||||||
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
||||||
|
|
||||||
db = FakeSession(
|
db = FakeSession(_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL))
|
||||||
_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL),
|
|
||||||
source_bans=[{"proxy_id": 1, "source": "avito"}],
|
|
||||||
)
|
|
||||||
with caplog.at_level(logging.DEBUG):
|
with caplog.at_level(logging.DEBUG):
|
||||||
caplog.clear()
|
caplog.clear()
|
||||||
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
|
||||||
|
|
@ -782,10 +749,6 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs(
|
||||||
assert MOBILEPROXY_KEY not in record.getMessage(), (
|
assert MOBILEPROXY_KEY not in record.getMessage(), (
|
||||||
f"scenario={name}: proxy_key leaked into log message args"
|
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(
|
async def test_mobileproxy_unknown_query_shape_still_masks_url_in_refusal(
|
||||||
|
|
|
||||||
|
|
@ -302,14 +302,7 @@ async def test_probe_clears_only_its_own_ban(monkeypatch: pytest.MonkeyPatch) ->
|
||||||
avito = db._ban(1, "avito")
|
avito = db._ban(1, "avito")
|
||||||
assert avito is not None, "чужой бан проба снимать не имеет права"
|
assert avito is not None, "чужой бан проба снимать не имеет права"
|
||||||
assert (avito["reason"], avito["ban_count"]) == ("banned:avito", 1), "и не переписывать"
|
assert (avito["reason"], avito["ban_count"]) == ("banned:avito", 1), "и не переписывать"
|
||||||
cian = db._ban(1, "cian")
|
assert db._ban(1, "cian") is None, "свой вердикт проба обязана снять"
|
||||||
# #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
|
assert counters["pair_cleared"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,403 +0,0 @@
|
||||||
"""#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,8 +52,7 @@ def _scalar_result(value: object) -> MagicMock:
|
||||||
|
|
||||||
|
|
||||||
def _cleared_bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock:
|
def _cleared_bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock:
|
||||||
"""Ответ на UPDATE ... RETURNING source (proxy_pool.clear_source_bans, #3404: гасит
|
"""Ответ на DELETE ... RETURNING source (proxy_pool.clear_source_bans, #2600 п.2).
|
||||||
строку — banned_until=now()/ban_count=0/cleared_at/cleared_reason, — а не DELETE).
|
|
||||||
|
|
||||||
PATCH enabled=true снимает баны узла по источникам — «ручное включение = чистый
|
PATCH enabled=true снимает баны узла по источникам — «ручное включение = чистый
|
||||||
лист», как и обнуление disabled_reason рядом.
|
лист», как и обнуление disabled_reason рядом.
|
||||||
|
|
@ -338,11 +337,8 @@ def test_patch_enable_clears_source_bans(client: TestClient, db: MagicMock) -> N
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
clear_sql = str(db.execute.call_args_list[1].args[0])
|
delete_sql = str(db.execute.call_args_list[1].args[0])
|
||||||
# #3404: гашение (UPDATE ... SET cleared_reason=...), а не DELETE — строка остаётся
|
assert "DELETE FROM scrape_proxy_source_bans" in delete_sql
|
||||||
# для трассируемости, но перестаёт блокировать выдачу.
|
|
||||||
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
|
assert db.execute.call_args_list[1].args[1]["proxy_id"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -355,10 +351,8 @@ def test_patch_disable_keeps_source_bans(client: TestClient, db: MagicMock) -> N
|
||||||
|
|
||||||
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
|
||||||
assert r.status_code == 200, r.text
|
assert r.status_code == 200, r.text
|
||||||
# #3404: clear_source_bans теперь UPDATE, а не DELETE — ищем по cleared_reason,
|
|
||||||
# уникальному для этого запроса маркеру.
|
|
||||||
assert not any(
|
assert not any(
|
||||||
"cleared_reason" in str(c.args[0]) for c in db.execute.call_args_list
|
"DELETE FROM scrape_proxy_source_bans" in str(c.args[0]) for c in db.execute.call_args_list
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,35 +0,0 @@
|
||||||
"""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,8 +45,6 @@ from typing import Any
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from scraper_kit.orchestration.run_context import current_run_id
|
|
||||||
|
|
||||||
try: # sentry опционален — kit не тянет его в зависимостях
|
try: # sentry опционален — kit не тянет его в зависимостях
|
||||||
import sentry_sdk
|
import sentry_sdk
|
||||||
except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует
|
except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует
|
||||||
|
|
@ -623,12 +621,6 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
||||||
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
||||||
дефолтом 'city_sweep' и подписывала им, например, proxy_healthcheck (#2674).
|
дефолтом '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).
|
Returns run_id (bigint).
|
||||||
"""
|
"""
|
||||||
row = db.execute(
|
row = db.execute(
|
||||||
|
|
@ -645,9 +637,7 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
).fetchone()
|
).fetchone()
|
||||||
db.commit()
|
db.commit()
|
||||||
assert row is not None, "scrape_runs INSERT returned no id"
|
assert row is not None, "scrape_runs INSERT returned no id"
|
||||||
run_id = int(row.id)
|
return 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:
|
def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = None) -> int:
|
||||||
|
|
@ -891,9 +881,6 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
if row is None:
|
if row is None:
|
||||||
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
|
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
|
||||||
db.commit()
|
db.commit()
|
||||||
# #3404: прогон завершён — сбрасываем канал атрибуции прокси, чтобы следующий
|
|
||||||
# прогон в том же треде (или health-check между ними) не унаследовал этот run_id.
|
|
||||||
current_run_id.set(None)
|
|
||||||
# #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая
|
# #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая
|
||||||
# для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort,
|
# для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort,
|
||||||
# после коммита — статус уже персистирован в БД.
|
# после коммита — статус уже персистирован в БД.
|
||||||
|
|
@ -936,8 +923,6 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
||||||
if row is None:
|
if row is None:
|
||||||
logger.warning("mark_failed no-op: run_id=%d not in 'running' state", run_id)
|
logger.warning("mark_failed no-op: run_id=%d not in 'running' state", run_id)
|
||||||
db.commit()
|
db.commit()
|
||||||
# #3404: см. mark_done — сброс канала атрибуции прокси.
|
|
||||||
current_run_id.set(None)
|
|
||||||
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
|
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
|
||||||
_alert_on_run_id(db, run_id)
|
_alert_on_run_id(db, run_id)
|
||||||
|
|
||||||
|
|
@ -1003,8 +988,6 @@ def mark_banned(
|
||||||
if row is None:
|
if row is None:
|
||||||
logger.warning("mark_banned no-op: run_id=%d not in 'running' state", run_id)
|
logger.warning("mark_banned no-op: run_id=%d not in 'running' state", run_id)
|
||||||
db.commit()
|
db.commit()
|
||||||
# #3404: см. mark_done — сброс канала атрибуции прокси.
|
|
||||||
current_run_id.set(None)
|
|
||||||
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
|
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
|
||||||
_alert_on_run_id(db, run_id)
|
_alert_on_run_id(db, run_id)
|
||||||
|
|
||||||
|
|
@ -1190,9 +1173,6 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
|
||||||
{"run_id": run_id},
|
{"run_id": run_id},
|
||||||
).fetchone()
|
).fetchone()
|
||||||
db.commit()
|
db.commit()
|
||||||
if result is not None:
|
|
||||||
# #3404: см. mark_done — сброс канала атрибуции прокси.
|
|
||||||
current_run_id.set(None)
|
|
||||||
return result is not None
|
return result is not None
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue