fix(tradein/proxy_pool): резерв выделенного узла считает узлы 'any' (#3299)

Запасной заход acquire() и защита последнего узла в mark_banned() спрашивали,
есть ли у выделенной привязки ВТОРОЙ узел той же привязки. Выделенный узел
штучный, поэтому ответ почти всегда «нет», и такой узел не выдавался никому,
хотя свой источник обслуживали 'any'-узлы. Прод 17.09: узел 15 (avito)
свободен и здоров, а cian/domclick его не получали; 30.08 так лёг добор
Домклика (прогоны 5449-5459).

Теперь резерв — любой узел той же привязки или 'any', здоровый
(consecutive_fails), с живой арендой порта и не забаненный источником
выделенной привязки. Предикат одинаковый в обоих местах. leased_by
намеренно не проверяется: аренда временная.

Тесты на живом Postgres (tests/test_3299_*): 6 случаев из приёмки, два
красные на старом SQL. Мок test_proxy_pool.py: маршрутизация primary/fallback
по полному фрагменту IN (:provider, 'any') и гейт нового предиката.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
bot-backend 2026-09-17 12:59:44 +05:00
parent 34642e1dd5
commit 92295971fe
3 changed files with 246 additions and 22 deletions

View file

@ -275,12 +275,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
чужая только запасной вариант, чтобы источник не голодал при живых свободных узлах чужая только запасной вариант, чтобы источник не голодал при живых свободных узлах
чужой affinity (#2600). чужой affinity (#2600).
Fallback НЕ трогает последний enabled-узел выделенной (не-'any') affinity: если Fallback НЕ трогает последний пригодный узел выделенной (не-'any') affinity: если
fallback заберёт его под чужой источник, «свой» останется без прокси вообще хуже, fallback заберёт его под чужой источник, «свой» останется без прокси вообще хуже,
чем голодание исходного источника, которое фикс призван устранить. Кандидат чем голодание исходного источника, которое фикс призван устранить. Кандидат
участвует в fallback, только если его affinity='any' ИЛИ у этой affinity есть ДРУГОЙ участвует в fallback, только если его affinity='any' ИЛИ у источника этой affinity
enabled-узел (EXISTS-подзапрос) т.е. выдача не обнулит доступность выделенной без него останется ДРУГОЙ кандидат (EXISTS-подзапрос): узел той же affinity или
affinity целиком. 'any', здоровый, с живой арендой порта и не забаненный этим источником (#3299 —
до него 'any'-узлы не считались, и единственный выделенный узел не выдавался никому).
Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql
единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR
@ -401,17 +402,29 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
-- fallback увести последний реально рабочий узел выделенной -- fallback увести последний реально рабочий узел выделенной
-- affinity и обрушить её (два domclick-узла, один забанен -- affinity и обрушить её (два domclick-узла, один забанен
-- domclick'ом → второй уходит под avito → domclick без прокси). -- domclick'ом → второй уходит под avito → domclick без прокси).
--
-- #3299: вопрос «останется ли у источника sp.provider_affinity
-- хоть один кандидат без sp», а не «есть ли ВТОРОЙ узел той же
-- привязки». 'any'-узлы основной запрос выдаёт выделенному
-- источнику наравне, значит и backup'ом они считаются; прежний
-- `= sp.provider_affinity` прятал единственный выделенный узел от
-- всех. Здоровье и срок аренды как в основном запросе.
-- leased_by НЕ проверяем: аренда вернётся через минуты, а
-- отказ из-за чужой аренды прятал бы узел при любом параллельном
-- прогоне (тот же довод, что у защиты в mark_banned).
OR EXISTS ( OR EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxies AS other 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.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 other.id <> sp.id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxy_source_bans b2 FROM scrape_proxy_source_bans b2
WHERE b2.proxy_id = other.id WHERE b2.proxy_id = other.id
AND b2.source = other.provider_affinity AND b2.source = sp.provider_affinity
AND b2.banned_until > now() AND b2.banned_until > now()
) )
) )
@ -535,8 +548,7 @@ def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None:
db.commit() db.commit()
if row is None: if row is None:
logger.debug( logger.debug(
"proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already " "proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already finalized?)",
"finalized?)",
run_id, run_id,
) )
except Exception: except Exception:
@ -999,17 +1011,22 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None =
-- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом -- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом
-- не считается (иначе защита сочла бы affinity живой, когда она -- не считается (иначе защита сочла бы affinity живой, когда она
-- уже нет). -- уже нет).
-- #3299: предикат backup'а — буква в букву как в fallback
-- acquire() (там же обоснование): 'any'-узлы в счёт, здоровье
-- и срок аренды проверяются.
OR EXISTS ( OR EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxies other 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.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 other.id <> sp.id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxy_source_bans b2 FROM scrape_proxy_source_bans b2
WHERE b2.proxy_id = other.id WHERE b2.proxy_id = other.id
AND b2.source = other.provider_affinity AND b2.source = sp.provider_affinity
AND b2.banned_until > now() AND b2.banned_until > now()
) )
) )

View file

@ -167,7 +167,9 @@ class FakeSession:
exp = row.get("expires_at") exp = row.get("expires_at")
return exp is None or exp > datetime.now(UTC) 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 = [ cands = [
r r
for r in self.rows for r in self.rows
@ -189,17 +191,27 @@ class FakeSession:
# тоже гейтим по подстроке (b2-подзапрос), иначе мок «чинил» бы # тоже гейтим по подстроке (b2-подзапрос), иначе мок «чинил» бы
# незащищённый SQL сам. # незащищённый SQL сам.
backup_must_be_usable = "b2.banned_until > now()" in 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: def _has_backup(row: dict[str, Any]) -> bool:
if row["provider_affinity"] == "any": if row["provider_affinity"] == "any":
return True return True
affinities = (
(row["provider_affinity"], "any")
if counts_any
else (row["provider_affinity"],)
)
return any( return any(
other["provider_affinity"] == row["provider_affinity"] other["provider_affinity"] in affinities
and other["enabled"] 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 other["id"] != row["id"]
and not ( and not (
backup_must_be_usable 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 for other in self.rows
) )
@ -392,13 +404,24 @@ class FakeSession:
# fallback-safe: другой ПРИГОДНЫЙ узел ТОЙ ЖЕ affinity (банимый узел # fallback-safe: другой ПРИГОДНЫЙ узел ТОЙ ЖЕ affinity (банимый узел
# остаётся enabled и тоже считается — бан теперь per-source; а вот # остаётся enabled и тоже считается — бан теперь per-source; а вот
# забаненный своим же источником backup'ом не считается). # забаненный своим же источником 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( return any(
other["provider_affinity"] == sp["provider_affinity"] other["provider_affinity"] in affinities
and other["enabled"] 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 other["id"] != sp["id"]
and not ( and not (
backup_must_be_usable 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 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: def test_acquire_fallback_prefers_non_expired_over_expired() -> None:
"""Просроченный узел пропускается, живой той же чужой affinity — выдан fallback'ом.""" """Просроченный узел пропускается, живой той же чужой affinity — выдан fallback'ом.
Узел 3 (занят чужим прогоном) резерв cian: с #3299 просроченный узел 1 резервом
не считается, и без узла 3 живой узел 2 был бы последним для cian и защищён.
"""
db = FakeSession( db = FakeSession(
[ [
_proxy(1, affinity="cian", expires_at=datetime.now(UTC) - timedelta(minutes=1)), _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(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] 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"): with caplog.at_level("WARNING", logger="app.services.proxy_pool"):
lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type] lease = acquire(db, "avito", run_id=1) # type: ignore[arg-type]
assert lease is not None and lease.id == 2 # живой узел всё равно выдан assert lease is not None and lease.id == 2 # живой узел всё равно выдан
assert any( assert any("expired" in r.message and "id=1" in r.message for r in caplog.records), (
"expired" in r.message and "id=1" in r.message for r in caplog.records "ожидался WARNING про просроченный proxy id=1"
), "ожидался WARNING про просроченный proxy id=1" )
# ── release ────────────────────────────────────────────────────────────────── # ── release ──────────────────────────────────────────────────────────────────
@ -1549,9 +1577,7 @@ async def test_probe_proxy_proxy_error_logs_single_line_without_traceback(
monkeypatch.setattr(httpx.AsyncClient, "get", _broken_get) monkeypatch.setattr(httpx.AsyncClient, "get", _broken_get)
with caplog.at_level("WARNING", logger="app.services.proxy_pool"): with caplog.at_level("WARNING", logger="app.services.proxy_pool"):
ok, exit_ip, latency_ms, fail_kind = await proxy_pool._probe_proxy( ok, exit_ip, latency_ms, fail_kind = await proxy_pool._probe_proxy("http://u:p@h1:8080")
"http://u:p@h1:8080"
)
assert ok is False assert ok is False
assert exit_ip is None assert exit_ip is None

View file

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