diff --git a/tradein-mvp/backend/app/services/proxy_egress.py b/tradein-mvp/backend/app/services/proxy_egress.py index 12880ea7..a7503ad0 100644 --- a/tradein-mvp/backend/app/services/proxy_egress.py +++ b/tradein-mvp/backend/app/services/proxy_egress.py @@ -16,7 +16,8 @@ reap_stale_leases) — она рассчитана на долгоживущие гарантированного `release` на каждом пути выхода; занимать под них lease значило бы дырявить пул фантомно занятыми узлами при малейшей утечке release. Резолвер ниже — ЧИСТО READ, той же таблицы `scrape_proxies` + `scrape_proxy_source_bans`, без блокировок -и без мутаций. +и без мутаций пула. Единственная запись — атрибуция прогона (`scrape_runs.proxy_id`, #3404), +своей короткой сессией, см. `resolve_proxy_url`. ПРАВИЛО ВЫБОРА: enabled=true, consecutive_fails < proxy_pool.MAX_CONSECUTIVE_FAILS (тот же карантинный порог, что у acquire), нет АКТИВНОЙ строки (banned_until > now()) @@ -69,7 +70,7 @@ from sqlalchemy.orm import Session from app.core.config import settings as _settings from app.core.db import SessionLocal as _SessionLocal -from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS +from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS, attribute_run_proxy logger = logging.getLogger(__name__) @@ -277,6 +278,21 @@ def resolve_proxy_url(db: Session, source: str) -> str | None: _safe_label(candidate.id, candidate.label, candidate.url), candidate.ban_count, ) + # #3404: прогоны на этом пути (yandex_detail_backfill, yandex_address_backfill, + # curl-ветка avito_detail_backfill) lease не берут, а атрибуцию раньше писал только + # acquire() — у yandex_detail_backfill proxy_id был NULL в 56 прогонах из 56. + # Своя сессия, не db вызывающего: attribute_run_proxy коммитит, а на сбое + # откатывает, db же — долгоживущая сессия прогона посреди работы (см. + # _pick_candidate). Сбой атрибуции проглатывается внутри — выдачу не роняет. + from scraper_kit.orchestration.run_context import current_run_id + + run_id = current_run_id.get() + if run_id is not None: + attr_db = _SessionLocal() + try: + attribute_run_proxy(attr_db, run_id, candidate.id) + finally: + attr_db.close() return candidate.url diag = _diagnose_no_candidate(db, source) diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index 667ef0f3..a57519d8 100644 --- a/tradein-mvp/backend/app/services/proxy_pool.py +++ b/tradein-mvp/backend/app/services/proxy_pool.py @@ -459,9 +459,9 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe ) 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. + # #3404: покрывает все пути АРЕНДЫ (curl — acquire на каждый вызов, браузер — + # sticky lease на весь прогон, ре-acquire при ротации узла mid-run). Путь без + # аренды (proxy_egress.resolve_proxy_url) пишет атрибуцию сам — см. docstring. attribute_run_proxy(db, run_id, proxy_id) if fallback_used: logger.warning( @@ -499,11 +499,12 @@ 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`) — то есть смена узла ЗА прогон фиксируется сама, - без отдельного вызова с чьей-либо стороны. + Два писателя. `acquire()` сразу после выдачи lease'а: покрывает и curl-путь + (acquire на каждый вызов), и браузерный sticky lease (один acquire на весь прогон), + и ре-acquire при ротации узла mid-run (`browser_fetcher`, + `_LEASE_ROTATE_AFTER_FAILS`) — то есть смена узла ЗА прогон фиксируется сама. + И `proxy_egress.resolve_proxy_url` — egress без аренды (yandex_detail_backfill и + др.); до этого у таких прогонов proxy_id оставался NULL (#3404, прод 17.09: 0 из 56). `scrape_runs.proxy_id` — ПОСЛЕДНИЙ использованный узел (перезаписывается при каждой новой выдаче); полная цепочка узлов, если она менялась, — в diff --git a/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py b/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py new file mode 100644 index 00000000..552bdfd6 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py @@ -0,0 +1,148 @@ +"""Прогон на egress-пути знает свой узел (#3404, хвост PR #3405). + +PR #3405 писал `scrape_runs.proxy_id` только из `proxy_pool.acquire()`. Прогоны, которые +берут прокси через `proxy_egress.resolve_proxy_url` (без аренды: yandex_detail_backfill, +yandex_address_backfill, curl-ветка avito_detail_backfill), атрибуцию не получали — +прод 17.09: у yandex_detail_backfill proxy_id NULL в 56 прогонах из 56. + +Живой Postgres, настоящие сессии (атрибуция идёт своей сессией, поэтому строки +коммитятся и удаляются в finally). Без БД — skip. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +from collections.abc import Iterator +from typing import Any + +import pytest +from scraper_kit.orchestration.run_context import current_run_id +from sqlalchemy import create_engine, text +from sqlalchemy.orm import Session, sessionmaker + +from app.services import proxy_egress + + +def _live_engine() -> Any | None: + dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "") + if not dsn or "localhost:5432/test" in dsn: + return None + try: + engine = create_engine(dsn, future=True) + with engine.connect() as conn: + conn.execute(text("SELECT proxy_id FROM scrape_runs LIMIT 1")) + return engine + except Exception: + return None + + +_ENGINE = _live_engine() +pytestmark = pytest.mark.skipif(_ENGINE is None, reason="no reachable Postgres test DB") + + +@pytest.fixture +def env(monkeypatch: pytest.MonkeyPatch) -> Iterator[dict[str, Any]]: + """Узел пула + два прогона: `run` — текущий, `other` — строка, которую вызывающий + держит изменённой без коммита.""" + assert _ENGINE is not None + factory = sessionmaker(bind=_ENGINE, future=True) + monkeypatch.setattr(proxy_egress, "_SessionLocal", factory) + setup = factory() + proxy_id = setup.execute( + text( + "INSERT INTO scrape_proxies (url, provider_affinity) " + "VALUES ('http://t3404-' || gen_random_uuid(), 'any') RETURNING id" + ) + ).scalar_one() + run_ids = [ + setup.execute( + text( + "INSERT INTO scrape_runs (source, status) " + "VALUES ('test_3404', 'running') RETURNING id" + ) + ).scalar_one() + for _ in range(2) + ] + setup.commit() + token = current_run_id.set(None) + try: + yield {"proxy_id": int(proxy_id), "run": int(run_ids[0]), "other": int(run_ids[1])} + finally: + current_run_id.reset(token) + setup.rollback() + setup.execute(text("DELETE FROM scrape_runs WHERE id = ANY(:ids)"), {"ids": run_ids}) + setup.execute(text("DELETE FROM scrape_proxies WHERE id = :id"), {"id": proxy_id}) + setup.commit() + setup.close() + + +def _run_row(run_id: int) -> Any: + assert _ENGINE is not None + with _ENGINE.connect() as conn: + return conn.execute( + text("SELECT proxy_id, counters FROM scrape_runs WHERE id = :id"), {"id": run_id} + ).one() + + +def _node_id_of(url: str) -> int: + assert _ENGINE is not None + with _ENGINE.connect() as conn: + return int( + conn.execute( + text("SELECT id FROM scrape_proxies WHERE url = :u"), {"u": url} + ).scalar_one() + ) + + +def test_resolve_within_run_writes_node_to_scrape_runs(env: dict[str, Any]) -> None: + """Главный случай: прогон идёт, egress выбран — прогон знает узел. + + На старом коде proxy_id остаётся NULL.""" + current_run_id.set(env["run"]) + caller = Session(bind=_ENGINE, future=True) + try: + url = proxy_egress.resolve_proxy_url(caller, "yandex") + finally: + caller.close() + + assert url is not None + node = _node_id_of(url) + row = _run_row(env["run"]) + assert row.proxy_id == node + counters = row.counters if isinstance(row.counters, dict) else json.loads(row.counters) + assert counters.get("proxy_ids") == [node] + + +def test_resolve_outside_run_writes_nothing(env: dict[str, Any]) -> None: + """Без прогона (админка, проверка кук) — scrape_runs не трогаем.""" + caller = Session(bind=_ENGINE, future=True) + try: + assert proxy_egress.resolve_proxy_url(caller, "yandex") is not None + finally: + caller.close() + + assert _run_row(env["run"]).proxy_id is None + + +def test_attribution_does_not_commit_callers_transaction(env: dict[str, Any]) -> None: + """db вызывающего — долгоживущая сессия прогона посреди работы. Атрибуция не имеет + права закоммитить его незавершённые изменения (или откатить их на сбое).""" + current_run_id.set(env["run"]) + caller = Session(bind=_ENGINE, future=True) + try: + caller.execute( + text("UPDATE scrape_runs SET counters = '{\"t3404\": 1}'::jsonb WHERE id = :id"), + {"id": env["other"]}, + ) + proxy_egress.resolve_proxy_url(caller, "yandex") + caller.rollback() + finally: + caller.close() + + other = _run_row(env["other"]).counters or {} + assert "t3404" not in other, "атрибуция закоммитила чужую транзакцию" + assert _run_row(env["run"]).proxy_id is not None