Compare commits
No commits in common. "f4099cb89a0542b0eea4555f190b042e0df8c440" and "06ad504505041c808ff5f0a16b11945b11748c74" have entirely different histories.
f4099cb89a
...
06ad504505
21 changed files with 59 additions and 1518 deletions
|
|
@ -1187,8 +1187,7 @@ class Settings(BaseSettings):
|
|||
|
||||
# #3283g: ротация exit-IP НА САМ БАН площадки, а не только по счётчику попыток.
|
||||
# Бан привязан к IP (замерено вживую: rotate_proxy() лечит забаненный узел за
|
||||
# секунды, clear_source_bans гасит бан в scrape_proxy_source_bans — строка живёт
|
||||
# до purge, #3404), но
|
||||
# секунды, clear_source_bans снимает запись из scrape_proxy_source_bans), но
|
||||
# #3251/#3212 запрещают сбрасывать browser-context на КАЖДЫЙ блок -- сброс без
|
||||
# смены IP выбрасывает пройденный QRATOR-PoW и запускает самоподдерживающийся
|
||||
# каскад блоков на том же адресе. rotate_on_ban МЕНЯЕТ IP вместе со сбросом,
|
||||
|
|
|
|||
|
|
@ -54,10 +54,7 @@ Self-healing (#2600):
|
|||
WARNING (пул надо пополнять, #2638).
|
||||
- Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс
|
||||
автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился
|
||||
IP, поэтому смена адреса делает строку недействительной. С #3404 «снятие» гасит
|
||||
строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`),
|
||||
а не удаляет её — строка живёт до штатного purge (SOURCE_BAN_PURGE_DAYS), но для
|
||||
выдачи и для эскалации следующего бана это неотличимо от прежнего DELETE.
|
||||
IP, поэтому смена адреса делает строку недействительной.
|
||||
|
||||
Ручное выключение vs авто-выключение (#2610):
|
||||
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
|
||||
|
|
@ -146,7 +143,6 @@ __all__ = [
|
|||
"STALE_LEASE_MINUTES",
|
||||
"ProxyLease",
|
||||
"acquire",
|
||||
"attribute_run_proxy",
|
||||
"clear_source_bans",
|
||||
"mark_banned",
|
||||
"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},
|
||||
)
|
||||
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 "
|
||||
|
|
@ -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:
|
||||
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
|
||||
db.execute(
|
||||
|
|
@ -1016,21 +935,15 @@ 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),
|
||||
-- #3404: строка могла быть погашена clear_source_bans (banned_until
|
||||
-- в прошлом попадает в WHERE ниже) — новый бан затирает её метки
|
||||
-- гашения, иначе на СНОВА забаненной паре висели бы cleared_at/
|
||||
-- cleared_reason от предыдущего, уже неактуального гашения.
|
||||
cleared_at = NULL,
|
||||
cleared_reason = NULL,
|
||||
updated_at = now()
|
||||
reason = CAST(:reason AS text),
|
||||
updated_at = now()
|
||||
-- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет),
|
||||
-- свою же (обычная эскалация) или перебиваем боевым сбором — он сильнее
|
||||
-- пробы. Иначе 0 rows и ветка "deferred" ниже (дефект #2803).
|
||||
|
|
@ -1145,25 +1058,12 @@ def clear_source_bans(
|
|||
банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса,
|
||||
держа узел вне выдачи уже без причины.
|
||||
|
||||
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).
|
||||
source=None — снять все баны узла; конкретный source — только его. DELETE, а не
|
||||
`banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации),
|
||||
а здесь мы как раз объявляем историю недействительной — новый бан начнётся с базовых
|
||||
SOURCE_BAN_BASE_HOURS.
|
||||
|
||||
Идемпотентно и в другую сторону: повторный вызов на уже погашенной строке
|
||||
(последний предикат в WHERE) её не трогает — 0 rows, `banned_until` НЕ
|
||||
сдвигается вперёд. Без этого условия повторный `PATCH enabled=true` двигал бы
|
||||
`banned_until` на каждый вызов и отодвигал бы purge на неопределённый срок.
|
||||
|
||||
`reason` идёт в лог (человекочитаемый повод — «manual enable», «ip rotated») и
|
||||
теперь ЕЩЁ в колонку `cleared_reason` — постоянный след того, кто и почему
|
||||
погасил бан.
|
||||
`reason` идёт только в лог (человекочитаемый повод — «manual enable», «ip rotated»).
|
||||
|
||||
`only_reason` — ФИЛЬТР по колонке reason, т.е. «снимать только строки, которые
|
||||
написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое
|
||||
|
|
@ -1175,22 +1075,14 @@ def clear_source_bans(
|
|||
rows = db.execute(
|
||||
text(
|
||||
"""
|
||||
UPDATE scrape_proxy_source_bans
|
||||
SET banned_until = now(),
|
||||
ban_count = 0,
|
||||
cleared_at = now(),
|
||||
cleared_reason = CAST(:reason AS text),
|
||||
updated_at = now()
|
||||
DELETE FROM scrape_proxy_source_bans
|
||||
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, "reason": reason},
|
||||
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason},
|
||||
).fetchall()
|
||||
db.commit()
|
||||
if rows:
|
||||
|
|
@ -1573,11 +1465,6 @@ 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,17 +221,10 @@ 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, run_id=run_id)
|
||||
lease = _proxy_pool.acquire(db, provider)
|
||||
finally:
|
||||
db.close()
|
||||
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;
|
||||
|
|
@ -490,7 +490,7 @@ def _sidecar_ban_page_error(upstream_status: int | None = 401) -> SidecarBanPage
|
|||
|
||||
@pytest.mark.asyncio
|
||||
async def test_fetch_detail_does_not_report_ban_again_on_sidecar_ban_page() -> None:
|
||||
"""Sidecar-бан рапортует сам фетчер (`report_platform_ban`, #3288) — здесь уже нет.
|
||||
"""Sidecar-бан рапортует сам фетчер (`_report_platform_ban`, #3288) — здесь уже нет.
|
||||
|
||||
Повтор отсюда приходит ПОСЛЕ возможной ротации lease по fail-streak и банил бы
|
||||
свежий узел. Ветка `parse_detail_html` ниже — другой случай: тот детект наш,
|
||||
|
|
|
|||
|
|
@ -438,38 +438,21 @@ class FakeSession:
|
|||
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
|
||||
)
|
||||
|
||||
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-ветке: иначе мок «чинил» бы код,
|
||||
# который фильтра не содержит, и тест на «успешная проба не гасит чужой бан»
|
||||
# остался бы зелёным на сломанном.
|
||||
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-ветке: иначе мок «чинил» бы код, который фильтра не содержит, и
|
||||
# тест на «успешная проба не гасит чужой бан» остался бы зелёным на сломанном.
|
||||
filters_reason = "reason = CAST(:only_reason AS text)" in sql
|
||||
only_reason = p.get("only_reason") if filters_reason else None
|
||||
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)
|
||||
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]
|
||||
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||
|
||||
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). Теперь бан
|
||||
# в отдельной таблице и истекает только по таймеру (до 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]:
|
||||
|
|
@ -1449,37 +1428,17 @@ 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_gates_all_bans_of_node() -> None:
|
||||
def test_clear_source_bans_removes_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
|
||||
# #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
|
||||
assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел
|
||||
# узел снова выдаётся источнику, который его банил
|
||||
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
|
||||
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"
|
||||
)
|
||||
assert cleared == 1
|
||||
# обе строки остаются (#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
|
||||
assert [b["source"] for b in db.bans] == ["cian"]
|
||||
|
||||
|
||||
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:
|
||||
"""Гашение (banned_until=now(), ban_count=0), а не удаление строки: следующий бан той
|
||||
же пары начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию (#3404)."""
|
||||
"""DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан
|
||||
начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию."""
|
||||
assert mark_banned is not None
|
||||
db = FakeSession(
|
||||
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
|
||||
|
|
|
|||
|
|
@ -105,19 +105,9 @@ class FakeSession:
|
|||
)
|
||||
return _FakeResult([])
|
||||
|
||||
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"]
|
||||
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]
|
||||
return _FakeResult([{"source": b["source"]} for b in cleared])
|
||||
|
||||
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]
|
||||
|
||||
assert result.ok is True
|
||||
# #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
|
||||
assert db.source_bans == [{"proxy_id": 2, "source": "avito"}]
|
||||
|
||||
|
||||
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)
|
||||
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]
|
||||
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:
|
||||
|
|
@ -593,10 +574,7 @@ 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),
|
||||
source_bans=[{"proxy_id": 1, "source": "avito"}],
|
||||
)
|
||||
db = FakeSession(_proxy_row(rotate_url=rotate_url))
|
||||
with caplog.at_level(logging.DEBUG):
|
||||
caplog.clear()
|
||||
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(), (
|
||||
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)
|
||||
|
|
@ -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]
|
||||
|
||||
assert result.ok is True
|
||||
# #3404: гашение, не удаление — строка остаётся, cleared_reason помечает снятие.
|
||||
assert len(db.source_bans) == 1
|
||||
assert db.source_bans[0]["cleared_reason"] == "exit ip rotated (mobileproxy status=200)"
|
||||
assert db.source_bans == []
|
||||
|
||||
|
||||
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)
|
||||
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
|
||||
|
||||
db = FakeSession(
|
||||
_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL),
|
||||
source_bans=[{"proxy_id": 1, "source": "avito"}],
|
||||
)
|
||||
db = FakeSession(_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL))
|
||||
with caplog.at_level(logging.DEBUG):
|
||||
caplog.clear()
|
||||
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(), (
|
||||
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,14 +302,7 @@ 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), "и не переписывать"
|
||||
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 db._ban(1, "cian") is None, "свой вердикт проба обязана снять"
|
||||
assert counters["pair_cleared"] == 1
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -99,7 +99,7 @@ async def test_fetch_detail_sidecar_ban_page_raises_platform_block_not_infra(
|
|||
async def test_fetch_detail_sidecar_ban_page_does_not_report_ban_again(
|
||||
sidecar_status: int,
|
||||
) -> None:
|
||||
"""Бан рапортует ФЕТЧЕР (`report_platform_ban`), провайдер — уже нет (#3288).
|
||||
"""Бан рапортует ФЕТЧЕР (`_report_platform_ban`), провайдер — уже нет (#3288).
|
||||
|
||||
До #3288 здесь стоял `assert_called_once()`, и это было верно, пока фетчер о
|
||||
бане не знал. Теперь рапорт идёт из `_post_fetch` — РАНЬШЕ и по правильному
|
||||
|
|
|
|||
|
|
@ -370,7 +370,7 @@ async def test_empty_pool_during_post_ban_rotation_keeps_platform_diagnosis() ->
|
|||
"""Ротация — best-effort: её `NoProxyAvailableError` не должна съесть бан.
|
||||
|
||||
Второго узла в пуле нет, поэтому `_acquire_lease()` внутри ротации падает.
|
||||
Фальсификация: убери `except NoProxyAvailableError` в `report_platform_ban` —
|
||||
Фальсификация: убери `except NoProxyAvailableError` в `_report_platform_ban` —
|
||||
наверх уедет она вместо `SidecarBanPageError`, провайдер завернёт её в
|
||||
AvitoSidecarUnavailableError, и подтверждённый отказ площадки попадёт в
|
||||
`scrape_runs.ban_kind` как 'infra' (ровно та подмена, что чинил #3283).
|
||||
|
|
|
|||
|
|
@ -1,200 +0,0 @@
|
|||
"""#3402: капча Циана приходит с HTTP 200 — это отказ площадки, а не «не разобрали».
|
||||
|
||||
Замер прода 06.09.2026: Циан отдаёт капчу (`<title>Captcha - база объявлений ЦИАН`,
|
||||
44 КБ; `<title>Вы не робот?`, 16 КБ) и страницу ошибки (`<title>Ошибка - Циан`, 374 КБ)
|
||||
с кодом **200**. Детектор сайдкара их не знал (`_REFUSAL_STATUSES` = {403,429}, маркеры
|
||||
сняты с Авито/Домклика), HTML уезжал клиенту как успех, `extract_state` возвращал None
|
||||
и провайдер печатал «defaultState extraction failed» — то есть отказ ПЛОЩАДКИ читался
|
||||
как дрейф НАШЕЙ разметки. Аренда не менялась: `fetch()` уже отрапортовал
|
||||
`mark_health(ok=True)` (HTTP-уровень успешен), fail-streak обнулялся, и один капча-узел
|
||||
сжигал весь батч — прогоны 6200: 0/210; 6123/6091/6052/6032/6010/5981: 0/400 при 161/162
|
||||
через здоровый узел на прогоне 13.
|
||||
|
||||
ДВА КЛАССА, а не один (проба прода 06.09.2026 09:25 UTC — одна карточка по узлам):
|
||||
* КАПЧА («Captcha - база объявлений ЦИАН», «Вы не робот?») — безусловный отказ
|
||||
площадки: `report_platform_ban` + `CianBlockedError`. Снимается за 1-2 часа (узел 14,
|
||||
час назад отдававший капчу, вернул настоящую карточку) — TTL бана по назначению;
|
||||
* «Ошибка - Циан» — ТОЛЬКО ЛОГ, прежний `None`. Природа не доказана: страница может
|
||||
быть транзиентной 5xx-заглушкой под кодом 200, а не отказом конкретному узлу (узел 1
|
||||
за час сменил её на «Вы не робот?»). `mark_banned` эскалирует TTL до часов — по этой
|
||||
догадке 20-минутный сбой Циана выбил бы из выдачи весь пул. Решение — по частоте в
|
||||
логах за цикл наблюдения.
|
||||
|
||||
Этот файл покрывает ВТОРОЙ слой (kit) — он нужен потому, что образы backend и browser
|
||||
деплоятся раздельно и расходятся на часы: в это окно сайдкар ещё отдаёт 200 с капчей.
|
||||
Первый слой (сайдкар: 403 + ban_page) — `browser/test_server_cian_captcha.py`.
|
||||
|
||||
Красное на main:
|
||||
* `fetch_detail` на капче возвращала None и НЕ звала `report_platform_ban`;
|
||||
* батч из двух объявлений давал 0 успехов — второй листинг шёл через тот же узел.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from scraper_kit.cian_exceptions import CianBlockedError
|
||||
from scraper_kit.providers.cian import detail as cian_detail
|
||||
|
||||
from app.tasks import cian_history_backfill
|
||||
|
||||
_URL = "https://ekb.cian.ru/sale/flat/327830237/"
|
||||
|
||||
# Фрагмент прод-страницы капчи: заголовок + то самое слово в теле.
|
||||
_CAPTCHA_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head><meta charset='utf-8'>"
|
||||
"<title>Captcha - база объявлений ЦИАН</title></head>"
|
||||
"<body><div id='captcha'></div>"
|
||||
"<script>window.__captcha__ = {sitekey: 'x'};</script></body></html>"
|
||||
)
|
||||
# Второй вариант капчи — тот же отказ, другая вёрстка (проба 06.09.2026, узел 1, 16 КБ).
|
||||
_ROBOT_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head><title>Вы не робот?</title></head>"
|
||||
"<body><div id='captcha-container'></div></body></html>"
|
||||
)
|
||||
_ERROR_PAGE_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head><title>Ошибка - Циан</title></head>"
|
||||
"<body><h1>Что-то пошло не так</h1></body></html>"
|
||||
)
|
||||
# Здоровая карточка в текущем (`.concat`) формате _cianConfig. Слово `captcha` в теле
|
||||
# ЕСТЬ — в нормальной карточке Циана оно встречается 11 раз (антифрод-скрипты), и
|
||||
# именно поэтому признаком отказа оно быть не может.
|
||||
_CARD_HTML = (
|
||||
"<html><head><title>Купить 1-комн. квартиру — ЦИАН</title></head><body>"
|
||||
"<script>window.captchaSettings = {}; /* captcha */</script><script>"
|
||||
"window._cianConfig['frontend-offer-card'] = "
|
||||
"(window._cianConfig['frontend-offer-card'] || []).concat("
|
||||
'[{"key":"defaultState","value":{"offerData":{"offer":{"cianId":327830237}}},'
|
||||
'"priority":0}]);'
|
||||
"</script></body></html>"
|
||||
)
|
||||
# Страница без состояния и без маркеров отказа — дрейф разметки, прежнее поведение.
|
||||
_NO_STATE_HTML = "<html><head><title>Купить квартиру — ЦИАН</title></head><body></body></html>"
|
||||
|
||||
|
||||
class _FakeFetcher:
|
||||
"""BrowserFetcher-заглушка: капча, пока не сменится «аренда».
|
||||
|
||||
`report_platform_ban` зеркалит боевую механику ровно в том, что нас интересует:
|
||||
рапорт по текущему lease + смена аренды (в бою — `report_ban` + fail-streak →
|
||||
`_LEASE_ROTATE_AFTER_FAILS`, browser_fetcher.py). Здесь смена мгновенная — тест
|
||||
про то, ЗОВЁТСЯ ли рапорт и продолжается ли батч на новом узле, а не про порог.
|
||||
"""
|
||||
|
||||
def __init__(self, *, refusal_html: str = _CAPTCHA_HTML, **kwargs: Any) -> None:
|
||||
self.node = 1
|
||||
self.refusal_html = refusal_html
|
||||
self.last_response_status: int | None = 200
|
||||
self.ban_reports: list[str] = []
|
||||
self.fetched: list[str] = []
|
||||
|
||||
async def __aenter__(self) -> _FakeFetcher:
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *_: object) -> None:
|
||||
return None
|
||||
|
||||
async def fetch(self, url: str, **kwargs: Any) -> str:
|
||||
self.fetched.append(url)
|
||||
return self.refusal_html if self.node == 1 else _CARD_HTML
|
||||
|
||||
def report_platform_ban(self, reason: str) -> None:
|
||||
self.ban_reports.append(reason)
|
||||
self.node += 1
|
||||
|
||||
|
||||
# ── Слой kit: капча → CianBlockedError + рапорт бана ──────────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("html", [_CAPTCHA_HTML, _ROBOT_HTML])
|
||||
async def test_captcha_page_is_a_platform_refusal(html: str) -> None:
|
||||
"""HTTP 200 + капча (оба варианта вёрстки) → CianBlockedError, а не тихий None."""
|
||||
fetcher = _FakeFetcher(refusal_html=html)
|
||||
|
||||
with pytest.raises(CianBlockedError) as exc_info:
|
||||
await cian_detail.fetch_detail(_URL, browser_fetcher=fetcher)
|
||||
|
||||
assert "капча Циана" in str(exc_info.value)
|
||||
# Рапорт бана — по ЖИВОМУ lease, внутри `async with BrowserFetcher(...)` caller'а.
|
||||
assert len(fetcher.ban_reports) == 1
|
||||
assert "cian detail" in fetcher.ban_reports[0]
|
||||
|
||||
|
||||
async def test_error_page_is_logged_but_not_banned(caplog: pytest.LogCaptureFixture) -> None:
|
||||
"""«Ошибка - Циан» → прежний None + WARNING; `report_platform_ban` НЕ зовётся.
|
||||
|
||||
Цикл наблюдения (#3402): страница может быть транзиентным сбоем площадки, а бан
|
||||
эскалирует TTL до часов — за 20-минутный сбой Циана пул вылетел бы из выдачи.
|
||||
"""
|
||||
fetcher = _FakeFetcher(refusal_html=_ERROR_PAGE_HTML)
|
||||
|
||||
with caplog.at_level(logging.WARNING):
|
||||
result = await cian_detail.fetch_detail(_URL, browser_fetcher=fetcher)
|
||||
|
||||
assert result is None
|
||||
assert fetcher.ban_reports == []
|
||||
assert "страница ошибки Циана" in caplog.text
|
||||
assert "не бан, только лог" in caplog.text
|
||||
|
||||
|
||||
async def test_page_without_state_and_without_markers_still_returns_none() -> None:
|
||||
"""Дрейф разметки (нет состояния, нет маркеров) — прежнее поведение, не бан."""
|
||||
fetcher = _FakeFetcher(refusal_html=_NO_STATE_HTML)
|
||||
|
||||
result = await cian_detail.fetch_detail(_URL, browser_fetcher=fetcher)
|
||||
|
||||
assert result is None
|
||||
assert fetcher.ban_reports == []
|
||||
|
||||
|
||||
async def test_healthy_card_with_the_word_captcha_parses() -> None:
|
||||
"""Слово `captcha` в теле здоровой карточки баном НЕ считается."""
|
||||
fetcher = _FakeFetcher()
|
||||
fetcher.node = 2 # «здоровый узел» → отдаёт карточку
|
||||
|
||||
result = await cian_detail.fetch_detail(_URL, browser_fetcher=fetcher)
|
||||
|
||||
assert result is not None
|
||||
assert result.cian_id == 327830237
|
||||
assert fetcher.ban_reports == []
|
||||
|
||||
|
||||
# ── Батч: аренда меняется ВНУТРИ прогона, а не на следующем ───────────────────
|
||||
|
||||
|
||||
def _db_with_rows(n: int) -> MagicMock:
|
||||
db = MagicMock()
|
||||
db.execute.return_value.mappings.return_value.all.return_value = [
|
||||
{"id": i, "source_url": f"https://ekb.cian.ru/sale/flat/{i}/"} for i in range(1, n + 1)
|
||||
]
|
||||
return db
|
||||
|
||||
|
||||
async def test_batch_rotates_after_captcha_and_enriches_the_next_listing() -> None:
|
||||
"""Капча на первом объявлении → рапорт бана → второе обогащено на новом узле.
|
||||
|
||||
На main тест красный по значению: `fetch_detail` возвращал None, рапорта не было,
|
||||
«узел» не менялся — оба объявления шли через капчу, listings_succeeded == 0.
|
||||
"""
|
||||
fetcher = _FakeFetcher()
|
||||
|
||||
with (
|
||||
patch.object(cian_history_backfill, "BrowserFetcher", lambda **kw: fetcher),
|
||||
patch.object(cian_history_backfill, "save_detail_enrichment", MagicMock()),
|
||||
patch("asyncio.sleep", new_callable=AsyncMock),
|
||||
):
|
||||
result = await cian_history_backfill.backfill_cian_history(
|
||||
_db_with_rows(2), do_listings=True, do_houses=False, do_valuations=False
|
||||
)
|
||||
|
||||
assert fetcher.ban_reports, "капча не отрапортована пулу — узел доработает батч"
|
||||
assert len(fetcher.fetched) == 2 # батч не оборван: следующее объявление взято
|
||||
assert result.listings_succeeded == 1 # второе обогащено ПОСЛЕ смены аренды
|
||||
assert result.listings_failed_fetch == 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:
|
||||
"""Ответ на UPDATE ... RETURNING source (proxy_pool.clear_source_bans, #3404: гасит
|
||||
строку — banned_until=now()/ban_count=0/cleared_at/cleared_reason, — а не DELETE).
|
||||
"""Ответ на DELETE ... RETURNING source (proxy_pool.clear_source_bans, #2600 п.2).
|
||||
|
||||
PATCH enabled=true снимает баны узла по источникам — «ручное включение = чистый
|
||||
лист», как и обнуление 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})
|
||||
assert r.status_code == 200, r.text
|
||||
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
|
||||
delete_sql = str(db.execute.call_args_list[1].args[0])
|
||||
assert "DELETE FROM scrape_proxy_source_bans" in delete_sql
|
||||
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})
|
||||
assert r.status_code == 200, r.text
|
||||
# #3404: clear_source_bans теперь UPDATE, а не DELETE — ищем по cleared_reason,
|
||||
# уникальному для этого запроса маркеру.
|
||||
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
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -1963,69 +1963,6 @@ def _is_domclick_refusal(html: str) -> bool:
|
|||
return any(marker in lower for marker in _DOMCLICK_REFUSAL_MARKERS)
|
||||
|
||||
|
||||
# ── Отказ Циана с кодом 200 (#3402, замер прода 06.09.2026) ────────────────────
|
||||
# Циан отдаёт капчу и страницу ошибки с HTTP 200, поэтому _REFUSAL_STATUSES их не
|
||||
# видит, а маркеров Авито (_CHALLENGE_MARKERS/_BAN_MARKERS) в них нет. Отказ уезжал
|
||||
# наверх как валидный HTML: парсер не находил defaultState, прогон
|
||||
# cian_detail_backfill держал ОДНУ аренду на весь батч и сжигал через капча-узел
|
||||
# 210-400 карточек подряд (прогоны 6200: 0/210; 6123/6091/6052/6032/6010/5981: 0/400,
|
||||
# при 161/162 через здоровый узел на прогоне 13).
|
||||
#
|
||||
# Опознаём по <title>, а не по слову "captcha" в теле: оно встречается 11 раз и в
|
||||
# НОРМАЛЬНОЙ карточке (антифрод-скрипты), 17 — на странице капчи; как признак оно
|
||||
# неотличимо. Тире в заголовке нормализуется (— и – → -): вёрстка Циана печатает
|
||||
# его по-разному, а различать заголовки по виду дефиса — заведомо хрупко.
|
||||
#
|
||||
# КАПЧА — безусловный отказ площадки: страница требует пройти проверку, контента за
|
||||
# ней нет, а узел, которому её показали, будет показывать её и дальше. Наверх идёт
|
||||
# BanPageDetectedError → бан пары «узел×cian» + ротация аренды. Проба прода
|
||||
# 06.09.2026 09:25 UTC по узлам: капча снимается за 1-2 часа (узел 14 через час
|
||||
# отдавал уже настоящую карточку), то есть TTL бана по назначению.
|
||||
_CIAN_CAPTCHA_TITLES: tuple[str, ...] = (
|
||||
"captcha - база объявлений циан",
|
||||
"вы не робот?",
|
||||
)
|
||||
|
||||
# «Ошибка - Циан» — НЕ бан (#3402, цикл наблюдения). Природа страницы не доказана:
|
||||
# она может быть транзиентной 5xx-заглушкой, отданной с кодом 200, а не отказом
|
||||
# конкретному узлу (та же проба 06.09.2026: узел 1 за час сменил «Ошибка - Циан»,
|
||||
# 374 КБ, на «Вы не робот?», 16 КБ). Цена ошибочного бана несимметрична: mark_banned
|
||||
# эскалирует TTL до часов, и 20-минутный сбой Циана выбил бы из выдачи весь пул.
|
||||
# Поэтому здесь только WARNING — решение принимаем по частоте в логах, а не по
|
||||
# догадке о причине.
|
||||
_CIAN_ERROR_TITLES: tuple[str, ...] = ("ошибка - циан",)
|
||||
|
||||
_TITLE_RE = re.compile(r"<title[^>]*>(.*?)</title>", re.IGNORECASE | re.DOTALL)
|
||||
|
||||
|
||||
def _page_title(html: str) -> str:
|
||||
"""Текст <title> — схлопнутые пробелы, нижний регистр, нормализованное тире."""
|
||||
match = _TITLE_RE.search(html)
|
||||
if match is None:
|
||||
return ""
|
||||
return " ".join(match.group(1).split()).lower().replace("—", "-").replace("–", "-")
|
||||
|
||||
|
||||
def _is_cian_captcha(html: str) -> bool:
|
||||
"""True, если HTML — капча Циана (см. _CIAN_CAPTCHA_TITLES) — отказ площадки."""
|
||||
title = _page_title(html)
|
||||
return any(marker in title for marker in _CIAN_CAPTCHA_TITLES)
|
||||
|
||||
|
||||
def _log_cian_error_page(html: str, url: str, status: int | None) -> None:
|
||||
"""WARNING на «Ошибка - Циан» (см. _CIAN_ERROR_TITLES); ни бана, ни исключения."""
|
||||
title = _page_title(html)
|
||||
if not any(marker in title for marker in _CIAN_ERROR_TITLES):
|
||||
return
|
||||
logger.warning(
|
||||
"tradein-browser[cian]: страница ошибки Циана (title=%r, upstream=%s) — "
|
||||
"не бан, только лог (#3402) url=%r",
|
||||
title,
|
||||
status,
|
||||
url,
|
||||
)
|
||||
|
||||
|
||||
# Маркеры исключения playwright «страница прямо сейчас перезагружается». Ловим по
|
||||
# тексту, а не по типу: сервис не импортирует playwright напрямую (page приходит
|
||||
# уже готовым), а Error/TimeoutError у него не образуют отдельной иерархии для
|
||||
|
|
@ -2398,13 +2335,6 @@ async def _fetch_once(
|
|||
raise BanPageDetectedError(
|
||||
f"tradein-browser[{provider}]: статический отказ площадки url={url!r}"
|
||||
)
|
||||
if provider == "cian":
|
||||
if _is_cian_captcha(text):
|
||||
raise BanPageDetectedError(
|
||||
f"tradein-browser[{provider}]: капча Циана "
|
||||
f"(title={_page_title(text)!r}) url={url!r}"
|
||||
)
|
||||
_log_cian_error_page(text, url, status)
|
||||
logger.info(
|
||||
"tradein-browser[%s]: %s → HTTP %s, тело %d Б url=%r",
|
||||
provider,
|
||||
|
|
@ -2440,17 +2370,6 @@ async def _fetch_once(
|
|||
raise BanPageDetectedError(
|
||||
f"tradein-browser[{provider}]: бан-страница (проблема с IP) url={url!r}"
|
||||
)
|
||||
# Капча Циана (#3402): приходит с HTTP 200, поэтому ни _REFUSAL_STATUSES, ни
|
||||
# маркеры Авито её не ловят — распознаём по <title> и отдаём наверх тем же
|
||||
# путём (403 + ban_page), что отказ DomClick. «Ошибка - Циан» — НЕ бан, только
|
||||
# лог: см. _CIAN_ERROR_TITLES.
|
||||
if provider == "cian":
|
||||
if _is_cian_captcha(html):
|
||||
raise BanPageDetectedError(
|
||||
f"tradein-browser[{provider}]: капча Циана "
|
||||
f"(title={_page_title(html)!r}) url={url!r}"
|
||||
)
|
||||
_log_cian_error_page(html, url, _last_response_status.get(provider))
|
||||
if provider == "domclick":
|
||||
# DomClick — своя ветка (#3196): нет отдельного маркера самого
|
||||
# рукопожатия (см. комментарий у _DOMCLICK_SUCCESS_MARKER), поэтому
|
||||
|
|
|
|||
|
|
@ -1,307 +0,0 @@
|
|||
"""test_server_cian_captcha.py — капча Циана приходит с HTTP 200 (#3402).
|
||||
|
||||
Замер прода 06.09.2026: Циан отдаёт капчу (`<title>Captcha - база объявлений ЦИАН`,
|
||||
44 КБ; `<title>Вы не робот?`, 16 КБ) и страницу ошибки (`<title>Ошибка - Циан`, 374 КБ)
|
||||
с кодом **200**, поэтому _REFUSAL_STATUSES {403,429} их не видит, а маркеры Авито
|
||||
(_CHALLENGE_MARKERS / _BAN_MARKERS) в них не встречаются. HTML уезжал клиенту как успех,
|
||||
парсер не находил defaultState, и cian_detail_backfill держал ОДНУ аренду на весь батч,
|
||||
сжигая через капча-узел 210-400 карточек подряд (прогоны 6200: 0/210;
|
||||
6123/6091/6052/6032/6010/5981: 0/400 — против 161/162 через здоровый узел на прогоне 13).
|
||||
|
||||
Слово `captcha` как признак не годится: в НОРМАЛЬНОЙ карточке Циана оно встречается 11
|
||||
раз (антифрод-скрипты), на странице капчи — 17. Отсюда детект по <title>.
|
||||
|
||||
ДВА КЛАССА, а не один (проба прода 06.09.2026 09:25 UTC по узлам через сайдкар):
|
||||
* КАПЧА («Captcha - база объявлений ЦИАН», «Вы не робот?») — отказ площадки: бан пары
|
||||
«узел×cian» + ротация аренды. Снимается за 1-2 часа (узел 14 через час отдавал уже
|
||||
настоящую карточку), то есть TTL бана по назначению;
|
||||
* «Ошибка - Циан» — ТОЛЬКО ЛОГ. Природа не доказана: может быть транзиентной
|
||||
5xx-заглушкой под кодом 200, а не отказом узлу (узел 1 за час сменил её на «Вы не
|
||||
робот?»). Цена ошибки несимметрична — mark_banned эскалирует TTL до часов, и
|
||||
20-минутный сбой Циана выбил бы из выдачи весь пул. Решение — по частоте в логах.
|
||||
|
||||
camoufox НЕ запускается: _browsers[provider] — поддельный browser/page (зеркалит
|
||||
test_server_http_status.py). wait_for_timeout на фейковой page — no-op.
|
||||
|
||||
Запуск (из tradein-mvp/browser/)::
|
||||
|
||||
python -m pytest test_server_cian_captcha.py -q
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import importlib.util
|
||||
import json
|
||||
import logging
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
from aiohttp.test_utils import make_mocked_request
|
||||
|
||||
# server.py — не пакет (отдельный сервис без __init__/pyproject). Грузим по пути.
|
||||
_SERVER_PATH = Path(__file__).resolve().parent / "server.py"
|
||||
_spec = importlib.util.spec_from_file_location("tradein_browser_server", _SERVER_PATH)
|
||||
assert _spec is not None and _spec.loader is not None
|
||||
server = importlib.util.module_from_spec(_spec)
|
||||
_spec.loader.exec_module(server)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_state(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Чистое per-provider состояние на каждый тест (зеркалит соседние тесты)."""
|
||||
monkeypatch.setattr(server, "_browsers", {})
|
||||
monkeypatch.setattr(server, "_browser_cms", {})
|
||||
monkeypatch.setattr(server, "_page_counters", {})
|
||||
monkeypatch.setattr(server, "_locks", {})
|
||||
monkeypatch.setattr(server, "_retry_tasks", {})
|
||||
monkeypatch.setattr(server, "_last_goto_at", {})
|
||||
monkeypatch.setattr(server, "_last_response_status", {})
|
||||
monkeypatch.setattr(server, "_launched_proxy", {})
|
||||
monkeypatch.setattr(server, "_locks_guard", asyncio.Lock())
|
||||
monkeypatch.delenv("SCRAPER_PROXY_URL", raising=False)
|
||||
|
||||
|
||||
# Обёртка капчи — фрагмент прод-страницы: заголовок + то самое слово в теле.
|
||||
_CAPTCHA_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head>"
|
||||
"<meta charset='utf-8'>"
|
||||
"<title>Captcha - база объявлений ЦИАН</title>"
|
||||
"</head><body><div id='captcha'></div>"
|
||||
"<script>window.__captcha__ = {sitekey: 'x'};</script>"
|
||||
"</body></html>"
|
||||
)
|
||||
# Второй вариант капчи — тот же отказ, другая вёрстка (проба 06.09.2026, узел 1, 16 КБ).
|
||||
_ROBOT_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head><title>Вы не робот?</title></head>"
|
||||
"<body><div id='captcha-container'></div></body></html>"
|
||||
)
|
||||
_ERROR_PAGE_HTML = (
|
||||
"<!DOCTYPE html><html lang='ru'><head><title>Ошибка - Циан</title></head>"
|
||||
"<body><h1>Что-то пошло не так</h1></body></html>"
|
||||
)
|
||||
# Нормальная карточка: слово captcha в теле есть (антифрод), состояние — на месте.
|
||||
_CARD_HTML = (
|
||||
"<html><head><title>Купить 1-комн. квартиру — ЦИАН</title></head><body>"
|
||||
"<script>window.captchaConfig = {}; /* captcha captcha captcha */</script>"
|
||||
"<script>window._cianConfig['frontend-offer-card'] = [];</script>"
|
||||
"</body></html>"
|
||||
)
|
||||
|
||||
|
||||
class _Response:
|
||||
"""Поддельный playwright Response — интересует только .status."""
|
||||
|
||||
def __init__(self, status: int) -> None:
|
||||
self.status = status
|
||||
|
||||
|
||||
class _Page:
|
||||
"""Поддельная page: goto отдаёт Response(200), content() — заданный HTML."""
|
||||
|
||||
def __init__(self, html: str, status: int = 200) -> None:
|
||||
self._html = html
|
||||
self._status = status
|
||||
self.goto_urls: list[str] = []
|
||||
self.closed = 0
|
||||
|
||||
async def route(self, pattern: str, handler: Any) -> None:
|
||||
return None
|
||||
|
||||
async def goto(self, url: str, **kwargs: Any) -> _Response:
|
||||
self.goto_urls.append(url)
|
||||
return _Response(self._status)
|
||||
|
||||
async def wait_for_timeout(self, ms: int) -> None:
|
||||
return None
|
||||
|
||||
async def content(self) -> str:
|
||||
return self._html
|
||||
|
||||
async def close(self) -> None:
|
||||
self.closed += 1
|
||||
|
||||
|
||||
class _Browser:
|
||||
def __init__(self, page: _Page) -> None:
|
||||
self._page = page
|
||||
|
||||
async def new_page(self) -> _Page:
|
||||
return self._page
|
||||
|
||||
|
||||
def _install(monkeypatch: pytest.MonkeyPatch, page: _Page, provider: str = "cian") -> None:
|
||||
server._browsers[provider] = _Browser(page)
|
||||
monkeypatch.setattr(
|
||||
server, "_RECYCLE_PAGES_BY_PROVIDER", dict.fromkeys(server.PROVIDERS, 10_000)
|
||||
)
|
||||
monkeypatch.setattr(server, "BROWSER_WAIT_MS", 0)
|
||||
monkeypatch.setattr(server, "_MIN_PAGE_INTERVAL_BY_PROVIDER", {})
|
||||
monkeypatch.setattr(server, "BROWSER_MIN_PAGE_INTERVAL_S", 0.0)
|
||||
|
||||
|
||||
def _json_body(response: Any) -> dict[str, Any]:
|
||||
return json.loads(response.body.decode())
|
||||
|
||||
|
||||
async def _coro(value: Any) -> Any:
|
||||
return value
|
||||
|
||||
|
||||
def _make_request(body: dict[str, Any]) -> Any:
|
||||
request = make_mocked_request("POST", "/fetch")
|
||||
request.json = lambda: _coro(body) # type: ignore[method-assign]
|
||||
return request
|
||||
|
||||
|
||||
# ── Детектор: что капча, что страница ошибки, что ни то ни другое ─────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("html", [_CAPTCHA_HTML, _ROBOT_HTML])
|
||||
def test_cian_captcha_pages_are_recognised(html: str) -> None:
|
||||
"""Оба варианта капчи — отказ площадки (второй, «Вы не робот?», добавлен 06.09)."""
|
||||
assert server._is_cian_captcha(html) is True
|
||||
|
||||
|
||||
def test_error_page_is_not_a_captcha() -> None:
|
||||
"""«Ошибка - Циан» баном НЕ считается — иначе 20-минутный сбой Циана выбивает пул."""
|
||||
assert server._is_cian_captcha(_ERROR_PAGE_HTML) is False
|
||||
|
||||
|
||||
def test_normal_card_with_the_word_captcha_is_not_a_refusal() -> None:
|
||||
"""Слово `captcha` в теле нормальной карточки признаком отказа НЕ является."""
|
||||
assert "captcha" in _CARD_HTML.lower()
|
||||
assert server._is_cian_captcha(_CARD_HTML) is False
|
||||
|
||||
|
||||
def test_em_dash_in_title_is_normalised(caplog: pytest.LogCaptureFixture) -> None:
|
||||
"""Вёрстка печатает тире по-разному — детект не должен зависеть от его вида."""
|
||||
with caplog.at_level(logging.WARNING):
|
||||
server._log_cian_error_page("<title>Ошибка — Циан</title>", "https://x", 200)
|
||||
|
||||
assert "страница ошибки Циана" in caplog.text
|
||||
|
||||
|
||||
# ── /fetch: 403 + ban_page на капче, отданной с HTTP 200 ──────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("html", [_CAPTCHA_HTML, _ROBOT_HTML])
|
||||
def test_fetch_once_raises_on_cian_captcha(monkeypatch: pytest.MonkeyPatch, html: str) -> None:
|
||||
"""HTTP 200 + капча → BanPageDetectedError, а не «валидный HTML» наверх."""
|
||||
page = _Page(html)
|
||||
_install(monkeypatch, page)
|
||||
|
||||
with pytest.raises(server.BanPageDetectedError):
|
||||
asyncio.run(server._fetch_once("cian", "https://ekb.cian.ru/sale/flat/1/"))
|
||||
|
||||
|
||||
def test_fetch_once_logs_error_page_without_banning(
|
||||
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
|
||||
) -> None:
|
||||
"""«Ошибка - Циан» → HTML уезжает наверх как есть + WARNING; исключения НЕТ.
|
||||
|
||||
Пока идёт цикл наблюдения (#3402): страница может быть транзиентным сбоем площадки,
|
||||
а `mark_banned` эскалирует TTL до часов — бан по догадке дороже пропущенного отказа.
|
||||
"""
|
||||
page = _Page(_ERROR_PAGE_HTML)
|
||||
_install(monkeypatch, page)
|
||||
|
||||
with caplog.at_level(logging.WARNING):
|
||||
html = asyncio.run(server._fetch_once("cian", "https://ekb.cian.ru/sale/flat/1/"))
|
||||
|
||||
assert html == _ERROR_PAGE_HTML
|
||||
assert "страница ошибки Циана" in caplog.text
|
||||
assert "не бан, только лог" in caplog.text
|
||||
|
||||
|
||||
def test_fetch_once_passes_normal_card_through(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Здоровая карточка (с тем же словом в теле) отдаётся как раньше."""
|
||||
page = _Page(_CARD_HTML)
|
||||
_install(monkeypatch, page)
|
||||
|
||||
html = asyncio.run(server._fetch_once("cian", "https://ekb.cian.ru/sale/flat/1/"))
|
||||
|
||||
assert html == _CARD_HTML
|
||||
|
||||
|
||||
@pytest.mark.parametrize("html", [_CAPTCHA_HTML, _ROBOT_HTML])
|
||||
def test_fetch_handler_returns_403_with_ban_page_on_captcha(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
html: str,
|
||||
) -> None:
|
||||
"""Тот же путь, что #3379/#3288 п.4: 403 + ban_page + ЧЕСТНЫЙ upstream-статус 200.
|
||||
|
||||
Статус передаём как есть: апстрим ответил 200, и врать про 403 площадки нельзя —
|
||||
клиент опознаёт бан по признаку `ban_page`, а не по коду (образы сайдкара и
|
||||
бэкенда деплоятся врозь).
|
||||
"""
|
||||
monkeypatch.setattr(server, "IS_PROD", False)
|
||||
page = _Page(html)
|
||||
_install(monkeypatch, page)
|
||||
|
||||
async def _ensure(provider: str, proxy_override: str | None = None) -> bool:
|
||||
return True
|
||||
|
||||
monkeypatch.setattr(server, "_ensure_browser", _ensure)
|
||||
|
||||
response = asyncio.run(
|
||||
server.fetch_handler(
|
||||
_make_request({"url": "https://ekb.cian.ru/sale/flat/1/", "source": "cian"})
|
||||
)
|
||||
)
|
||||
|
||||
body = _json_body(response)
|
||||
assert response.status == 403
|
||||
assert body["ban_page"] is True
|
||||
assert body["status"] == 200
|
||||
assert "BanPageDetectedError" in body["error"]
|
||||
|
||||
|
||||
def test_fetch_handler_returns_200_without_ban_page_on_error_page(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""«Ошибка - Циан» доезжает клиенту как обычный ответ: 200, без `ban_page`."""
|
||||
monkeypatch.setattr(server, "IS_PROD", False)
|
||||
page = _Page(_ERROR_PAGE_HTML)
|
||||
_install(monkeypatch, page)
|
||||
|
||||
async def _ensure(provider: str, proxy_override: str | None = None) -> bool:
|
||||
return True
|
||||
|
||||
monkeypatch.setattr(server, "_ensure_browser", _ensure)
|
||||
|
||||
with caplog.at_level(logging.WARNING):
|
||||
response = asyncio.run(
|
||||
server.fetch_handler(
|
||||
_make_request({"url": "https://ekb.cian.ru/sale/flat/1/", "source": "cian"})
|
||||
)
|
||||
)
|
||||
|
||||
body = _json_body(response)
|
||||
assert response.status == 200
|
||||
assert "ban_page" not in body
|
||||
assert body["html"] == _ERROR_PAGE_HTML
|
||||
assert "не бан, только лог" in caplog.text
|
||||
|
||||
|
||||
def test_fetch_handler_does_not_ban_other_providers_on_the_same_html(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""Детект per-provider: та же страница у Авито проходит прежним путём (нет ложных банов)."""
|
||||
monkeypatch.setattr(server, "IS_PROD", False)
|
||||
page = _Page(_CAPTCHA_HTML)
|
||||
_install(monkeypatch, page, provider="avito")
|
||||
|
||||
async def _ensure(provider: str, proxy_override: str | None = None) -> bool:
|
||||
return True
|
||||
|
||||
monkeypatch.setattr(server, "_ensure_browser", _ensure)
|
||||
|
||||
response = asyncio.run(
|
||||
server.fetch_handler(_make_request({"url": "https://www.avito.ru/x", "source": "avito"}))
|
||||
)
|
||||
|
||||
assert response.status == 200
|
||||
assert _json_body(response)["html"] == _CAPTCHA_HTML
|
||||
|
|
@ -823,7 +823,7 @@ class BrowserFetcher:
|
|||
шаг — так отказ ПЛОЩАДКИ (подтверждённая бан-страница) не копит глобальный
|
||||
счётчик здоровья узла: он исправен, его отбил конкретный источник, и его
|
||||
судьбу решает `mark_banned(source=...)` по паре «узел×источник»
|
||||
(см. `report_platform_ban`). До #3288 узел с тремя бан-страницами Авито
|
||||
(см. `_report_platform_ban`). До #3288 узел с тремя бан-страницами Авито
|
||||
уходил из выдачи ВСЕМ источникам — прод-замер 31.08: yandex/cian/domclick
|
||||
получали ProxyPoolExhaustedError при banned_for_source=0 и трёх живых узлах;
|
||||
- ok=False копит `_lease_fail_streak`; после `_LEASE_ROTATE_AFTER_FAILS`
|
||||
|
|
@ -883,13 +883,8 @@ class BrowserFetcher:
|
|||
)
|
||||
self._lease = self._acquire_lease()
|
||||
|
||||
def report_platform_ban(self, reason: str) -> None:
|
||||
"""Исход /fetch, который опознан как бан-страница площадки (#3288).
|
||||
|
||||
Публичный метод (#3402): тем же путём обязан идти отказ, распознанный не
|
||||
сайдкаром, а провайдером — по телу ответа с HTTP 200 (капча Циана). Разница с
|
||||
`report_ban` существенна: тот только пишет бан пары «узел×источник», а сменить
|
||||
сожжённую аренду ВНУТРИ батча позволяет только fail-streak ниже.
|
||||
def _report_platform_ban(self, reason: str) -> None:
|
||||
"""Исход /fetch, который сайдкар опознал как бан-страницу площадки (#3288).
|
||||
|
||||
Отличается от обычного провала РОВНО одним: узел не получает `mark_health(False)`.
|
||||
Бан — приговор паре «узел×источник» (`mark_banned`, строка в
|
||||
|
|
@ -987,7 +982,9 @@ class BrowserFetcher:
|
|||
# httpx.HTTPStatusError, и общий `except Exception` ниже забирал её себе,
|
||||
# отправляя подтверждённый отказ ПЛОЩАДКИ в глобальный счётчик здоровья узла.
|
||||
self.last_response_status = None
|
||||
self.report_platform_ban(f"sidecar ban page (upstream={exc.upstream_status}) for {url}")
|
||||
self._report_platform_ban(
|
||||
f"sidecar ban page (upstream={exc.upstream_status}) for {url}"
|
||||
)
|
||||
raise
|
||||
except Exception:
|
||||
# NoProxyAvailableError (пустой пул) сюда НЕ приходит: он поднимается в
|
||||
|
|
@ -1044,7 +1041,7 @@ class BrowserFetcher:
|
|||
except SidecarBanPageError as exc:
|
||||
# Ветка ДО общего except по той же причине, что в _post_fetch (#3288):
|
||||
# бан-страница — отказ площадки, а не отказ узла.
|
||||
self.report_platform_ban(
|
||||
self._report_platform_ban(
|
||||
f"sidecar ban page on fetch-json (upstream={exc.upstream_status}) for {url}"
|
||||
)
|
||||
raise
|
||||
|
|
|
|||
|
|
@ -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.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 присутствует
|
||||
|
|
@ -623,12 +621,6 @@ 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(
|
||||
|
|
@ -645,9 +637,7 @@ 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"
|
||||
run_id = int(row.id)
|
||||
current_run_id.set(run_id)
|
||||
return run_id
|
||||
return int(row.id)
|
||||
|
||||
|
||||
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:
|
||||
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,
|
||||
# после коммита — статус уже персистирован в БД.
|
||||
|
|
@ -936,8 +923,6 @@ 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)
|
||||
|
||||
|
|
@ -1003,8 +988,6 @@ 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)
|
||||
|
||||
|
|
@ -1190,9 +1173,6 @@ 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
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -585,7 +585,7 @@ async def fetch_detail(
|
|||
# DomClick (#3239/#3283, providers/domclick/detail.py).
|
||||
#
|
||||
# report_ban ЗДЕСЬ НЕТ намеренно (#3288): бан рапортует сам fetcher, в
|
||||
# report_platform_ban, — раньше и по ПРАВИЛЬНОМУ lease. К моменту этого
|
||||
# _report_platform_ban, — раньше и по ПРАВИЛЬНОМУ lease. К моменту этого
|
||||
# кадра fetcher мог уже сротироваться на свежий узел по fail-streak, и
|
||||
# повторный рапорт отсюда банил бы того, кто к площадке не ходил.
|
||||
raise AvitoBlockedError(
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ from typing import TYPE_CHECKING, Any
|
|||
from sqlalchemy import text
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from scraper_kit.browser_fetcher import SidecarBanPageError
|
||||
from scraper_kit.ceiling_height import plausible_ceiling_m
|
||||
from scraper_kit.cian_exceptions import CianBlockedError
|
||||
from scraper_kit.cian_state_parser import extract_all_states, extract_state
|
||||
|
|
@ -100,44 +99,6 @@ def _raise_if_blocked(offer_url: str, status_code: int) -> None:
|
|||
raise CianBlockedError(f"Cian detail {offer_url} → HTTP 403 (WAF-блок узла)")
|
||||
|
||||
|
||||
# ── Капча/страница ошибки Циана с HTTP 200 (#3402) ──────────────────────────────
|
||||
# Дубль детектора сайдкара (browser/server.py::_is_cian_captcha) — намеренный:
|
||||
# образы backend и browser деплоятся раздельно и расходятся на часы, а в этот
|
||||
# промежуток отказ площадки приходит сюда ровно так же, как до правки — 200 + HTML
|
||||
# капчи. Опознаём по <title>, а не по слову "captcha": оно есть и в нормальной
|
||||
# карточке (11 вхождений против 17 на капче), то есть как признак не различает.
|
||||
#
|
||||
# КАПЧА — безусловный отказ площадки (бан пары «узел×cian» + ротация аренды): за ней
|
||||
# нет контента, и узел, которому её показали, будет получать её дальше. Проба прода
|
||||
# 06.09.2026 09:25 UTC по узлам: капча снимается за 1-2 часа, то есть TTL бана по
|
||||
# назначению.
|
||||
_CAPTCHA_TITLES: tuple[str, ...] = (
|
||||
"captcha - база объявлений циан",
|
||||
"вы не робот?",
|
||||
)
|
||||
|
||||
# «Ошибка - Циан» — НЕ бан (цикл наблюдения): природа страницы не доказана, она может
|
||||
# быть транзиентной 5xx-заглушкой Циана под кодом 200, а не отказом конкретному узлу.
|
||||
# Цена ошибки несимметрична — mark_banned эскалирует TTL до часов, и 20-минутный сбой
|
||||
# площадки выбил бы из выдачи весь пул. Здесь только WARNING и прежний None; решение
|
||||
# принимаем по частоте в логах.
|
||||
_ERROR_PAGE_TITLES: tuple[str, ...] = ("ошибка - циан",)
|
||||
|
||||
_TITLE_RE = re.compile(r"<title[^>]*>(.*?)</title>", re.IGNORECASE | re.DOTALL)
|
||||
|
||||
|
||||
def _page_title(html: str) -> str:
|
||||
"""Текст <title>: схлопнутые пробелы, нижний регистр, нормализованное тире.
|
||||
|
||||
Тире нормализуется (— и – → -): вёрстка печатает его по-разному, а различать
|
||||
заголовки по виду дефиса — заведомо хрупко.
|
||||
"""
|
||||
match = _TITLE_RE.search(html)
|
||||
if match is None:
|
||||
return ""
|
||||
return " ".join(match.group(1).split()).lower().replace("—", "-").replace("–", "-")
|
||||
|
||||
|
||||
async def fetch_detail(
|
||||
offer_url: str,
|
||||
*,
|
||||
|
|
@ -166,12 +127,7 @@ async def fetch_detail(
|
|||
CianBlockedError: HTTP 403 на curl-путях — WAF Циана отбил узел, с которого мы
|
||||
пришли (#2700). Оба вызывающих в orchestration/pipeline.py уже считают
|
||||
исключение в `errors_count`, а на own-session-пути оно дополнительно снимает
|
||||
узел с выдачи Циану через `curl_proxy_url`. Тем же исключением приезжает
|
||||
КАПЧА Циана (#3402): её сайдкар отдаёт как `ban_page` + 403, а если образ
|
||||
сайдкара старее — она распознаётся здесь, по <title>, уже после HTTP 200
|
||||
(`report_platform_ban` на живом lease + raise вместо тихого None). Страница
|
||||
«Ошибка - Циан» баном НЕ считается — WARNING и прежний None, см.
|
||||
`_ERROR_PAGE_TITLES`.
|
||||
узел с выдачи Циану через `curl_proxy_url`.
|
||||
NoProxyAvailableError: пул прокси пуст (#2616) — пробрасывается со ВСЕХ путей, а
|
||||
не гасится в None: запрос не уходил, и следующий вызов упрётся в то же самое,
|
||||
поэтому решение «оборвать батч» принимает вызывающий (#3197).
|
||||
|
|
@ -180,18 +136,6 @@ async def fetch_detail(
|
|||
# Browser path: get fully JS-rendered HTML; same parse path follows.
|
||||
try:
|
||||
html = await browser_fetcher.fetch(offer_url)
|
||||
except SidecarBanPageError as exc:
|
||||
# #3402: сайдкар опознал отказ площадки по маркерам тела (капча Циана,
|
||||
# `ban_page` + 403). ПОРЯДОК ВЕТОК ВАЖЕН — SidecarBanPageError
|
||||
# подкласс httpx.HTTPStatusError, и общий `except` ниже увёл бы
|
||||
# подтверждённый отказ Циана в `return None`, то есть в «не смогли
|
||||
# разобрать». Зеркалит avito/detail.py и domclick/detail.py (#3283/#3239).
|
||||
#
|
||||
# report_ban ЗДЕСЬ НЕТ намеренно (#3288): бан рапортует сам fetcher, в
|
||||
# report_platform_ban, — раньше и по ПРАВИЛЬНОМУ lease.
|
||||
raise CianBlockedError(
|
||||
f"Cian detail: сайдкар опознал отказ площадки для {offer_url}: {exc}"
|
||||
) from exc
|
||||
except Exception as exc:
|
||||
# «Пул пуст» — не отказ страницы: запрос не уходил вовсе, и вызывающий обязан
|
||||
# оборвать батч (#3197). Проглоченный здесь, он приезжал наверх как None, то
|
||||
|
|
@ -245,32 +189,6 @@ async def fetch_detail(
|
|||
# NOTE: detail pages use 'defaultState', SERP uses 'initialState'
|
||||
offer_state = extract_state(html, mfe="frontend-offer-card", key="defaultState")
|
||||
if offer_state is None:
|
||||
# #3402: страницы БЕЗ состояния бывают трёх разных родов, и до этой правки все
|
||||
# печатались как «extraction failed» — то есть отказ площадки читался как дрейф
|
||||
# нашей разметки. Капча приходит с HTTP 200, диагностировать её по статусу
|
||||
# (ban_kind_from_status) нечем — только по телу. Три рода: капча (бан+ротация),
|
||||
# страница ошибки (только лог, природа не доказана), дрейф разметки (как было).
|
||||
title = _page_title(html)
|
||||
if any(marker in title for marker in _CAPTCHA_TITLES):
|
||||
if browser_fetcher is not None:
|
||||
# Детект НАШ, сайдкар его не видел (старый образ отдал 200) — рапорт
|
||||
# обязателен здесь, иначе тот же узел доработает батч до конца. lease
|
||||
# ещё жив: caller держит `async with BrowserFetcher(...)` снаружи.
|
||||
# report_platform_ban (а не report_ban): бан пары «узел×cian» ПЛЮС
|
||||
# fail-streak, по которому аренда меняется ВНУТРИ батча — ровно то же,
|
||||
# что делает сайдкар-путь через SidecarBanPageError.
|
||||
browser_fetcher.report_platform_ban(f"cian detail: капча Циана для {offer_url}")
|
||||
logger.warning("Cian detail %s: капча Циана (title=%r)", offer_url, title)
|
||||
raise CianBlockedError(f"Cian detail: капча Циана (title={title!r}) для {offer_url}")
|
||||
if any(marker in title for marker in _ERROR_PAGE_TITLES):
|
||||
# Ни рапорта, ни исключения — см. _ERROR_PAGE_TITLES: пока это наблюдение,
|
||||
# а не диагноз. Возврат прежний (None) — вызывающий считает «не разобрали».
|
||||
logger.warning(
|
||||
"Cian detail %s: страница ошибки Циана (title=%r) — не бан, только лог (#3402)",
|
||||
offer_url,
|
||||
title,
|
||||
)
|
||||
return None
|
||||
logger.warning("Cian detail %s: defaultState extraction failed", offer_url)
|
||||
return None
|
||||
|
||||
|
|
|
|||
|
|
@ -639,7 +639,7 @@ async def fetch_detail(
|
|||
# вместо 'platform' и ротация IP не запустилась бы вовсе.
|
||||
#
|
||||
# report_ban ЗДЕСЬ НЕТ намеренно (#3288): его делает сам fetcher в
|
||||
# report_platform_ban — раньше и по ПРАВИЛЬНОМУ lease (после ротации по
|
||||
# _report_platform_ban — раньше и по ПРАВИЛЬНОМУ lease (после ротации по
|
||||
# fail-streak этот кадр забанил бы уже СЛЕДУЮЩИЙ узел). Ветка
|
||||
# parse_detail_html ниже — другой случай: там детект НАШ, fetcher его не
|
||||
# видит, и report_ban остаётся.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue