diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index 437e937c..667ef0f3 100644 --- a/tradein-mvp/backend/app/services/proxy_pool.py +++ b/tradein-mvp/backend/app/services/proxy_pool.py @@ -275,12 +275,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe чужая — только запасной вариант, чтобы источник не голодал при живых свободных узлах чужой affinity (#2600). - Fallback НЕ трогает последний enabled-узел выделенной (не-'any') affinity: если + Fallback НЕ трогает последний пригодный узел выделенной (не-'any') affinity: если fallback заберёт его под чужой источник, «свой» останется без прокси вообще — хуже, чем голодание исходного источника, которое фикс призван устранить. Кандидат - участвует в fallback, только если его affinity='any' ИЛИ у этой affinity есть ДРУГОЙ - enabled-узел (EXISTS-подзапрос) — т.е. выдача не обнулит доступность выделенной - affinity целиком. + участвует в fallback, только если его affinity='any' ИЛИ у источника этой affinity + без него останется ДРУГОЙ кандидат (EXISTS-подзапрос): узел той же affinity или + 'any', здоровый, с живой арендой порта и не забаненный этим источником (#3299 — + до него 'any'-узлы не считались, и единственный выделенный узел не выдавался никому). Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql — единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR @@ -401,17 +402,29 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe -- fallback увести последний реально рабочий узел выделенной -- affinity и обрушить её (два domclick-узла, один забанен -- domclick'ом → второй уходит под avito → domclick без прокси). + -- + -- #3299: вопрос «останется ли у источника sp.provider_affinity + -- хоть один кандидат без sp», а не «есть ли ВТОРОЙ узел той же + -- привязки». 'any'-узлы основной запрос выдаёт выделенному + -- источнику наравне, значит и backup'ом они считаются; прежний + -- `= sp.provider_affinity` прятал единственный выделенный узел от + -- всех. Здоровье и срок аренды — как в основном запросе. + -- leased_by НЕ проверяем: аренда вернётся через минуты, а + -- отказ из-за чужой аренды прятал бы узел при любом параллельном + -- прогоне (тот же довод, что у защиты в mark_banned). OR EXISTS ( SELECT 1 FROM scrape_proxies AS other - WHERE other.provider_affinity = sp.provider_affinity + WHERE other.provider_affinity IN (sp.provider_affinity, 'any') AND other.enabled + AND other.consecutive_fails < CAST(:max_fails AS integer) + AND (other.expires_at IS NULL OR other.expires_at > now()) AND other.id <> sp.id AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b2 WHERE b2.proxy_id = other.id - AND b2.source = other.provider_affinity + AND b2.source = sp.provider_affinity AND b2.banned_until > now() ) ) @@ -535,8 +548,7 @@ def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None: db.commit() if row is None: logger.debug( - "proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already " - "finalized?)", + "proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already finalized?)", run_id, ) except Exception: @@ -999,17 +1011,22 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None = -- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом -- не считается (иначе защита сочла бы affinity живой, когда она -- уже нет). + -- #3299: предикат backup'а — буква в букву как в fallback + -- acquire() (там же обоснование): 'any'-узлы в счёт, здоровье + -- и срок аренды проверяются. OR EXISTS ( SELECT 1 FROM scrape_proxies other - WHERE other.provider_affinity = sp.provider_affinity + WHERE other.provider_affinity IN (sp.provider_affinity, 'any') AND other.enabled + AND other.consecutive_fails < CAST(:max_fails AS integer) + AND (other.expires_at IS NULL OR other.expires_at > now()) AND other.id <> sp.id AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b2 WHERE b2.proxy_id = other.id - AND b2.source = other.provider_affinity + AND b2.source = sp.provider_affinity AND b2.banned_until > now() ) ) diff --git a/tradein-mvp/backend/tests/services/test_proxy_pool.py b/tradein-mvp/backend/tests/services/test_proxy_pool.py index c5466431..6b34e725 100644 --- a/tradein-mvp/backend/tests/services/test_proxy_pool.py +++ b/tradein-mvp/backend/tests/services/test_proxy_pool.py @@ -167,7 +167,9 @@ class FakeSession: exp = row.get("expires_at") return exp is None or exp > datetime.now(UTC) - if "provider_affinity IN" in sql: # primary: своя affinity ИЛИ 'any' + # primary: своя affinity ИЛИ 'any'. Полный фрагмент, а не "provider_affinity IN": + # с #3299 такой же IN есть и в backup-подзапросе fallback'а. + if "provider_affinity IN (:provider, 'any')" in sql: cands = [ r for r in self.rows @@ -189,17 +191,27 @@ class FakeSession: # тоже гейтим по подстроке (b2-подзапрос), иначе мок «чинил» бы # незащищённый SQL сам. backup_must_be_usable = "b2.banned_until > now()" in sql + # #3299: backup — 'any' или та же affinity, здоровый, с живой арендой; гейт + # по подстроке, как выше — иначе мок сам «чинил» бы старый предикат. + counts_any = "IN (sp.provider_affinity, 'any')" in sql def _has_backup(row: dict[str, Any]) -> bool: if row["provider_affinity"] == "any": return True + affinities = ( + (row["provider_affinity"], "any") + if counts_any + else (row["provider_affinity"],) + ) return any( - other["provider_affinity"] == row["provider_affinity"] + other["provider_affinity"] in affinities and other["enabled"] + and (not counts_any or other["consecutive_fails"] < max_fails) + and (not counts_any or _not_expired(other)) and other["id"] != row["id"] and not ( backup_must_be_usable - and self._has_active_ban(other["id"], other["provider_affinity"]) + and self._has_active_ban(other["id"], row["provider_affinity"]) ) for other in self.rows ) @@ -392,13 +404,24 @@ class FakeSession: # fallback-safe: другой ПРИГОДНЫЙ узел ТОЙ ЖЕ affinity (банимый узел # остаётся enabled и тоже считается — бан теперь per-source; а вот # забаненный своим же источником backup'ом не считается). + # #3299: тот же гейт, что в acquire-ветке мока. + counts_any = "IN (sp.provider_affinity, 'any')" in sql + affinities = ( + (sp["provider_affinity"], "any") if counts_any else (sp["provider_affinity"],) + ) return any( - other["provider_affinity"] == sp["provider_affinity"] + other["provider_affinity"] in affinities and other["enabled"] + and (not counts_any or other["consecutive_fails"] < max_fails) + and ( + not counts_any + or other.get("expires_at") is None + or other["expires_at"] > datetime.now(UTC) + ) and other["id"] != sp["id"] and not ( backup_must_be_usable - and self._has_active_ban(other["id"], other["provider_affinity"]) + and self._has_active_ban(other["id"], sp["provider_affinity"]) ) for other in self.rows ) @@ -712,11 +735,16 @@ def test_acquire_fallback_skips_expired_lease() -> None: def test_acquire_fallback_prefers_non_expired_over_expired() -> None: - """Просроченный узел пропускается, живой той же чужой affinity — выдан fallback'ом.""" + """Просроченный узел пропускается, живой той же чужой affinity — выдан fallback'ом. + + Узел 3 (занят чужим прогоном) — резерв cian: с #3299 просроченный узел 1 резервом + не считается, и без узла 3 живой узел 2 был бы последним для cian и защищён. + """ db = FakeSession( [ _proxy(1, affinity="cian", expires_at=datetime.now(UTC) - timedelta(minutes=1)), _proxy(2, affinity="cian", expires_at=datetime.now(UTC) + timedelta(hours=1)), + _proxy(3, affinity="cian", leased_by=99), ] ) lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type] @@ -735,9 +763,9 @@ def test_acquire_warns_on_expired_proxy(caplog: pytest.LogCaptureFixture) -> Non with caplog.at_level("WARNING", logger="app.services.proxy_pool"): lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type] assert lease is not None and lease.id == 2 # живой узел всё равно выдан - assert any( - "expired" in r.message and "id=1" in r.message for r in caplog.records - ), "ожидался WARNING про просроченный proxy id=1" + assert any("expired" in r.message and "id=1" in r.message for r in caplog.records), ( + "ожидался WARNING про просроченный proxy id=1" + ) # ── release ────────────────────────────────────────────────────────────────── @@ -1549,9 +1577,7 @@ async def test_probe_proxy_proxy_error_logs_single_line_without_traceback( monkeypatch.setattr(httpx.AsyncClient, "get", _broken_get) with caplog.at_level("WARNING", logger="app.services.proxy_pool"): - ok, exit_ip, latency_ms, fail_kind = await proxy_pool._probe_proxy( - "http://u:p@h1:8080" - ) + ok, exit_ip, latency_ms, fail_kind = await proxy_pool._probe_proxy("http://u:p@h1:8080") assert ok is False assert exit_ip is None diff --git a/tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py b/tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py new file mode 100644 index 00000000..002ce1a5 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py @@ -0,0 +1,181 @@ +"""Защита «последнего узла выделенной affinity» считает узлы 'any' (#3299). + +Запасной заход `acquire()` и защита в `mark_banned()` спрашивали «есть ли у привязки +ВТОРОЙ выделенный узел». Выделенный узел по смыслу штучный, поэтому ответ почти всегда +«нет», и такой узел не выдавался НИКОМУ, даже когда свой источник спокойно обслуживали +'any'-узлы. Прод 17.09: узел 15 (avito) свободен и здоров, для cian/domclick fallback +его не отдавал; 30.08 так же лёг добор Домклика (прогоны 5449-5459). + +Проверяется SQL, а не его пересказ: живой Postgres со схемой из миграций (в CI его +поднимает ci-tradein.yml), всё в одной внешней транзакции с откатом — `commit()` внутри +proxy_pool освобождает savepoint. Чужие строки scrape_proxies на время теста выключены +в той же транзакции. Без БД — skip, как соседние live-тесты. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from collections.abc import Iterator +from typing import Any + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.orm import Session + +from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS, acquire, mark_banned + + +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 1 FROM scrape_proxies 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 db() -> Iterator[Session]: + assert _ENGINE is not None + conn = _ENGINE.connect() + outer = conn.begin() + session = Session(bind=conn, join_transaction_mode="create_savepoint") + session.execute(text("UPDATE scrape_proxies SET enabled = false")) + # commit здесь и в хелперах ниже — это RELEASE savepoint'а, внешняя транзакция цела. + # Без него db.rollback() внутри acquire (пустой отбор) откатил бы и подготовку. + session.commit() + try: + yield session + finally: + session.close() + outer.rollback() + conn.close() + + +def _node(db: Session, affinity: str, *, fails: int = 0) -> int: + proxy_id = db.execute( + text( + """ + INSERT INTO scrape_proxies (url, provider_affinity, consecutive_fails) + VALUES ('http://t3299-' || gen_random_uuid(), :aff, :fails) + RETURNING id + """ + ), + {"aff": affinity, "fails": fails}, + ).scalar_one() + db.commit() + return int(proxy_id) + + +def _ban(db: Session, proxy_id: int, source: str) -> None: + db.execute( + text( + """ + INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason) + VALUES (:id, :source, now() + interval '6 hours', 'banned:' || :source) + """ + ), + {"id": proxy_id, "source": source}, + ) + db.commit() + + +def _leased_by(db: Session, proxy_id: int) -> int | None: + return db.execute( + text("SELECT leased_by FROM scrape_proxies WHERE id = :id"), {"id": proxy_id} + ).scalar_one() + + +# ── acquire: запасной заход ────────────────────────────────────────────────── + + +def test_dedicated_node_goes_to_other_source_when_any_node_backs_it(db: Session) -> None: + """Замеренный расклад issue: выделенный avito-узел + 'any'-узел, которого домклик не + получает (забанен им). Домклик берёт выделенный узел, avito остаётся с 'any'. + + На старом предикате acquire('domclick') возвращал None.""" + dedicated = _node(db, "avito") + any_node = _node(db, "any") + _ban(db, any_node, "domclick") + + lease = acquire(db, "domclick", run_id=None) + + assert lease is not None and lease.id == dedicated, f"домклику выдан {lease}" + own = acquire(db, "avito", run_id=None) + assert own is not None and own.id == any_node, f"avito остался без узла: {own}" + + +def test_only_dedicated_node_without_any_nodes_stays_protected(db: Session) -> None: + """Исходный смысл защиты (#2600): единственный узел привязки и ни одного 'any' — + fallback его не забирает, узел остаётся за своим источником.""" + dedicated = _node(db, "avito") + + assert acquire(db, "domclick", run_id=None) is None + assert _leased_by(db, dedicated) is None + own = acquire(db, "avito", run_id=None) + assert own is not None and own.id == dedicated + + +def test_any_node_banned_by_dedicated_source_is_not_a_backup(db: Session) -> None: + """'any'-узел, забаненный самим avito, avito не обслужит — резервом не считается.""" + dedicated = _node(db, "avito") + any_node = _node(db, "any") + _ban(db, any_node, "domclick") + _ban(db, any_node, "avito") + + assert acquire(db, "domclick", run_id=None) is None + assert _leased_by(db, dedicated) is None + + +def test_unhealthy_any_node_is_not_a_backup(db: Session) -> None: + """Узел в карантине по consecutive_fails acquire не выдаёт — резервом он тоже не + считается (раньше здоровье резерва не проверялось вовсе).""" + dedicated = _node(db, "avito") + _node(db, "any", fails=MAX_CONSECUTIVE_FAILS) + + assert acquire(db, "domclick", run_id=None) is None + assert _leased_by(db, dedicated) is None + + +# ── mark_banned: защита считает так же, как acquire ────────────────────────── + + +def test_mark_banned_counts_dedicated_node_reachable_via_any_backup(db: Session) -> None: + """Банится 'any'-узел для cian. У cian остаётся выделенный avito-узел, за которым + стоит 'any'-резерв — значит бан безопасен, и cian действительно получает этот узел. + + На старом предикате защита ответила бы 'protected' и оставила бы cian ходить через + отбитый узел, хотя живой узел для него был.""" + banned = _node(db, "any") + dedicated = _node(db, "avito") + backup = _node(db, "any") + _ban(db, backup, "cian") + + assert mark_banned(db, banned, source="cian") == "banned" + lease = acquire(db, "cian", run_id=None) + assert lease is not None and lease.id == dedicated + + +def test_mark_banned_still_protects_when_dedicated_node_has_no_usable_backup( + db: Session, +) -> None: + """Зеркало: резерв выделенного узла забанен самим avito → для cian узла нет, бан не + пишется, и отбитый узел продолжает выдаваться cian (обещание из лога защиты).""" + banned = _node(db, "any") + _node(db, "avito") + _ban(db, banned, "avito") + + assert mark_banned(db, banned, source="cian") == "protected" + lease = acquire(db, "cian", run_id=None) + assert lease is not None and lease.id == banned