Выбор оператора мобильного прокси опирался на две ненадёжные опоры. Первая: `scrape_runs` не знала, через какой узел шёл прогон — колонка `proxy_id` была только у банов и ротаций. «Какой узел собрал 5 карточек из 21» не выяснялось ни одним запросом. Вторая: `clear_source_bans` делала DELETE, а зовётся она после КАЖДОЙ успешной ротации exit-IP. У #540723 (МегаФон) 23 успешные ротации и ноль строк банов, у #540722 (Tele2) ротаций почти не было и 7 банов. «7 против 0» читалось как «Tele2 хуже», хотя в той же мере это «у МегаФона историю стёрли 23 раза». Теперь: - `scrape_runs.proxy_id` — последний выданный прогону узел; полная цепочка (если узел менялся mid-run) копится в `counters.proxy_ids`. Пишет `proxy_pool.attribute_run_proxy` из единственной точки — сразу после выдачи лиза в `acquire()`, поэтому curl-путь, браузерный sticky lease и ре-acquire при ротации покрыты одинаково. `run_id` доходит до адаптера через ContextVar (`scraper_kit.orchestration.run_context`): протокол `ProxyProvider.acquire` его не несёт, а `RealProxyProvider` живёт одним объектом на весь планировщик. Best-effort: `lock_timeout` 2с и проглоченное исключение — диагностика не вправе ронять выдачу прокси или ждать на блокировке строки прогона. - `clear_source_bans` гасит строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`) вместо удаления. Эскалация сохраняется 1:1: формула в `mark_banned` берёт ПРЕДЫДУЩИЙ `ban_count` показателем степени, при нуле это ровно `SOURCE_BAN_BASE_HOURS` — как после DELETE. Строка доживает до штатного purge по `SOURCE_BAN_PURGE_DAYS`. Для всех читателей `scrape_proxy_source_bans` погашенная строка неотличима от отсутствующей: acquire, оба guard-подзапроса `mark_banned`, `proxy_egress` (ранжирование по `ban_count` даёт 0, как у узла без истории), admin `_active_ban` — все гейтятся по `banned_until > now()`. Ничего не бэкфиллится: связать прошедшие прогоны с узлами нечем (`leased_by` исторически = NON_RUN_LEASE_MARKER), врать восстановленным значением нельзя. Миграция 287. Тесты: 9 новых на обе части (главный — эскалация после гашения даёт базовые 6ч, а не удвоенные) + 14 существующих переведены с DELETE-семантики на гашение, включая проверку, что секрет ротации не утекает в новое `cleared_reason`. Полный прогон бэкенда: 5600 passed, 37 skipped. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011WHFxVPWoBnSZihkdH1Uou
403 lines
18 KiB
Python
403 lines
18 KiB
Python
"""#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] # не должно бросить
|