Merge pull request 'feat(tradein/proxy): прогон знает свой узел, а снятый бан перестаёт стирать историю (#3404)' (#3405) from feat/3404-proxy-run-attribution into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m14s
Deploy Trade-In / build-backend (push) Successful in 1m39s
Deploy Trade-In / deploy (push) Successful in 1m53s
Deploy Trade-In / deploy-status (push) Successful in 2s
Deploy Trade-In / perimeter-smoke (push) Successful in 11s

This commit is contained in:
lekss361 2026-09-06 10:18:14 +00:00
commit 72e3bc9e24
11 changed files with 832 additions and 46 deletions

View file

@ -1187,7 +1187,8 @@ class Settings(BaseSettings):
# #3283g: ротация exit-IP НА САМ БАН площадки, а не только по счётчику попыток.
# Бан привязан к IP (замерено вживую: rotate_proxy() лечит забаненный узел за
# секунды, clear_source_bans снимает запись из scrape_proxy_source_bans), но
# секунды, clear_source_bans гасит бан в scrape_proxy_source_bans — строка живёт
# до purge, #3404), но
# #3251/#3212 запрещают сбрасывать browser-context на КАЖДЫЙ блок -- сброс без
# смены IP выбрасывает пройденный QRATOR-PoW и запускает самоподдерживающийся
# каскад блоков на том же адресе. rotate_on_ban МЕНЯЕТ IP вместе со сбросом,

View file

@ -54,7 +54,10 @@ Self-healing (#2600):
WARNING (пул надо пополнять, #2638).
- Ручное снятие `clear_source_bans` (ложный бан детектора капчи, #2642) плюс
автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился
IP, поэтому смена адреса делает строку недействительной.
IP, поэтому смена адреса делает строку недействительной. С #3404 «снятие» гасит
строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`),
а не удаляет её строка живёт до штатного purge (SOURCE_BAN_PURGE_DAYS), но для
выдачи и для эскалации следующего бана это неотличимо от прежнего DELETE.
Ручное выключение vs авто-выключение (#2610):
- scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины
@ -143,6 +146,7 @@ __all__ = [
"STALE_LEASE_MINUTES",
"ProxyLease",
"acquire",
"attribute_run_proxy",
"clear_source_bans",
"mark_banned",
"mark_browser_health",
@ -441,6 +445,11 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
{"run_id": lease_marker, "id": proxy_id},
)
db.commit()
if run_id is not None and run_id != NON_RUN_LEASE_MARKER:
# #3404: одна точка, покрывающая ВСЕ пути выдачи (curl — acquire на каждый
# вызов, браузер — sticky lease на весь прогон, ре-acquire при ротации узла
# mid-run) — см. attribute_run_proxy docstring.
attribute_run_proxy(db, run_id, proxy_id)
if fallback_used:
logger.warning(
"proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity "
@ -474,6 +483,78 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
)
def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None:
"""Записать узел, через который идёт прогон run_id, в scrape_runs (#3404).
Единственный писатель `acquire()` сразу после выдачи lease'а: покрывает и
curl-путь (acquire на каждый вызов), и браузерный sticky lease (один acquire на
весь прогон), и ре-acquire при ротации узла mid-run (`browser_fetcher`,
`_LEASE_ROTATE_AFTER_FAILS`) то есть смена узла ЗА прогон фиксируется сама,
без отдельного вызова с чьей-либо стороны.
`scrape_runs.proxy_id` ПОСЛЕДНИЙ использованный узел (перезаписывается при
каждой новой выдаче); полная цепочка узлов, если она менялась, в
`counters.proxy_ids` (список id, без дублей). Пишем через `||`-мерж
`counters` (тот же контракт, что у `runs.update_heartbeat`/`mark_done`)
чужие ключи (чекпоинт, метка interrupted) не затираются.
Идемпотентно: повторная выдача ТОГО ЖЕ узла не дублирует его в `proxy_ids`
(`@>`-проверка перед append). Best-effort: любой сбой (например, run_id уже
не существует гонка с финализацией) логируется WARNING и проглатывается
атрибуция прогону не должна ронять выдачу прокси, это диагностика, а не
часть контракта lease'а. 0 rows (run_id не найден) — DEBUG, не ошибка: сама
выдача при этом уже произошла и коммитнута предыдущим db.commit() в acquire().
"""
try:
# Строку прогона параллельно обновляет heartbeat/финализатор из ДРУГОЙ сессии
# (короткие транзакции, каждая со своим commit). Пересечение маловероятно, но
# ждать на блокировке в пути выдачи прокси нельзя — диагностика не должна
# тормозить сбор. Не дождались за 2с — уходим в except ниже (WARNING, lease цел).
db.execute(text("SET LOCAL lock_timeout = '2s'"))
row = db.execute(
text(
"""
UPDATE scrape_runs
SET proxy_id = CAST(:proxy_id AS bigint),
counters = COALESCE(counters, '{}'::jsonb) || jsonb_build_object(
'proxy_ids',
CASE
WHEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
@> to_jsonb(CAST(:proxy_id AS bigint))
THEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
ELSE COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
|| jsonb_build_array(CAST(:proxy_id AS bigint))
END
)
WHERE id = CAST(:run_id AS bigint)
RETURNING id
"""
),
{"proxy_id": proxy_id, "run_id": run_id},
).first()
db.commit()
if row is None:
logger.debug(
"proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already "
"finalized?)",
run_id,
)
except Exception:
# Best-effort (см. docstring) — атрибуция диагностическая, не часть
# контракта lease'а: lease уже выдан и не должен теряться из-за неё.
logger.warning(
"proxy_pool: attribute_run_proxy failed run_id=%d proxy_id=%d — lease "
"issued regardless",
run_id,
proxy_id,
exc_info=True,
)
try:
db.rollback()
except Exception:
pass
def release(db: Session, proxy_id: int) -> None:
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
db.execute(
@ -935,15 +1016,21 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None =
)
)
ON CONFLICT (proxy_id, source) DO UPDATE
SET ban_count = scrape_proxy_source_bans.ban_count + 1,
banned_until = now() + make_interval(hours => CAST(
SET ban_count = scrape_proxy_source_bans.ban_count + 1,
banned_until = now() + make_interval(hours => CAST(
LEAST(
CAST(:base_hours AS integer)
* power(2, LEAST(scrape_proxy_source_bans.ban_count, 16)),
CAST(:max_hours AS integer)
) AS integer)),
reason = CAST(:reason AS text),
updated_at = now()
reason = CAST(:reason AS text),
-- #3404: строка могла быть погашена clear_source_bans (banned_until
-- в прошлом попадает в WHERE ниже) новый бан затирает её метки
-- гашения, иначе на СНОВА забаненной паре висели бы cleared_at/
-- cleared_reason от предыдущего, уже неактуального гашения.
cleared_at = NULL,
cleared_reason = NULL,
updated_at = now()
-- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет),
-- свою же (обычная эскалация) или перебиваем боевым сбором он сильнее
-- пробы. Иначе 0 rows и ветка "deferred" ниже (дефект #2803).
@ -1058,12 +1145,25 @@ def clear_source_bans(
банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса,
держа узел вне выдачи уже без причины.
source=None снять все баны узла; конкретный source только его. DELETE, а не
`banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации),
а здесь мы как раз объявляем историю недействительной новый бан начнётся с базовых
SOURCE_BAN_BASE_HOURS.
source=None снять все баны узла; конкретный source только его. Гасим строку
(`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`), а НЕ
удаляем (#3404, было DELETE): строка доживает до штатного purge'а
(`run_proxy_healthcheck`, SOURCE_BAN_PURGE_DAYS), но перестаёт блокировать
выдачу немедленно и перестаёт нести историю эскалации `ban_count = 0` даёт
следующему бану той же пары ровно те же SOURCE_BAN_BASE_HOURS, что и раньше
после DELETE (формула `mark_banned` берёт ПРЕДЫДУЩИЙ ban_count показателем
степени: 0 база, без множителя). Причина держать строку трассируемость
(видно, что бан БЫЛ и когда/кем снят), а не поведение: для читателей ниже
погашенная строка неотличима от отсутствующей (см. риски в шапке PR #3404).
`reason` идёт только в лог (человекочитаемый повод «manual enable», «ip rotated»).
Идемпотентно и в другую сторону: повторный вызов на уже погашенной строке
(последний предикат в WHERE) её не трогает 0 rows, `banned_until` НЕ
сдвигается вперёд. Без этого условия повторный `PATCH enabled=true` двигал бы
`banned_until` на каждый вызов и отодвигал бы purge на неопределённый срок.
`reason` идёт в лог (человекочитаемый повод «manual enable», «ip rotated») и
теперь ЕЩЁ в колонку `cleared_reason` постоянный след того, кто и почему
погасил бан.
`only_reason` ФИЛЬТР по колонке reason, т.е. «снимать только строки, которые
написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое
@ -1075,14 +1175,22 @@ def clear_source_bans(
rows = db.execute(
text(
"""
DELETE FROM scrape_proxy_source_bans
UPDATE scrape_proxy_source_bans
SET banned_until = now(),
ban_count = 0,
cleared_at = now(),
cleared_reason = CAST(:reason AS text),
updated_at = now()
WHERE proxy_id = CAST(:proxy_id AS bigint)
AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text))
AND (CAST(:only_reason AS text) IS NULL OR reason = CAST(:only_reason AS text))
-- Уже погашенная строка (гейт покоя, см. докстринг) не трогаем: без
-- него повторный вызов сдвигал бы banned_until вперёд и отодвигал purge.
AND NOT (cleared_at IS NOT NULL AND ban_count = 0 AND banned_until <= now())
RETURNING source
"""
),
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason},
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason, "reason": reason},
).fetchall()
db.commit()
if rows:
@ -1465,6 +1573,11 @@ async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
# Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с
# 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса:
# неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать».
# #3404: под этот же порог теперь попадают и ПОГАШЕННЫЕ clear_source_bans строки —
# для них banned_until == момент гашения (== cleared_at), т.е. таймер до purge
# отсчитывается от гашения, а не от исходного истечения бана. ban_count у них уже
# 0 к моменту гашения, так что покидающий purge их не «сбрасывает» повторно —
# он просто убирает уже неактуальный след из таблицы.
purged = len(
db.execute(
text(

View file

@ -221,10 +221,17 @@ class RealProxyProvider:
def acquire(self, provider: str) -> ProxyLease | None:
from scraper_kit.contracts import ProxyLease as _KitProxyLease
from scraper_kit.orchestration.run_context import current_run_id
# #3404: протокол ProxyProvider.acquire(provider) не несёт run_id (этот
# адаптер — один объект на весь scheduler_main.py), поэтому берём его из
# ContextVar, который выставляет runs.create_run. None — вызов вне прогона
# (health-check, эстиматор) — proxy_pool.acquire в этом случае лизит под
# NON_RUN_LEASE_MARKER, как и раньше, атрибуцию в scrape_runs не пишет.
run_id = current_run_id.get()
db = _SessionLocal()
try:
lease = _proxy_pool.acquire(db, provider)
lease = _proxy_pool.acquire(db, provider, run_id=run_id)
finally:
db.close()
if lease is None:

View file

@ -0,0 +1,107 @@
-- 287_proxy_run_attribution.sql
-- scrape_runs.proxy_id — узел прогона (#3404 A) + soft-clear банов по источнику (#3404 B).
--
-- Dependencies: 015_scrape_runs.sql (scrape_runs), 157_scrape_proxies.sql (scrape_proxies),
-- 210_scrape_proxy_source_bans.sql (scrape_proxy_source_bans).
-- Apply after: 286_offer_price_history_decimal_slips.sql
--
-- ЧАСТЬ A — WHY:
-- До сих пор ни одна строка scrape_runs не знала, через какой узел пула шёл прогон:
-- ProxyProvider.acquire(provider) run_id не принимает, а leased_by у боевого пути —
-- NON_RUN_LEASE_MARKER. Разбор исхода прогона по узлу (кто плодит баны/провалы)
-- был возможен только вручную, по времени. Код-часть (app/services/proxy_pool.py:
-- attribute_run_proxy) пишет сюда после каждой выдачи лиза; здесь только схема.
--
-- ЧАСТЬ A — WHAT:
-- proxy_id — узел, через который шёл прогон. Если за прогон узел МЕНЯЛСЯ (ротация
-- при повторных провалах в browser_fetcher, либо curl-путь берёт лиз на каждый вызов
-- в providers/_proxy.py), здесь остаётся ПОСЛЕДНИЙ; полная цепочка узлов копится в
-- scrape_runs.counters->'proxy_ids' (jsonb-массив, пишет тот же attribute_run_proxy).
-- ON DELETE SET NULL, а не CASCADE — узел из пула может быть выведен/удалён оператором,
-- история прогонов (аналитика, отчёты) не должна пропадать вместе с ним.
-- Индекс (proxy_id, started_at DESC) — под разрез «исход прогона по узлу за период»
-- (WHERE proxy_id = ... ORDER BY started_at DESC); partial по proxy_id IS NOT NULL не
-- делаем, потому что колонка сортировки (started_at) в самом индексе — Postgres и так
-- не будет использовать индекс без него для прогонов без узла.
--
-- НИЧЕГО НЕ БЭКФИЛЛИТСЯ: связать уже прошедшие прогоны с конкретным узлом задним
-- числом нечем — leased_by исторически = NON_RUN_LEASE_MARKER, а лог выдачи лизов
-- не хранит run_id. Для всех строк scrape_runs, созданных ДО этой миграции,
-- proxy_id остаётся NULL навсегда — это не «прогон без прокси», а «прогон, для
-- которого атрибуция не собиралась». Врать восстановленным/угаданным значением
-- нельзя, поэтому backfill-UPDATE здесь сознательно отсутствует.
--
-- ЧАСТЬ B — WHY:
-- clear_source_bans() (proxy_pool.py) сейчас делает DELETE строки
-- scrape_proxy_source_bans. Это стирает историю эскалации (ban_count) и не оставляет
-- следа, что бан был снят ДОСРОЧНО (оператором/успешной ротацией exit-IP), в отличие
-- от бана, который просто истёк сам. Код-часть переводит функцию на UPDATE
-- (гашение: banned_until=now(), ban_count=0, cleared_at/cleared_reason проставляются),
-- строка доживает до штатного purge в run_proxy_healthcheck (SOURCE_BAN_PURGE_DAYS).
-- Эскалация при повторном бане той же пары сохраняется 1:1: формула в mark_banned
-- берёт СТАРЫЙ ban_count как показатель степени (base * 2^ban_count), при
-- ban_count=0 после гашения это ровно SOURCE_BAN_BASE_HOURS=6ч — байт-в-байт как
-- свежий INSERT после DELETE.
--
-- ЧАСТЬ B — WHAT:
-- cleared_at — момент досрочного снятия бана (не путать с истечением banned_until
-- само по себе: NULL значит «бан снят не был / истёк сам», не-NULL — снят
-- оператором или ротацией exit-IP до истечения срока или сразу после).
-- cleared_reason — свободный текст причины снятия (тот же 'reason', что передаётся
-- в clear_source_bans).
--
-- ИДЕМПОТЕНТНОСТЬ:
-- ADD COLUMN IF NOT EXISTS × 3, CREATE INDEX IF NOT EXISTS, COMMENT ON COLUMN
-- (безусловны, но идемпотентны сами по себе — просто перезаписывают тот же текст).
-- Backfill-DML в файле нет вовсе, поэтому повторный прогон — чистый no-op.
BEGIN;
SET LOCAL lock_timeout = '5s';
-- ── Часть A: scrape_runs.proxy_id ────────────────────────────────────────────
ALTER TABLE scrape_runs
ADD COLUMN IF NOT EXISTS proxy_id bigint REFERENCES scrape_proxies(id) ON DELETE SET NULL;
CREATE INDEX IF NOT EXISTS idx_scrape_runs_proxy_id_started_at
ON scrape_runs (proxy_id, started_at DESC);
COMMENT ON COLUMN scrape_runs.proxy_id IS
'Узел пула (scrape_proxies.id), через который шёл прогон. NULL = прогон без '
'прокси (эстиматор, admin-инициированные вызовы с NON_RUN_LEASE_MARKER) либо '
'прогон ДО применения миграции 287 (backfill не делался — связать нечем). '
'Если узел менялся mid-run (ротация после серии провалов в browser_fetcher, '
'либо curl-путь берёт лиз заново на каждый вызов) — здесь ПОСЛЕДНИЙ выданный '
'узел, полная цепочка — counters->''proxy_ids'' (jsonb-массив id, в порядке '
'первой выдачи). ON DELETE SET NULL: удаление узла из пула не должно уносить '
'историю прогонов.';
-- ── Часть B: soft-clear в scrape_proxy_source_bans ───────────────────────────
ALTER TABLE scrape_proxy_source_bans
ADD COLUMN IF NOT EXISTS cleared_at timestamptz,
ADD COLUMN IF NOT EXISTS cleared_reason text;
COMMENT ON COLUMN scrape_proxy_source_bans.cleared_at IS
'Момент досрочного снятия бана (proxy_pool.clear_source_bans, #3404) — '
'оператором или успешной ротацией exit-IP. NULL = бан не снимался вручную '
'(либо ещё активен, либо истёк сам по banned_until). Строка при гашении НЕ '
'удаляется — доживает до штатного purge (SOURCE_BAN_PURGE_DAYS), таймер '
'которого для погашенных строк отсчитывается от banned_until = момент гашения.';
COMMENT ON COLUMN scrape_proxy_source_bans.cleared_reason IS
'Причина досрочного снятия бана (тот же текст, что передан в '
'clear_source_bans(reason=...)). NULL, если строка не гасилась вручную.';
COMMENT ON COLUMN scrape_proxy_source_bans.ban_count IS
'Сколько раз эта пара банилась. Срок ТЕКУЩЕГО бана (banned_until - banned_at) = '
'base * 2^(ban_count-1), потолок SOURCE_BAN_MAX_HOURS: ban_count=1 → 6ч, 2 → 12ч, '
'3 → 24ч и т.д. Сбрасывается либо purge''ем через SOURCE_BAN_PURGE_DAYS после '
'истечения, либо proxy_pool.clear_source_bans (#3404: досрочное ГАШЕНИЕ строки —'
' banned_until=now(), ban_count=0, cleared_at/cleared_reason проставляются; '
'строка НЕ удаляется, живёт до purge). Оба пути одинаково обнуляют ban_count, '
'поэтому следующий бан той же пары в обоих случаях стартует заново с '
'SOURCE_BAN_BASE_HOURS.';
COMMIT;

View file

@ -438,21 +438,38 @@ class FakeSession:
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
)
if "DELETE FROM scrape_proxy_source_bans" in sql and "proxy_id = CAST" in sql:
# clear_source_bans: снять баны узла (все либо один source), #2600 п.2.
# Фильтр по reason (#2800) гейтим по подстроке боевого SQL — как ban-фильтры
# в acquire-ветке: иначе мок «чинил» бы код, который фильтра не содержит, и
# тест на «успешная проба не гасит чужой бан» остался бы зелёным на сломанном.
if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql:
# clear_source_bans (#3404): гасим строку (banned_until=now(), ban_count=0,
# cleared_at/cleared_reason) вместо DELETE — трассируемость снятия бана, строка
# доживает до штатного purge. Фильтр по reason (#2800) гейтим по подстроке
# боевого SQL — как ban-фильтры в acquire-ветке: иначе мок «чинил» бы код,
# который фильтра не содержит, и тест на «успешная проба не гасит чужой бан»
# остался бы зелёным на сломанном.
filters_reason = "reason = CAST(:only_reason AS text)" in sql
only_reason = p.get("only_reason") if filters_reason else None
cleared = [
b
for b in self.bans
if b["proxy_id"] == p["proxy_id"]
and (p["source"] is None or b["source"] == p["source"])
and (only_reason is None or b.get("reason") == only_reason)
]
self.bans = [b for b in self.bans if b not in cleared]
now = datetime.now(UTC)
cleared: list[dict[str, Any]] = []
for b in self.bans:
if b["proxy_id"] != p["proxy_id"]:
continue
if p["source"] is not None and b["source"] != p["source"]:
continue
if only_reason is not None and b.get("reason") != only_reason:
continue
# Гейт покоя (реальный WHERE): уже погашенную строку повторно не трогаем —
# без него повторный вызов сдвигал бы banned_until вперёд.
if (
b.get("cleared_at") is not None
and b["ban_count"] == 0
and b["banned_until"] <= now
):
continue
b["banned_until"] = now
b["ban_count"] = 0
b["cleared_at"] = now
b["cleared_reason"] = p["reason"]
b["updated_at"] = now
cleared.append(b)
return _FakeResult([{"source": b["source"]} for b in cleared])
if "DELETE FROM scrape_proxy_source_bans" in sql: # purge истёкших (#2600 п.2)
@ -1419,6 +1436,10 @@ async def test_healthcheck_purges_long_expired_bans_only(
# До п.2 ложный бан лечился PATCH enabled=true (обнулял disabled_reason). Теперь бан
# в отдельной таблице и истекает только по таймеру (до 72ч при эскалации) — без этой
# ручки ложное срабатывание детектора капчи (#2642) снималось бы только руками в SQL.
#
# #3404: снятие гасит строку (banned_until=now(), ban_count=0, cleared_at/cleared_reason),
# а не удаляет её — строка живёт для трассируемости до штатного purge, но для выдачи и
# для эскалации следующего бана неотличима от прежнего DELETE.
def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
@ -1428,17 +1449,37 @@ def _active_ban(pid: int, source: str, ban_count: int = 1) -> dict[str, Any]:
"ban_count": ban_count,
"banned_until": datetime.now(UTC) + timedelta(hours=SOURCE_BAN_MAX_HOURS),
"reason": f"banned:{source}",
"cleared_at": None,
"cleared_reason": None,
}
def test_clear_source_bans_removes_all_bans_of_node() -> None:
def test_clear_source_bans_gates_all_bans_of_node() -> None:
db = FakeSession(
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],
bans=[_active_ban(1, "avito"), _active_ban(1, "cian"), _active_ban(2, "avito")],
)
cleared = proxy_pool.clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
assert cleared == 2
assert [(b["proxy_id"], b["source"]) for b in db.bans] == [(2, "avito")] # чужой цел
# #3404: строки НЕ удаляются — все три остаются в таблице (трассируемость).
assert {(b["proxy_id"], b["source"]) for b in db.bans} == {
(1, "avito"),
(1, "cian"),
(2, "avito"),
}
now = datetime.now(UTC)
node1_bans = [b for b in db.bans if b["proxy_id"] == 1]
assert len(node1_bans) == 2
for b in node1_bans:
assert b["banned_until"] <= now
assert b["ban_count"] == 0
assert b["cleared_at"] is not None
assert b["cleared_reason"] == "manual enable"
# чужой бан (proxy_id=2) не тронут — остаётся активным
other_ban = next(b for b in db.bans if b["proxy_id"] == 2)
assert other_ban["banned_until"] > now
assert other_ban["ban_count"] == 1
assert other_ban["cleared_at"] is None
# узел снова выдаётся источнику, который его банил
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
assert lease is not None and lease.id == 1
@ -1453,7 +1494,16 @@ def test_clear_source_bans_single_source_keeps_others() -> None:
db, 1, source="avito", reason="ip rotated"
)
assert cleared == 1
assert [b["source"] for b in db.bans] == ["cian"]
# обе строки остаются (#3404), но только "avito" погашена
avito_ban = db._ban(1, "avito")
cian_ban = db._ban(1, "cian")
now = datetime.now(UTC)
assert avito_ban["banned_until"] <= now
assert avito_ban["ban_count"] == 0
assert avito_ban["cleared_reason"] == "ip rotated"
assert cian_ban["banned_until"] > now
assert cian_ban["ban_count"] == 1
assert cian_ban["cleared_at"] is None
def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
@ -1462,8 +1512,8 @@ def test_clear_source_bans_noop_when_nothing_to_clear() -> None:
def test_clear_source_bans_resets_escalation() -> None:
"""DELETE, а не banned_until=now(): снятие обнуляет и ban_count — следующий бан
начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию."""
"""Гашение (banned_until=now(), ban_count=0), а не удаление строки: следующий бан той
же пары начинается с базовых SOURCE_BAN_BASE_HOURS, а не продолжает эскалацию (#3404)."""
assert mark_banned is not None
db = FakeSession(
[_proxy(1, affinity="any"), _proxy(2, affinity="any")],

View file

@ -105,9 +105,19 @@ class FakeSession:
)
return _FakeResult([])
if "DELETE FROM scrape_proxy_source_bans" in sql: # clear_source_bans (#2600 п.2)
cleared = [b for b in self.source_bans if b["proxy_id"] == p["proxy_id"]]
self.source_bans = [b for b in self.source_bans if b not in cleared]
if "UPDATE scrape_proxy_source_bans" in sql and "cleared_reason" in sql:
# clear_source_bans (#3404): гасит строку (cleared_reason=reason), а НЕ
# удаляет — строка остаётся для трассируемости. Фейк мутирует найденные
# записи в месте, не убирая их из self.source_bans.
target_source = p.get("source")
cleared = [
b
for b in self.source_bans
if b["proxy_id"] == p["proxy_id"]
and (target_source is None or b["source"] == target_source)
]
for b in cleared:
b["cleared_reason"] = p["reason"]
return _FakeResult([{"source": b["source"]} for b in cleared])
if "UPDATE scrape_proxies" in sql and "exit_ip" in sql:
@ -404,7 +414,13 @@ async def test_successful_rotation_clears_source_bans(monkeypatch: pytest.Monkey
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
assert result.ok is True
assert db.source_bans == [{"proxy_id": 2, "source": "avito"}]
# #3404: строки не удаляются — гасятся (cleared_reason проставлен), остаются для
# трассируемости. Чужой узел (proxy_id=2) не тронут вообще.
proxy1_bans = [b for b in db.source_bans if b["proxy_id"] == 1]
assert len(proxy1_bans) == 2
assert all(b.get("cleared_reason") == "exit ip rotated (status=200)" for b in proxy1_bans)
other_ban = next(b for b in db.source_bans if b["proxy_id"] == 2)
assert other_ban.get("cleared_reason") is None
async def test_failed_rotation_keeps_source_bans(monkeypatch: pytest.MonkeyPatch) -> None:
@ -471,10 +487,13 @@ async def test_token_never_appears_in_reason_success(monkeypatch: pytest.MonkeyP
fake_client, _ = _fake_async_client(response=(200, {"ip": "1.1.1.1"}), exception=None)
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
db = FakeSession(_proxy_row())
db = FakeSession(_proxy_row(), source_bans=[{"proxy_id": 1, "source": "avito"}])
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
assert SECRET_TOKEN not in (result.reason or "")
assert SECRET_TOKEN not in (result.new_ip or "")
# #3404: успешная ротация гасит бан и пишет cleared_reason — секрет не должен
# попасть и туда (та же гигиена, что для reason/note/логов).
assert SECRET_TOKEN not in (db.source_bans[0].get("cleared_reason") or "")
async def test_token_never_appears_in_reason_on_401(monkeypatch: pytest.MonkeyPatch) -> None:
@ -574,7 +593,10 @@ async def test_token_never_appears_in_log_messages_or_sentry_text(
else:
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", _no_http_allowed())
db = FakeSession(_proxy_row(rotate_url=rotate_url))
db = FakeSession(
_proxy_row(rotate_url=rotate_url),
source_bans=[{"proxy_id": 1, "source": "avito"}],
)
with caplog.at_level(logging.DEBUG):
caplog.clear()
await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
@ -583,6 +605,12 @@ async def test_token_never_appears_in_log_messages_or_sentry_text(
assert SECRET_TOKEN not in record.getMessage(), (
f"scenario={name}: token leaked into log message args"
)
if name == "success":
# #3404: успешная ротация гасит бан и пишет cleared_reason — секрет не
# должен попасть и туда.
assert SECRET_TOKEN not in (db.source_bans[0].get("cleared_reason") or ""), (
f"scenario={name}: token leaked into cleared_reason"
)
assert sentry_texts, "expected at least one Sentry capture (401 scenario)"
assert all(SECRET_TOKEN not in text for text in sentry_texts)
@ -646,7 +674,9 @@ async def test_mobileproxy_success_clears_source_bans(monkeypatch: pytest.Monkey
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
assert result.ok is True
assert db.source_bans == []
# #3404: гашение, не удаление — строка остаётся, cleared_reason помечает снятие.
assert len(db.source_bans) == 1
assert db.source_bans[0]["cleared_reason"] == "exit ip rotated (mobileproxy status=200)"
async def test_mobileproxy_non_ok_status_is_failure_and_consumes_quota(
@ -737,7 +767,10 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs(
fake_client, _ = _fake_async_client_get(response=response, exception=exception)
monkeypatch.setattr(proxy_rotation.httpx, "AsyncClient", fake_client)
db = FakeSession(_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL))
db = FakeSession(
_proxy_row(rotate_url=_MOBILEPROXY_ROTATE_URL),
source_bans=[{"proxy_id": 1, "source": "avito"}],
)
with caplog.at_level(logging.DEBUG):
caplog.clear()
result = await proxy_rotation.rotate_proxy(db, 1) # type: ignore[arg-type]
@ -749,6 +782,10 @@ async def test_mobileproxy_key_never_leaks_in_reason_or_db_or_logs(
assert MOBILEPROXY_KEY not in record.getMessage(), (
f"scenario={name}: proxy_key leaked into log message args"
)
if name == "success":
# #3404: успешная ротация гасит бан и пишет cleared_reason — ключ не
# должен попасть и туда.
assert MOBILEPROXY_KEY not in (db.source_bans[0].get("cleared_reason") or ""), name
async def test_mobileproxy_unknown_query_shape_still_masks_url_in_refusal(

View file

@ -302,7 +302,14 @@ async def test_probe_clears_only_its_own_ban(monkeypatch: pytest.MonkeyPatch) ->
avito = db._ban(1, "avito")
assert avito is not None, "чужой бан проба снимать не имеет права"
assert (avito["reason"], avito["ban_count"]) == ("banned:avito", 1), "и не переписывать"
assert db._ban(1, "cian") is None, "свой вердикт проба обязана снять"
cian = db._ban(1, "cian")
# #3404: снятие теперь ГАСИТ строку, а не удаляет — вердикт пробы перестаёт
# блокировать выдачу (banned_until в прошлом, ban_count обнулён), но остаётся
# виден в истории как снятый досрочно.
assert cian is not None, "погашенная строка живёт до purge (#3404)"
assert cian["ban_count"] == 0, "свой вердикт проба обязана снять"
assert cian["banned_until"] <= datetime.now(UTC), "и снять немедленно"
assert cian["cleared_at"] is not None, "снятие должно быть отмечено"
assert counters["pair_cleared"] == 1

View file

@ -0,0 +1,403 @@
"""#3404: атрибуция прогона к узлу (scrape_runs.proxy_id) + гашение бана вместо DELETE.
Offline-тесты (без live БД) в стиле test_proxy_pool.py / test_3390_single_runs_module.py:
stateful fake-сессия эмулирует таблицы scrape_proxies / scrape_proxy_source_bans /
scrape_runs, интерпретируя SQL по ключевым фрагментам, так что acquire/mark_banned/
clear_source_bans/attribute_run_proxy проверяются по фактическому изменению состояния,
а не по замоканному возврату.
Покрытие:
1. clear_source_bans гасит строку (banned_until<=now, ban_count=0, cleared_at
заполнен), а НЕ удаляет строка остаётся в таблице.
2. Эскалация 1:1: бан clear бан той же пары снова стартует с ровно
SOURCE_BAN_BASE_HOURS, а не удвоенного срока (главный тест issue).
3. Погашенная строка не мешает acquire(source) узел выдаётся как обычно.
4. attribute_run_proxy доводит proxy_id прогона до scrape_runs через acquire();
путь без run_id (health-checker, run_id=None) оставляет scrape_runs нетронутым.
5. Смена узла за прогон (повторный acquire тем же run_id на другой узел)
отражается в counters['proxy_ids'] без дублей.
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from datetime import UTC, datetime, timedelta
from typing import Any
from app.services.proxy_pool import (
SOURCE_BAN_BASE_HOURS,
acquire,
attribute_run_proxy,
clear_source_bans,
mark_banned,
)
# ── stateful fake session (scrape_proxies + scrape_proxy_source_bans + scrape_runs) ──
class _FakeResult:
def __init__(self, rows: list[Any]) -> None:
self._rows = rows
def mappings(self) -> _FakeResult:
return self
def fetchone(self) -> Any:
return self._rows[0] if self._rows else None
def first(self) -> Any:
return self._rows[0] if self._rows else None
def fetchall(self) -> list[Any]:
return [type("Row", (), r)() if isinstance(r, dict) else r for r in self._rows]
def all(self) -> list[Any]:
return list(self._rows)
def _proxy(
pid: int,
*,
affinity: str = "avito",
enabled: bool = True,
leased_by: int | None = None,
) -> dict[str, Any]:
return {
"id": pid,
"url": f"http://proxy{pid}",
"kind": "residential",
"rotate_url": None,
"browser_unfit_since": None,
"enabled": enabled,
"consecutive_fails": 0,
"provider_affinity": affinity,
"leased_by": leased_by,
"leased_at": None,
"expires_at": None,
"last_ok_at": None,
}
class RunAttributionDb:
"""Мини-Postgres: scrape_proxies + scrape_proxy_source_bans + scrape_runs.
Гейты по SQL-подстрокам скопированы из проверяемого кода (proxy_pool.py) как в
test_proxy_pool.py/test_3390: значение проверяется по факту исполнения реального
statement'а, а не зашито ожиданием теста.
"""
def __init__(self, proxies: list[dict[str, Any]]) -> None:
self.proxies = proxies
self.bans: list[dict[str, Any]] = []
self.runs: dict[int, dict[str, Any]] = {}
self.set_local_statements: list[str] = []
def add_run(self, run_id: int) -> None:
self.runs[run_id] = {"proxy_id": None, "counters": {}}
def _by_id(self, pid: int) -> dict[str, Any] | None:
return next((r for r in self.proxies if r["id"] == pid), None)
def _ban(self, pid: int, source: str) -> dict[str, Any] | None:
return next((b for b in self.bans if b["proxy_id"] == pid and b["source"] == source), None)
def _has_active_ban(self, pid: int, source: str) -> bool:
b = self._ban(pid, source)
return b is not None and b["banned_until"] > datetime.now(UTC)
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = " ".join(str(stmt).split())
p = params or {}
if "pg_advisory_xact_lock" in sql:
return _FakeResult([])
if "expires_at <= now()" in sql and "FOR UPDATE SKIP LOCKED" not in sql:
return _FakeResult([]) # без просроченных узлов в этих тестах
if "FOR UPDATE SKIP LOCKED" in sql: # acquire SELECT (primary/fallback)
provider = p["provider"]
max_fails = p["max_fails"]
primary = "provider_affinity IN" in sql
cands = [
r
for r in self.proxies
if r["enabled"]
and r["consecutive_fails"] < max_fails
and r["leased_by"] is None
and not self._has_active_ban(r["id"], provider)
and (r["provider_affinity"] in (provider, "any") if primary else True)
]
cands.sort(
key=lambda r: (
r["browser_unfit_since"] is not None,
r["last_ok_at"] is None,
r["last_ok_at"] or datetime.min.replace(tzinfo=UTC),
r["id"],
)
)
return _FakeResult(cands[:1])
if "SET leased_by = CAST(:run_id" in sql: # acquire lease UPDATE
row = self._by_id(p["id"])
if row is not None:
row["leased_by"] = p["run_id"]
row["leased_at"] = datetime.now(UTC)
return _FakeResult([])
if sql.startswith("SET LOCAL"):
# attribute_run_proxy ставит lock_timeout перед UPDATE scrape_runs:
# ждать на блокировке строки прогона в пути выдачи прокси нельзя.
# Для фейка это no-op, но проглатывать молча нечестно — гейт ниже
# ловит любой ДРУГОЙ незнакомый SQL.
self.set_local_statements.append(sql)
return _FakeResult([])
if "SET leased_by = NULL" in sql: # release
row = self._by_id(p["id"])
if row is not None:
row["leased_by"] = None
row["leased_at"] = None
return _FakeResult([])
if "INSERT INTO scrape_proxy_source_bans" in sql: # mark_banned upsert
proxy_id, source = p["proxy_id"], p["source"]
if self._by_id(proxy_id) is None:
return _FakeResult([])
max_fails = p["max_fails"]
still_available = any(
r["id"] != proxy_id
and r["enabled"]
and r["consecutive_fails"] < max_fails
and not self._has_active_ban(r["id"], source)
for r in self.proxies
)
if not still_available:
return _FakeResult([]) # protected — последний узел
now = datetime.now(UTC)
ban = self._ban(proxy_id, source)
reason = p["reason"]
if ban is None:
ban = {
"proxy_id": proxy_id,
"source": source,
"ban_count": 1,
"banned_until": now + timedelta(hours=p["base_hours"]),
"reason": reason,
"cleared_at": None,
"cleared_reason": None,
}
self.bans.append(ban)
else:
is_active_foreign = (
ban["banned_until"] > now
and ban.get("reason") != reason
and reason != p.get("live_reason")
)
if is_active_foreign:
return _FakeResult([]) # deferred — чужой активный владелец
ban["ban_count"] += 1
hours = min(p["base_hours"] * 2 ** (ban["ban_count"] - 1), p["max_hours"])
ban["banned_until"] = now + timedelta(hours=hours)
ban["reason"] = reason
ban["cleared_at"] = None
ban["cleared_reason"] = None
return _FakeResult(
[{"ban_count": ban["ban_count"], "banned_until": ban["banned_until"]}]
)
if "UPDATE scrape_proxy_source_bans" in sql and "cleared_at" in sql: # clear_source_bans
proxy_id = p["proxy_id"]
source = p.get("source")
only_reason = p.get("only_reason")
reason = p["reason"]
now = datetime.now(UTC)
cleared = []
for b in self.bans:
if b["proxy_id"] != proxy_id:
continue
if source is not None and b["source"] != source:
continue
if only_reason is not None and b.get("reason") != only_reason:
continue
if (
b.get("cleared_at") is not None
and b["ban_count"] == 0
and b["banned_until"] <= now
):
continue # гейт покоя — уже погашена, no-op
b["banned_until"] = now
b["ban_count"] = 0
b["cleared_at"] = now
b["cleared_reason"] = reason
cleared.append({"source": b["source"]})
return _FakeResult(cleared)
if "UPDATE scrape_runs" in sql: # attribute_run_proxy (#3404)
run = self.runs.get(p["run_id"])
if run is None:
return _FakeResult([])
run["proxy_id"] = p["proxy_id"]
proxy_ids: list[int] = run["counters"].get("proxy_ids", [])
if p["proxy_id"] not in proxy_ids:
proxy_ids = [*proxy_ids, p["proxy_id"]]
run["counters"]["proxy_ids"] = proxy_ids
return _FakeResult([{"id": p["run_id"]}])
raise AssertionError(f"unhandled SQL in fake session: {sql[:120]}")
def commit(self) -> None:
pass
def rollback(self) -> None:
pass
# ── 1. clear_source_bans гасит, а не удаляет ────────────────────────────────
def test_clear_source_bans_gates_not_deletes() -> None:
db = RunAttributionDb([_proxy(1), _proxy(2)])
assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type]
assert len(db.bans) == 1
cleared = clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
assert cleared == 1
assert len(db.bans) == 1, "строка должна остаться на месте, не удалиться"
ban = db.bans[0]
assert ban["banned_until"] <= datetime.now(UTC)
assert ban["ban_count"] == 0
assert ban["cleared_at"] is not None
assert ban["cleared_reason"] == "manual enable"
# ── 2. Эскалация сохранилась 1:1 (главный тест issue) ───────────────────────
def test_escalation_resets_to_base_after_clear() -> None:
db = RunAttributionDb([_proxy(1), _proxy(2)])
assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type]
first_ban = db.bans[0]
assert first_ban["ban_count"] == 1
first_span = first_ban["banned_until"] - datetime.now(UTC)
assert timedelta(hours=SOURCE_BAN_BASE_HOURS - 1) < first_span <= timedelta(
hours=SOURCE_BAN_BASE_HOURS
)
clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
assert db.bans[0]["ban_count"] == 0
assert mark_banned(db, 1, source="avito", reason="banned:avito") == "banned" # type: ignore[arg-type]
second_ban = db.bans[0]
assert second_ban["ban_count"] == 1, "после clear эскалация обязана стартовать заново"
second_span = second_ban["banned_until"] - datetime.now(UTC)
assert timedelta(hours=SOURCE_BAN_BASE_HOURS - 1) < second_span <= timedelta(
hours=SOURCE_BAN_BASE_HOURS
), f"срок {second_span} обязан быть БАЗОВЫМ ({SOURCE_BAN_BASE_HOURS}ч), а не удвоенным"
# ── 3. Погашенная строка не мешает выдаче узла ──────────────────────────────
def test_acquire_ignores_cleared_ban() -> None:
db = RunAttributionDb([_proxy(1), _proxy(2)])
mark_banned(db, 1, source="avito", reason="banned:avito") # type: ignore[arg-type]
clear_source_bans(db, 1, reason="manual enable") # type: ignore[arg-type]
# узел 2 занят чужим прогоном — иначе acquire мог бы честно выдать его вместо
# узла 1, и тест перестал бы проверять именно "погашенный бан не блокирует".
db.proxies[1]["leased_by"] = 999
lease = acquire(db, "avito") # type: ignore[arg-type]
assert lease is not None, "погашенный бан не должен блокировать acquire"
assert lease.id == 1
def test_acquire_still_blocks_active_ban_after_clear_gate_check() -> None:
"""Контроль ложноположительного теста выше: АКТИВНЫЙ (не погашенный) бан acquire
по-прежнему блокирует это не сломано этим же изменением."""
db = RunAttributionDb([_proxy(1), _proxy(2)])
mark_banned(db, 1, source="avito", reason="banned:avito") # type: ignore[arg-type]
lease = acquire(db, "avito") # type: ignore[arg-type]
assert lease is not None
assert lease.id == 2, "узел 1 активно забанен по avito — выдан должен быть узел 2"
# ── 4. proxy_id прогона доезжает до scrape_runs ─────────────────────────────
def test_acquire_attributes_run_proxy() -> None:
db = RunAttributionDb([_proxy(1)])
db.add_run(42)
lease = acquire(db, "avito", run_id=42) # type: ignore[arg-type]
assert lease is not None
assert db.runs[42]["proxy_id"] == lease.id == 1
assert db.runs[42]["counters"]["proxy_ids"] == [1]
# Атрибуция обязана ограничить ожидание блокировки: строку прогона параллельно
# пишет heartbeat из другой сессии, а путь выдачи прокси ждать не может.
assert any("lock_timeout" in stmt for stmt in db.set_local_statements)
def test_acquire_without_run_id_leaves_scrape_runs_untouched() -> None:
"""health-checker и прочие не-run вызовы (run_id=None) не трогают scrape_runs."""
db = RunAttributionDb([_proxy(1)])
db.add_run(42)
lease = acquire(db, "avito") # type: ignore[arg-type]
assert lease is not None
assert db.runs[42]["proxy_id"] is None
assert db.runs[42]["counters"] == {}
# ── 5. Смена узла за прогон отражается в counters ───────────────────────────
def test_mid_run_proxy_rotation_recorded_in_counters() -> None:
"""Ре-acquire тем же run_id на другой узел (старый lease ещё держится, как при
ротации в browser_fetcher `_LEASE_ROTATE_AFTER_FAILS`) оба узла в proxy_ids,
proxy_id несёт ПОСЛЕДНИЙ (сверено с attribute_run_proxy docstring, а не с ТЗ)."""
db = RunAttributionDb([_proxy(1), _proxy(2)])
db.add_run(7)
first = acquire(db, "avito", run_id=7) # type: ignore[arg-type]
assert first is not None
# старый lease НЕ освобождён (leased_by=7) — второй acquire тем же run_id
# неизбежно возьмёт другой свободный узел, детерминированно.
assert db._by_id(first.id)["leased_by"] == 7
second = acquire(db, "avito", run_id=7) # type: ignore[arg-type]
assert second is not None
assert second.id != first.id, "тестовая обвязка обязана взять ДРУГОЙ узел"
assert db.runs[7]["proxy_id"] == second.id
assert set(db.runs[7]["counters"]["proxy_ids"]) == {first.id, second.id}
def test_attribute_run_proxy_idempotent_same_proxy() -> None:
"""Повторная атрибуция ТЕМ ЖЕ узлом не дублирует id в proxy_ids."""
db = RunAttributionDb([_proxy(1)])
db.add_run(9)
attribute_run_proxy(db, 9, 1) # type: ignore[arg-type]
attribute_run_proxy(db, 9, 1) # type: ignore[arg-type]
assert db.runs[9]["counters"]["proxy_ids"] == [1]
def test_attribute_run_proxy_best_effort_on_missing_run() -> None:
"""run_id не найден (гонка с финализацией) — best-effort no-op, исключение не летит."""
db = RunAttributionDb([_proxy(1)])
# run 999 никогда не создавался в этой fake-БД
attribute_run_proxy(db, 999, 1) # type: ignore[arg-type] # не должно бросить

View file

@ -52,7 +52,8 @@ def _scalar_result(value: object) -> MagicMock:
def _cleared_bans_result(rows: list[dict[str, Any]] | None = None) -> MagicMock:
"""Ответ на DELETE ... RETURNING source (proxy_pool.clear_source_bans, #2600 п.2).
"""Ответ на UPDATE ... RETURNING source (proxy_pool.clear_source_bans, #3404: гасит
строку banned_until=now()/ban_count=0/cleared_at/cleared_reason, а не DELETE).
PATCH enabled=true снимает баны узла по источникам «ручное включение = чистый
лист», как и обнуление disabled_reason рядом.
@ -337,8 +338,11 @@ def test_patch_enable_clears_source_bans(client: TestClient, db: MagicMock) -> N
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": True})
assert r.status_code == 200, r.text
delete_sql = str(db.execute.call_args_list[1].args[0])
assert "DELETE FROM scrape_proxy_source_bans" in delete_sql
clear_sql = str(db.execute.call_args_list[1].args[0])
# #3404: гашение (UPDATE ... SET cleared_reason=...), а не DELETE — строка остаётся
# для трассируемости, но перестаёт блокировать выдачу.
assert "UPDATE scrape_proxy_source_bans" in clear_sql
assert "cleared_reason" in clear_sql
assert db.execute.call_args_list[1].args[1]["proxy_id"] == 1
@ -351,8 +355,10 @@ def test_patch_disable_keeps_source_bans(client: TestClient, db: MagicMock) -> N
r = client.patch("/api/v1/admin/proxies/1", json={"enabled": False})
assert r.status_code == 200, r.text
# #3404: clear_source_bans теперь UPDATE, а не DELETE — ищем по cleared_reason,
# уникальному для этого запроса маркеру.
assert not any(
"DELETE FROM scrape_proxy_source_bans" in str(c.args[0]) for c in db.execute.call_args_list
"cleared_reason" in str(c.args[0]) for c in db.execute.call_args_list
)

View file

@ -0,0 +1,35 @@
"""ContextVar с id текущего прогона (#3404) — канал для мест, куда `run_id` не доходит
аргументом.
ЗАЧЕМ ЭТО СУЩЕСТВУЕТ
--------------------
`RealProxyProvider` (`app.services.scraper_adapters`) ОДИН объект на весь жизненный
цикл `scheduler_main.py`, реализующий kit-протокол `ProxyProvider.acquire(provider)`
(`scraper_kit/contracts.py`) сигнатура протокола не несёт `run_id`, а конструктору
адаптера неоткуда его взять: он создаётся один раз при старте планировщика, задолго до
того, как какой-либо конкретный прогон начнётся. Пробрасывать `run_id` через цепочку
вызовов от `create_run` до `proxy_pool.acquire` означало бы менять сигнатуры всех
провайдеров пайплайна (curl-путь `providers/_proxy.py`, `browser_fetcher`) ради одного
параметра, который в 99% вызовов (health-check, эстиматор) не нужен вовсе.
`ContextVar`, а не глобальная переменная: `asyncio`-таски и синхронные Celery-контексты
не делят один поток предсказуемо, а `ContextVar` копируется в `asyncio.to_thread`/
дочерние таски и не путается между КОНКУРЕНТНЫМИ прогонами в одном процессе тот же
довод, что и у `app.core.public_request.is_public_request`.
Что это НЕ гарантирует: если прогон и вызов `acquire()` идут в разных тредах/тасках,
созданных ДО `current_run_id.set()`, значение не унаследуется атрибуция для такого
прогона молча останется NULL (не ошибка, см. `proxy_pool.attribute_run_proxy`).
"""
from __future__ import annotations
from contextvars import ContextVar
#: id прогона (`scrape_runs.id`), который сейчас выполняется в этом контексте.
#: None — контекст не привязан ни к какому прогону (health-check, эстиматор, admin).
#: `runs.py` управляет им напрямую через `.set()`/`.set(None)` (не через
#: context-manager: жизненный цикл прогона не вложен в один `with`-блок — он
#: пересекает несколько функций/финализаторов, см. `create_run`/`mark_done`/
#: `mark_failed`).
current_run_id: ContextVar[int | None] = ContextVar("scraper_kit_current_run_id", default=None)

View file

@ -45,6 +45,8 @@ from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
from scraper_kit.orchestration.run_context import current_run_id
try: # sentry опционален — kit не тянет его в зависимостях
import sentry_sdk
except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует
@ -621,6 +623,12 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / );
отдельной колонки run_type больше нет она 3244 прогона подряд молчала
дефолтом 'city_sweep' и подписывала им, например, proxy_healthcheck (#2674).
После коммита выставляет `current_run_id` (#3404) — канал, которым `run_id`
доходит до `RealProxyProvider.acquire`, живущего одним объектом на весь
планировщик (см. `run_context.py`). Финализаторы (`mark_done`/`mark_failed`/
`mark_banned`/`mark_cancelled`) сбрасывают его обратно в None, чтобы следующий
прогон в том же треде не унаследовал чужой id.
Returns run_id (bigint).
"""
row = db.execute(
@ -637,7 +645,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
).fetchone()
db.commit()
assert row is not None, "scrape_runs INSERT returned no id"
return int(row.id)
run_id = int(row.id)
current_run_id.set(run_id)
return run_id
def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = None) -> int:
@ -881,6 +891,9 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
if row is None:
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# #3404: прогон завершён — сбрасываем канал атрибуции прокси, чтобы следующий
# прогон в том же треде (или health-check между ними) не унаследовал этот run_id.
current_run_id.set(None)
# #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая
# для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort,
# после коммита — статус уже персистирован в БД.
@ -923,6 +936,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
if row is None:
logger.warning("mark_failed no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# #3404: см. mark_done — сброс канала атрибуции прокси.
current_run_id.set(None)
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
_alert_on_run_id(db, run_id)
@ -988,6 +1003,8 @@ def mark_banned(
if row is None:
logger.warning("mark_banned no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# #3404: см. mark_done — сброс канала атрибуции прокси.
current_run_id.set(None)
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
_alert_on_run_id(db, run_id)
@ -1173,6 +1190,9 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
{"run_id": run_id},
).fetchone()
db.commit()
if result is not None:
# #3404: см. mark_done — сброс канала атрибуции прокси.
current_run_id.set(None)
return result is not None