fix(tradein/proxy_egress): прогон на egress-пути тоже знает свой узел (#3404)

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 0 из 56, yandex_address_backfill 0 из 1.

resolve_proxy_url при выбранном узле и выставленном current_run_id зовёт
attribute_run_proxy. Своей короткой сессией: db вызывающего — долгоживущая
сессия прогона посреди работы, а атрибуция коммитит и на сбое откатывает.

Тесты на живом Postgres: атрибуция внутри прогона (proxy_id и
counters.proxy_ids), вне прогона scrape_runs не тронут, незакоммиченная
работа вызывающего не коммитится. Красные прогоны: на старом коде
«assert None == 70»; с атрибуцией через db вызывающего — «атрибуция
закоммитила чужую транзакцию».

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
bot-backend 2026-09-17 13:03:22 +05:00
parent 611878df4c
commit ab003877dc
3 changed files with 175 additions and 10 deletions

View file

@ -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)

View file

@ -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` ПОСЛЕДНИЙ использованный узел (перезаписывается при
каждой новой выдаче); полная цепочка узлов, если она менялась, в

View file

@ -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