Merge remote-tracking branch 'origin/main' into fix/backtest-house-type
All checks were successful
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 21s
CI / changes (pull_request) Successful in 33s
CI Trade-In / backend-tests (pull_request) Successful in 8m41s
All checks were successful
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 21s
CI / changes (pull_request) Successful in 33s
CI Trade-In / backend-tests (pull_request) Successful in 8m41s
This commit is contained in:
commit
86629dc9f4
21 changed files with 1045 additions and 88 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
)
|
||||
)
|
||||
|
|
@ -446,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(
|
||||
|
|
@ -486,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` — ПОСЛЕДНИЙ использованный узел (перезаписывается при
|
||||
каждой новой выдаче); полная цепочка узлов, если она менялась, — в
|
||||
|
|
@ -535,8 +549,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:
|
||||
|
|
@ -983,6 +996,10 @@ def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None =
|
|||
WHERE sp.id <> CAST(:proxy_id AS bigint)
|
||||
AND sp.enabled
|
||||
AND sp.consecutive_fails < CAST(:max_fails AS integer)
|
||||
-- Как в acquire(): узел с истёкшей арендой порта остаётся enabled,
|
||||
-- но не выдаётся. Без этой строки он «спасал» источник, и бан уходил
|
||||
-- последнему живому узлу (ревью #3565 к #3299).
|
||||
AND (sp.expires_at IS NULL OR sp.expires_at > now())
|
||||
AND NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM scrape_proxy_source_bans b
|
||||
|
|
@ -999,17 +1016,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()
|
||||
)
|
||||
)
|
||||
|
|
|
|||
|
|
@ -0,0 +1,29 @@
|
|||
-- 320_listings_domclick_zero_area_to_null.sql
|
||||
-- Issue #3252: Домклик отдаёт `livingArea: 0` / `kitchenArea: 0` как «не указано»
|
||||
-- (живая карточка 2078257603, 17.09: обе площади 0, ремонт пустой). Парсер карточки
|
||||
-- писал этот 0 в колонку как настоящую площадь; COALESCE(:new, old) в UPDATE при этом
|
||||
-- затирал нулём ранее известное значение. Писатель починен в том же PR
|
||||
-- (providers/domclick/detail.py::_pos_float), но починка разбора строк не чинит:
|
||||
-- COALESCE(NULL, 0) оставит 0 навсегда.
|
||||
--
|
||||
-- ЗАМЕР ПРОДА 2026-09-17 (SELECT source, count(*) FILTER (WHERE kitchen_area_m2=0) ...):
|
||||
-- domklik: kitchen_area_m2 = 0 — 303 строки, living_area_m2 = 0 — 84 строки;
|
||||
-- все обогащены карточкой 26.08–14.09 (detail_enriched_at), ни одной до этого.
|
||||
-- avito / cian / yandex / n1: нулей нет вовсе.
|
||||
-- Ожидаемо тронуто: 303 + 84 строки. Отрицательных площадей нет ни у кого.
|
||||
--
|
||||
-- Idempotent: повторный прогон — 0 строк.
|
||||
|
||||
BEGIN;
|
||||
|
||||
UPDATE listings
|
||||
SET kitchen_area_m2 = NULL
|
||||
WHERE source = 'domklik'
|
||||
AND kitchen_area_m2 = 0;
|
||||
|
||||
UPDATE listings
|
||||
SET living_area_m2 = NULL
|
||||
WHERE source = 'domklik'
|
||||
AND living_area_m2 = 0;
|
||||
|
||||
COMMIT;
|
||||
|
|
@ -83,9 +83,9 @@ _SSR_LITERAL = """{
|
|||
},
|
||||
"legalOptions": {"saleType": "Свободная продажа"},
|
||||
"egrnData": {
|
||||
"area": 38.2,
|
||||
"floor": 5,
|
||||
"owners_count": 1,
|
||||
"area": {"status": "success", "value": 38.2},
|
||||
"floor": {"status": "success", "value": 5},
|
||||
"owners_count": {"status": "success", "value": 1},
|
||||
"collateral": true,
|
||||
"collateral_sber": false
|
||||
},
|
||||
|
|
@ -264,7 +264,7 @@ def test_parse_detail_html_raw_extra() -> None:
|
|||
# wall/floor type живут ТОЛЬКО в raw_extra, НЕ в listings.house_type.
|
||||
assert "house_type" not in raw
|
||||
assert raw["domclick_building_guid"] == "abc-guid-123"
|
||||
assert raw["egrn_area"] == 38.2
|
||||
assert raw["egrn_area"] == {"status": "success", "value": 38.2}
|
||||
assert raw["demand"]["calls"] == 5
|
||||
assert raw["demand"]["favorites"] == 12
|
||||
# AVM (Layer C, top-level pricePrediction) → raw_extra.avm
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -147,6 +147,34 @@ tests/test_3463_db_timeouts.py::test_lock_wait_over_ceiling_is_aborted
|
|||
tests/test_3463_db_timeouts.py::test_set_local_statement_timeout_overrides_session_ceiling
|
||||
tests/test_3463_db_timeouts.py::test_set_local_is_scoped_to_its_transaction
|
||||
|
||||
# Жилая площадь и балконы переживают переобход выдачи (#3252) — тот же `_live_session()`.
|
||||
# Проверяют ПОВЕДЕНИЕ апсерта на живой схеме: SERP без полей не стирает то, что
|
||||
# добыла карточка (COALESCE), и настоящее новое значение всё ещё перезаписывает.
|
||||
# В ci-tradein.yml бегут по-настоящему (postgres-сервис); без БД — skip.
|
||||
tests/test_3252_domclick_card_fields.py::test_serp_rescrape_keeps_detail_living_area_and_balconies
|
||||
tests/test_3252_domclick_card_fields.py::test_real_new_value_still_overwrites
|
||||
|
||||
# Пул прокси (#3299, #3310, #3404) — предикаты выдачи/защиты и атрибуция прогона живут
|
||||
# в SQL, мок их не проверяет. Live-тесты на схеме из миграций: #3299 и #3310 — в одной
|
||||
# внешней транзакции с откатом (commit пула = RELEASE savepoint), #3404 — настоящими
|
||||
# сессиями, строки t3404-* удаляются в finally. В ci-tradein.yml идут по-настоящему
|
||||
# (postgres-сервис, #2745); на машине без БД и в deploy-tradein.yml — skip.
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_dedicated_node_goes_to_other_source_when_any_node_backs_it
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_only_dedicated_node_without_any_nodes_stays_protected
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_any_node_banned_by_dedicated_source_is_not_a_backup
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_unhealthy_any_node_is_not_a_backup
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_expired_any_node_is_not_a_backup
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_mark_banned_counts_dedicated_node_reachable_via_any_backup
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_mark_banned_still_protects_when_dedicated_node_has_no_usable_backup
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_mark_banned_unhealthy_backup_does_not_count
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_mark_banned_expired_backup_does_not_count
|
||||
tests/test_3299_fallback_counts_any_nodes.py::test_mark_banned_expired_node_does_not_save_the_source
|
||||
tests/test_3310_curl_ban_keeps_last_node.py::test_three_bans_on_last_node_keep_it_in_the_pool
|
||||
tests/test_3310_curl_ban_keeps_last_node.py::test_transport_failure_still_counts_against_node_health
|
||||
tests/test_3404_egress_run_attribution.py::test_resolve_within_run_writes_node_to_scrape_runs
|
||||
tests/test_3404_egress_run_attribution.py::test_resolve_outside_run_writes_nothing
|
||||
tests/test_3404_egress_run_attribution.py::test_attribution_does_not_commit_callers_transaction
|
||||
|
||||
# Миграция 310 (#3385): настоящий SQL-файл с DELETE и RAISE-остановкой исполняет только
|
||||
# Postgres. В ci-tradein.yml бегут по-настоящему (postgres-сервис, #2745); краснеют от
|
||||
# снятия подписи батча, сужения окна до ×10 и снятия порога — проверено вручную 17.09.
|
||||
|
|
|
|||
|
|
@ -103,7 +103,7 @@ async def test_403_bans_the_node_for_cian_only() -> None:
|
|||
with pytest.raises(CianBlockedError):
|
||||
await _fetch(403, spy)
|
||||
assert spy.mark_banned_calls == [(1, "cian")]
|
||||
assert spy.mark_health_calls == [(1, False)]
|
||||
assert spy.mark_health_calls == [] # #3310: отказ площадки не копит consecutive_fails
|
||||
assert spy.release_calls == [1] # lease не течёт даже на бане
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -179,7 +179,7 @@ async def test_price_history_403_bans_the_node_for_cian() -> None:
|
|||
pool = _SpyPool()
|
||||
result = await _run_price_history(pool, status_code=403)
|
||||
assert pool.mark_banned_calls == [(9, "cian")]
|
||||
assert pool.mark_health_calls == [(9, False)]
|
||||
assert pool.mark_health_calls == [] # #3310: бан вместо mark_health(False)
|
||||
assert result.errors == 1 # прогон честен: отказ посчитан
|
||||
|
||||
|
||||
|
|
@ -335,7 +335,7 @@ async def test_zhk_resolve_403_reaches_the_pool() -> None:
|
|||
with pytest.raises(CianBlockedError):
|
||||
await _resolve(403, spy)
|
||||
assert spy.mark_banned_calls == [(9, "cian")]
|
||||
assert spy.mark_health_calls == [(9, False)]
|
||||
assert spy.mark_health_calls == [] # #3310: бан вместо mark_health(False)
|
||||
assert spy.release_calls == [9]
|
||||
|
||||
|
||||
|
|
|
|||
192
tradein-mvp/backend/tests/test_3252_domclick_card_fields.py
Normal file
192
tradein-mvp/backend/tests/test_3252_domclick_card_fields.py
Normal file
|
|
@ -0,0 +1,192 @@
|
|||
"""#3252: три тихие дыры в полях карточки Домклика.
|
||||
|
||||
1. owners_count. egrnData отдаёт поля ЕГРН обёрткой ``{"status", "value"}``, а парсер
|
||||
делал ``int()`` от словаря и молча получал None: на проде 0 из ~2970 карточек,
|
||||
обогащённых с 29.08 (при 5075 из 6296 у ручного прогона 18.07, который обёртку
|
||||
разворачивал).
|
||||
2. Нулевые площади. ``livingArea: 0`` / ``kitchenArea: 0`` у Домклика значит «не
|
||||
указано», а в колонку ложился 0.00 (303 кухни и 84 жилые площади).
|
||||
3. Жилая площадь и балконы стирались переобходом выдачи: апсерт в scraper_kit.base
|
||||
писал их сырым ``EXCLUDED``, а SERP Домклика и Авито этих полей не отдаёт. У
|
||||
domklik, переобойдённых после обогащения, living_area_m2 = 0 из 5217.
|
||||
|
||||
Состояние карточек ниже — урезанный ``__SSR_STATE__`` живых карточек 2074362051 и
|
||||
2078257603, снятый 17.09 (значения как есть, лишние ветки выкинуты).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import uuid
|
||||
from decimal import Decimal
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pytest
|
||||
from scraper_kit.base import ScrapedLot, save_listings
|
||||
from scraper_kit.providers.domclick.detail import parse_detail_html
|
||||
from sqlalchemy import text
|
||||
|
||||
_LIVE_2074362051 = {
|
||||
"productCard": {
|
||||
"objectInfo": {
|
||||
"area": 59.8,
|
||||
"balconies": 1,
|
||||
"kitchenArea": 14.9,
|
||||
"livingArea": 25.8,
|
||||
"renovation": "Косметический",
|
||||
},
|
||||
"egrnData": {
|
||||
"area": {"status": "success", "value": 59.8},
|
||||
"collateral": False,
|
||||
"collateral_sber": False,
|
||||
"floor": {"status": "success", "value": 14},
|
||||
"owners_count": {"status": "error", "value": 6},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
_LIVE_2078257603 = {
|
||||
"productCard": {
|
||||
"objectInfo": {
|
||||
"area": 67.4,
|
||||
"balconies": 0,
|
||||
"kitchenArea": 0,
|
||||
"livingArea": 0,
|
||||
"renovation": "",
|
||||
},
|
||||
"egrnData": {
|
||||
"area": {"status": "error", "value": 65.3},
|
||||
"collateral": False,
|
||||
"collateral_sber": False,
|
||||
"floor": {"status": "success", "value": 21},
|
||||
"owners_count": {"status": "success", "value": 1},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def _parse(state: dict, card_id: str):
|
||||
html = f"<script>window.__SSR_STATE__ = {json.dumps(state, ensure_ascii=False)};</script>"
|
||||
return parse_detail_html(html, f"https://ekaterinburg.domclick.ru/card/sale__flat__{card_id}")
|
||||
|
||||
|
||||
# ── 1-2. Разбор живой формы ──────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_owners_count_is_read_from_egrn_wrapper_whatever_the_status() -> None:
|
||||
# status=error у owners_count = «много собственников», value — сам факт ЕГРН.
|
||||
assert _parse(_LIVE_2074362051, "2074362051").owners_count == 6
|
||||
assert _parse(_LIVE_2078257603, "2078257603").owners_count == 1
|
||||
|
||||
|
||||
def test_egrn_area_keeps_wrapper_in_raw_payload() -> None:
|
||||
e = _parse(_LIVE_2078257603, "2078257603")
|
||||
assert e.raw_extra["egrn_area"] == {"status": "error", "value": 65.3}
|
||||
|
||||
|
||||
def test_zero_areas_mean_unknown_not_zero() -> None:
|
||||
zero = _parse(_LIVE_2078257603, "2078257603")
|
||||
assert (zero.living_area_m2, zero.kitchen_area_m2) == (None, None)
|
||||
real = _parse(_LIVE_2074362051, "2074362051")
|
||||
assert (real.living_area_m2, real.kitchen_area_m2) == (25.8, 14.9)
|
||||
|
||||
|
||||
# ── 3. Переобход выдачи не стирает то, что добыла карточка (живой Postgres) ──
|
||||
|
||||
|
||||
def _live_session() -> Any | None:
|
||||
try:
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "")
|
||||
if not dsn or "localhost:5432/test" in dsn:
|
||||
return None
|
||||
engine = create_engine(dsn, future=True)
|
||||
with engine.connect() as conn:
|
||||
conn.execute(text("SELECT 1"))
|
||||
return sessionmaker(bind=engine, future=True)()
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _matcher() -> MagicMock:
|
||||
m = MagicMock()
|
||||
m.match_or_create_house.return_value = (None, 0.0, "no_address")
|
||||
m.upsert_listing_source.return_value = None
|
||||
return m
|
||||
|
||||
|
||||
def _serp_lot(sid: str, **detail: Any) -> ScrapedLot:
|
||||
return ScrapedLot(
|
||||
source="domklik",
|
||||
source_url=f"https://ekaterinburg.domclick.ru/card/sale__flat__t3252{sid}",
|
||||
source_id=f"t3252-{sid}",
|
||||
price_rub=7_000_000,
|
||||
**detail,
|
||||
)
|
||||
|
||||
|
||||
def _row(db: Any, sid: str) -> Any:
|
||||
return db.execute(
|
||||
text(
|
||||
"SELECT living_area_m2, balconies_count FROM listings "
|
||||
"WHERE source='domklik' AND source_id = :sid"
|
||||
),
|
||||
{"sid": f"t3252-{sid}"},
|
||||
).fetchone()
|
||||
|
||||
|
||||
def _cleanup(db: Any) -> None:
|
||||
try:
|
||||
db.rollback()
|
||||
ids = "(SELECT id FROM listings WHERE source='domklik' AND source_id LIKE 't3252-%')"
|
||||
db.execute(text(f"DELETE FROM listings_snapshots WHERE listing_id IN {ids}"))
|
||||
db.execute(text("DELETE FROM listing_sources WHERE ext_id LIKE 't3252-%'"))
|
||||
db.execute(text("DELETE FROM listings WHERE source='domklik' AND source_id LIKE 't3252-%'"))
|
||||
db.commit()
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||
def test_serp_rescrape_keeps_detail_living_area_and_balconies() -> None:
|
||||
db = _live_session()
|
||||
sid = uuid.uuid4().hex[:8]
|
||||
try:
|
||||
save_listings(
|
||||
db,
|
||||
[_serp_lot(sid, living_area_m2=25.8, balconies_count=1)],
|
||||
matcher=_matcher(),
|
||||
region_code=66,
|
||||
)
|
||||
save_listings(db, [_serp_lot(sid)], matcher=_matcher(), region_code=66) # выдача: полей нет
|
||||
row = _row(db, sid)
|
||||
assert (row.living_area_m2, row.balconies_count) == (Decimal("25.80"), 1)
|
||||
finally:
|
||||
_cleanup(db)
|
||||
|
||||
|
||||
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||
def test_real_new_value_still_overwrites() -> None:
|
||||
db = _live_session()
|
||||
sid = uuid.uuid4().hex[:8]
|
||||
try:
|
||||
save_listings(
|
||||
db,
|
||||
[_serp_lot(sid, living_area_m2=25.8, balconies_count=1)],
|
||||
matcher=_matcher(),
|
||||
region_code=66,
|
||||
)
|
||||
save_listings(
|
||||
db,
|
||||
[_serp_lot(sid, living_area_m2=30.5, balconies_count=2)],
|
||||
matcher=_matcher(),
|
||||
region_code=66,
|
||||
)
|
||||
row = _row(db, sid)
|
||||
assert (row.living_area_m2, row.balconies_count) == (Decimal("30.50"), 2)
|
||||
finally:
|
||||
_cleanup(db)
|
||||
231
tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py
Normal file
231
tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py
Normal file
|
|
@ -0,0 +1,231 @@
|
|||
"""Защита «последнего узла выделенной 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, expired: bool = False) -> int:
|
||||
proxy_id = db.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO scrape_proxies (url, provider_affinity, consecutive_fails, expires_at)
|
||||
VALUES ('http://t3299-' || gen_random_uuid(), :aff, :fails,
|
||||
CASE WHEN :expired THEN now() - interval '1 minute' END)
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
{"aff": affinity, "fails": fails, "expired": expired},
|
||||
).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
|
||||
|
||||
|
||||
def test_expired_any_node_is_not_a_backup(db: Session) -> None:
|
||||
"""Узел с истёкшей арендой порта acquire не выдаёт — резервом он тоже не считается.
|
||||
|
||||
Отдельно от карантина: у предиката резерва два независимых условия пригодности, и
|
||||
каждое проверяется своим узлом (ревью PR #3565: без этого теста снятие проверки
|
||||
`expires_at` из подзапроса оставляло сьют зелёным)."""
|
||||
dedicated = _node(db, "avito")
|
||||
_node(db, "any", expired=True)
|
||||
|
||||
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
|
||||
|
||||
|
||||
def test_mark_banned_unhealthy_backup_does_not_count(db: Session) -> None:
|
||||
"""Резерв выделенного узла в карантине по consecutive_fails — для cian узла нет,
|
||||
защита держит. Сам резерв кандидатом cian не считается по той же причине (карантин
|
||||
проверяется и во внешнем отборе), так что ответ решает именно предикат резерва."""
|
||||
banned = _node(db, "any")
|
||||
_node(db, "avito")
|
||||
_ban(db, banned, "avito")
|
||||
_node(db, "any", fails=MAX_CONSECUTIVE_FAILS)
|
||||
|
||||
assert mark_banned(db, banned, source="cian") == "protected"
|
||||
|
||||
|
||||
def test_mark_banned_expired_backup_does_not_count(db: Session) -> None:
|
||||
"""Резерв выделенного узла с истёкшей арендой — для cian узла нет, защита держит."""
|
||||
banned = _node(db, "any")
|
||||
_node(db, "avito")
|
||||
_ban(db, banned, "avito")
|
||||
_node(db, "any", expired=True)
|
||||
|
||||
assert mark_banned(db, banned, source="cian") == "protected"
|
||||
|
||||
|
||||
def test_mark_banned_expired_node_does_not_save_the_source(db: Session) -> None:
|
||||
"""Внешний отбор: у cian остался только 'any'-узел с истёкшей арендой. acquire его не
|
||||
выдаёт, значит банимый узел последний — бан не пишется, и cian получает этот узел.
|
||||
|
||||
До правки защита засчитывала просроченный узел (пока healthcheck не наберёт ему
|
||||
отказов), бан уходил, а cian оставался ни с чем."""
|
||||
banned = _node(db, "any")
|
||||
_node(db, "any", expired=True)
|
||||
|
||||
assert mark_banned(db, banned, source="cian") == "protected"
|
||||
lease = acquire(db, "cian", run_id=None)
|
||||
assert lease is not None and lease.id == banned
|
||||
126
tradein-mvp/backend/tests/test_3310_curl_ban_keeps_last_node.py
Normal file
126
tradein-mvp/backend/tests/test_3310_curl_ban_keeps_last_node.py
Normal file
|
|
@ -0,0 +1,126 @@
|
|||
"""Защита последнего узла держит слово и на curl-пути (#3310).
|
||||
|
||||
`mark_banned` бережёт узел от бан-строки, если он последний для источника, и пишет
|
||||
«узел продолжит выдаваться». Но `curl_proxy_url` на том же `ProxyBanError` вслед за
|
||||
баном звал `mark_health(ok=False)`: после трёх таких ответов `consecutive_fails = 3`, и
|
||||
`acquire` отсекал узел для ВСЕХ источников — пул объявлял себя пустым при живом узле.
|
||||
Браузерный путь это уже не делает (#3288, `report_platform_ban`), curl-путь — делал.
|
||||
|
||||
Проверка по значению на живом Postgres, настоящим трактом: `curl_proxy_url` →
|
||||
`RealProxyProvider` → `proxy_pool`. Сессии адаптера привязаны к одной внешней
|
||||
транзакции с откатом (commit внутри пула = RELEASE savepoint). Без БД — skip.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from collections.abc import Iterator
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
from scraper_kit.cian_exceptions import CianBlockedError
|
||||
from scraper_kit.orchestration.run_context import current_run_id
|
||||
from scraper_kit.providers._proxy import curl_proxy_url
|
||||
from sqlalchemy import create_engine, text
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.services import proxy_pool, scraper_adapters
|
||||
|
||||
|
||||
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")
|
||||
|
||||
|
||||
@dataclass
|
||||
class _Config:
|
||||
use_proxy_pool_curl: bool = True
|
||||
environment: str = "production"
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def node(monkeypatch: pytest.MonkeyPatch) -> Iterator[tuple[int, Session]]:
|
||||
"""Единственный включённый узел пула ('any'); адаптер ходит в ту же транзакцию."""
|
||||
assert _ENGINE is not None
|
||||
conn = _ENGINE.connect()
|
||||
outer = conn.begin()
|
||||
|
||||
def _session() -> Session:
|
||||
return Session(bind=conn, join_transaction_mode="create_savepoint")
|
||||
|
||||
db = _session()
|
||||
db.execute(text("UPDATE scrape_proxies SET enabled = false"))
|
||||
proxy_id = db.execute(
|
||||
text(
|
||||
"INSERT INTO scrape_proxies (url, provider_affinity) "
|
||||
"VALUES ('http://t3310-' || gen_random_uuid(), 'any') RETURNING id"
|
||||
)
|
||||
).scalar_one()
|
||||
db.commit()
|
||||
monkeypatch.setattr(scraper_adapters, "_SessionLocal", _session)
|
||||
token = current_run_id.set(None)
|
||||
try:
|
||||
yield int(proxy_id), db
|
||||
finally:
|
||||
current_run_id.reset(token)
|
||||
db.close()
|
||||
outer.rollback()
|
||||
conn.close()
|
||||
|
||||
|
||||
def _fails(db: Session, proxy_id: int) -> int:
|
||||
return int(
|
||||
db.execute(
|
||||
text("SELECT consecutive_fails FROM scrape_proxies WHERE id = :id"), {"id": proxy_id}
|
||||
).scalar_one()
|
||||
)
|
||||
|
||||
|
||||
def test_three_bans_on_last_node_keep_it_in_the_pool(node: tuple[int, Session]) -> None:
|
||||
"""Приёмка issue: единственный узел + три бана подряд → acquire не пуст.
|
||||
|
||||
На старом коде: consecutive_fails == 3 и acquire('cian') → None."""
|
||||
proxy_id, db = node
|
||||
provider = scraper_adapters.RealProxyProvider()
|
||||
|
||||
for _ in range(3):
|
||||
with pytest.raises(CianBlockedError):
|
||||
with curl_proxy_url(_Config(), provider, "cian", env_fallback_url=None):
|
||||
raise CianBlockedError("HTTP 403")
|
||||
|
||||
assert _fails(db, proxy_id) == 0
|
||||
ban_rows = db.execute(
|
||||
text("SELECT count(*) FROM scrape_proxy_source_bans WHERE proxy_id = :id"),
|
||||
{"id": proxy_id},
|
||||
).scalar_one()
|
||||
assert ban_rows == 0, "защита последнего узла не сработала — тест проверял бы не то"
|
||||
lease = proxy_pool.acquire(db, "cian", run_id=None)
|
||||
assert lease is not None and lease.id == proxy_id, f"пул пуст при живом узле: {lease}"
|
||||
|
||||
|
||||
def test_transport_failure_still_counts_against_node_health(node: tuple[int, Session]) -> None:
|
||||
"""Обратная сторона: сетевой сбой — не бан, узлу по-прежнему засчитывается отказ."""
|
||||
proxy_id, db = node
|
||||
provider = scraper_adapters.RealProxyProvider()
|
||||
|
||||
with pytest.raises(OSError):
|
||||
with curl_proxy_url(_Config(), provider, "cian", env_fallback_url=None):
|
||||
raise OSError("proxy 407")
|
||||
|
||||
assert _fails(db, proxy_id) == 1
|
||||
|
|
@ -255,7 +255,7 @@ async def test_captcha_on_curl_path_bans_the_node_for_cian() -> None:
|
|||
|
||||
assert spy.mark_banned_calls == [(1, "cian")]
|
||||
assert (1, True) not in spy.mark_health_calls, "узел с капчей записан здоровым"
|
||||
assert spy.mark_health_calls == [(1, False)]
|
||||
assert spy.mark_health_calls == [] # #3310: бан вместо mark_health(False)
|
||||
assert spy.release_calls == [1] # lease не течёт
|
||||
|
||||
|
||||
|
|
|
|||
148
tradein-mvp/backend/tests/test_3404_egress_run_attribution.py
Normal file
148
tradein-mvp/backend/tests/test_3404_egress_run_attribution.py
Normal 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
|
||||
|
|
@ -0,0 +1,90 @@
|
|||
"""Детальная Циана на своей сессии не держит event loop операциями пула (#3408 п.3).
|
||||
|
||||
`providers/cian/detail.py::fetch_detail` без `session`/`browser_fetcher` входил в
|
||||
синхронный `curl_proxy_url`: `RealProxyProvider.acquire/mark_health/release` ходят в БД
|
||||
прямо на loop'е. Вызывающий — админ-ручка истории цен Циана в публичном
|
||||
`tradein-backend` (один воркер uvicorn, #3083): пара блокирующих вызовов на каждый
|
||||
листинг батча. #3398 перевёл так три сайта `/estimate`, этот остался.
|
||||
|
||||
Проверка по значению тем же способом, что test_3398_pool_ops_off_event_loop.py: пока
|
||||
`acquire` спит 0.3 с в потоке, соседняя корутина тикает. На синхронном входе — ноль.
|
||||
Бан/здоровье/release на этом пути по-прежнему проверяет test_2700_cian_detail_403_node.py.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import time
|
||||
from contextlib import suppress
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from scraper_kit.contracts import ProxyLease
|
||||
from scraper_kit.providers.cian import detail as cian_detail
|
||||
|
||||
_ACQUIRE_SLEEP_S = 0.3
|
||||
_LEASE = ProxyLease(id=5, url="http://pool-node:3128", kind="http", rotate_url=None)
|
||||
|
||||
|
||||
class _SlowProvider:
|
||||
def __init__(self) -> None:
|
||||
self.released: list[int] = []
|
||||
self.health: list[bool] = []
|
||||
|
||||
def acquire(self, provider: str) -> ProxyLease:
|
||||
time.sleep(_ACQUIRE_SLEEP_S) # как checkout коннекта из пула
|
||||
return _LEASE
|
||||
|
||||
def release(self, lease: ProxyLease) -> None:
|
||||
self.released.append(lease.id)
|
||||
|
||||
def mark_health(self, lease: ProxyLease, ok: bool, **kw: Any) -> None:
|
||||
self.health.append(ok)
|
||||
|
||||
def mark_banned(self, lease: ProxyLease, *, source: str) -> None: # pragma: no cover
|
||||
pass
|
||||
|
||||
|
||||
@dataclass
|
||||
class _Config:
|
||||
use_proxy_pool_curl: bool = True
|
||||
cian_proxy_url: str | None = None
|
||||
environment: str = "production"
|
||||
|
||||
|
||||
async def test_acquire_does_not_block_event_loop() -> None:
|
||||
provider = _SlowProvider()
|
||||
session = MagicMock()
|
||||
session.get = AsyncMock(return_value=MagicMock(status_code=404, text=""))
|
||||
session.close = AsyncMock()
|
||||
ticks = 0
|
||||
|
||||
async def _ticker() -> None:
|
||||
nonlocal ticks
|
||||
while True:
|
||||
await asyncio.sleep(0)
|
||||
ticks += 1
|
||||
|
||||
task = asyncio.create_task(_ticker())
|
||||
await asyncio.sleep(0)
|
||||
try:
|
||||
with patch.object(cian_detail, "build_curl_cffi_session", return_value=session):
|
||||
got = await cian_detail.fetch_detail(
|
||||
"https://ekb.cian.ru/sale/flat/1/",
|
||||
config=_Config(), # type: ignore[arg-type]
|
||||
proxy_provider=provider, # type: ignore[arg-type]
|
||||
)
|
||||
finally:
|
||||
task.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
|
||||
assert got is None
|
||||
assert session.get.await_count == 1, "запрос не ушёл — тест ничего не проверил"
|
||||
assert provider.released == [_LEASE.id] and provider.health == [True]
|
||||
# Свободный loop за 0.3 с успевает десятки тысяч тиков, заблокированный — единицы.
|
||||
assert ticks > 100, f"loop простоял всё время acquire: тиков всего {ticks}"
|
||||
|
|
@ -26,7 +26,6 @@ from unittest.mock import MagicMock, patch
|
|||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from app.core.config import settings
|
||||
from app.services import estimator
|
||||
from app.services.geocoder import GeocodeResult
|
||||
|
||||
|
|
@ -45,7 +44,7 @@ _ANCHOR_PPM2 = (190_000, 195_000, 200_000, 205_000, 210_000)
|
|||
|
||||
|
||||
def _cap() -> float:
|
||||
return _CORRIDOR["high_ppm2"] * (1.0 + settings.estimate_corridor_clamp_slack)
|
||||
return _CORRIDOR["high_ppm2"] * (1.0 + estimator.CORRIDOR_CLAMP_SLACK)
|
||||
|
||||
|
||||
def _call(
|
||||
|
|
|
|||
|
|
@ -8,8 +8,8 @@
|
|||
- radius median выше dkp_low × factor → no-op (медиана не изменена)
|
||||
- dkp_raw is None → no-op (нет базы для floor)
|
||||
- anchor-путь (anchor_tier != None) → не затронут floor'ом
|
||||
- коридор ниже estimate_corridor_clamp_min_n → floor не применяется (#3466)
|
||||
- коридор ровно estimate_corridor_clamp_min_n → floor применяется (#3466)
|
||||
- коридор ниже CORRIDOR_CLAMP_MIN_N → floor не применяется (#3466)
|
||||
- коридор ровно CORRIDOR_CLAMP_MIN_N → floor применяется (#3466)
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -194,23 +194,23 @@ def test_no_dkp_raw_no_floor() -> None:
|
|||
|
||||
|
||||
def test_floor_not_applied_below_clamp_min_n() -> None:
|
||||
"""n = estimate_corridor_clamp_min_n − 1: тот же floor, что в тесте 1, но выключен.
|
||||
"""n = CORRIDOR_CLAMP_MIN_N − 1: тот же floor, что в тесте 1, но выключен.
|
||||
|
||||
`DkpCorridor.advisory_only` (#3452) обосновывает себя общим порогом у клампа
|
||||
headline И у radius-floor. Половина про кламп стережётся test_3452_*, эта —
|
||||
здесь: без гейта по count медиана 80k поднялась бы до 120k.
|
||||
"""
|
||||
from app.core.config import settings
|
||||
from app.core.config import CORRIDOR_CLAMP_MIN_N
|
||||
|
||||
analogs = _six(80_000.0)
|
||||
dkp_raw = {
|
||||
"count": settings.estimate_corridor_clamp_min_n - 1,
|
||||
"count": CORRIDOR_CLAMP_MIN_N - 1,
|
||||
"low_ppm2": 150_000,
|
||||
"median_ppm2": 180_000,
|
||||
"high_ppm2": 220_000,
|
||||
"period_months": 12,
|
||||
}
|
||||
est = _run_estimate(analogs, dkp_raw, radius_floor_factor=0.8)
|
||||
est = _run_estimate(analogs, dkp_raw)
|
||||
|
||||
assert est.median_price_per_m2 < 100_000, (
|
||||
f"median_ppm2={est.median_price_per_m2}: коридор из "
|
||||
|
|
@ -222,18 +222,18 @@ def test_floor_not_applied_below_clamp_min_n() -> None:
|
|||
|
||||
|
||||
def test_floor_applied_at_exactly_clamp_min_n() -> None:
|
||||
"""n = estimate_corridor_clamp_min_n: advisory_only=False → floor обязан поднять.
|
||||
"""n = CORRIDOR_CLAMP_MIN_N: advisory_only=False → floor обязан поднять.
|
||||
|
||||
Тест 4 стережёт «ниже порога — нет», тест 1 — «15 сделок — да». Граница между
|
||||
ними (`>=` против `>`) не стереглась: при n == min_n витрина не пишет подпись
|
||||
«справочно», значит коридор обязан войти в цену.
|
||||
"""
|
||||
from app.core.config import settings
|
||||
from app.core.config import CORRIDOR_CLAMP_MIN_N
|
||||
from app.schemas.trade_in import DkpCorridor
|
||||
|
||||
analogs = _six(80_000.0)
|
||||
dkp_raw = {
|
||||
"count": settings.estimate_corridor_clamp_min_n,
|
||||
"count": CORRIDOR_CLAMP_MIN_N,
|
||||
"low_ppm2": 150_000,
|
||||
"median_ppm2": 180_000,
|
||||
"high_ppm2": 220_000,
|
||||
|
|
@ -242,7 +242,7 @@ def test_floor_applied_at_exactly_clamp_min_n() -> None:
|
|||
# Та же граница с другой стороны: на ней коридор уже не справочный.
|
||||
assert DkpCorridor(**dkp_raw).advisory_only is False
|
||||
|
||||
est = _run_estimate(analogs, dkp_raw, radius_floor_factor=0.8)
|
||||
est = _run_estimate(analogs, dkp_raw)
|
||||
|
||||
assert est.median_price_per_m2 == 150_000 * 0.8, (
|
||||
f"median_ppm2={est.median_price_per_m2}: коридор из {dkp_raw['count']} сделок "
|
||||
|
|
|
|||
|
|
@ -127,10 +127,10 @@ def test_flag_on_exception_marks_fail_and_still_releases() -> None:
|
|||
assert spy.release_calls == [7]
|
||||
|
||||
|
||||
def test_ban_exception_calls_mark_banned_in_addition_to_mark_health() -> None:
|
||||
def test_ban_exception_calls_mark_banned_instead_of_mark_health() -> None:
|
||||
"""Исключение — подкласс ProxyBanError (напр. AvitoBlockedError) внутри блока —
|
||||
вызывает mark_banned(lease, source=provider) В ДОПОЛНЕНИЕ к mark_health(ok=False)
|
||||
(#2600 п.1 curl-путь). Zero изменений в вызывающем коде — сигнал детектируется
|
||||
вызывает mark_banned(lease, source=provider) ВМЕСТО mark_health(ok=False)
|
||||
(#2600 п.1 curl-путь, #3310). Zero изменений в вызывающем коде — сигнал детектируется
|
||||
по ТИПУ исключения, а не явным вызовом."""
|
||||
|
||||
class _FakeBlockedError(ProxyBanError):
|
||||
|
|
@ -143,7 +143,8 @@ def test_ban_exception_calls_mark_banned_in_addition_to_mark_health() -> None:
|
|||
assert url == _LEASE.url
|
||||
raise _FakeBlockedError("firewall page detected")
|
||||
assert spy.mark_banned_calls == [(7, "avito")]
|
||||
assert spy.mark_health_calls == [(7, False)] # оба сигнала, не взаимоисключающие
|
||||
# #3310: бан ВМЕСТО mark_health(False) — иначе три бана выводили узел из выдачи всем
|
||||
assert spy.mark_health_calls == []
|
||||
assert spy.release_calls == [7] # lease всё равно освобождён
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -587,9 +587,18 @@ def save_listings(
|
|||
price_rub = EXCLUDED.price_rub,
|
||||
price_per_m2 = EXCLUDED.price_per_m2,
|
||||
-- Cian-specific: обновляем при каждом re-scrape
|
||||
living_area_m2 = EXCLUDED.living_area_m2,
|
||||
-- #3252: жилая площадь и балконы — COALESCE по той же причине, что
|
||||
-- #3063 ниже: их добывает detail (avito, domklik), а SERP этих
|
||||
-- источников отдаёт NULL. Прод 17.09: у domklik, переобойдённых
|
||||
-- после обогащения, living_area_m2 = 0 из 5217 (без переобхода —
|
||||
-- 2990 из 4058), у avito — 1 из 8639 (4122 из 6505).
|
||||
living_area_m2 = COALESCE(
|
||||
EXCLUDED.living_area_m2, listings.living_area_m2
|
||||
),
|
||||
bedrooms_count = EXCLUDED.bedrooms_count,
|
||||
balconies_count = EXCLUDED.balconies_count,
|
||||
balconies_count = COALESCE(
|
||||
EXCLUDED.balconies_count, listings.balconies_count
|
||||
),
|
||||
loggias_count = EXCLUDED.loggias_count,
|
||||
description_minhash = EXCLUDED.description_minhash,
|
||||
cadastral_number = EXCLUDED.cadastral_number,
|
||||
|
|
@ -742,8 +751,11 @@ def save_listings(
|
|||
listings.price_previous_rub, listings.newbuilding_id,
|
||||
listings.newbuilding_url, listings.card_hash, listings.is_active
|
||||
) IS DISTINCT FROM (
|
||||
EXCLUDED.price_rub, EXCLUDED.price_per_m2, EXCLUDED.living_area_m2,
|
||||
EXCLUDED.bedrooms_count, EXCLUDED.balconies_count, EXCLUDED.loggias_count,
|
||||
EXCLUDED.price_rub, EXCLUDED.price_per_m2,
|
||||
COALESCE(EXCLUDED.living_area_m2, listings.living_area_m2),
|
||||
EXCLUDED.bedrooms_count,
|
||||
COALESCE(EXCLUDED.balconies_count, listings.balconies_count),
|
||||
EXCLUDED.loggias_count,
|
||||
EXCLUDED.description_minhash, EXCLUDED.cadastral_number,
|
||||
EXCLUDED.building_cadastral_number,
|
||||
-- #3063: те же COALESCE, что в SET выше. Правая часть гейта ОБЯЗАНА
|
||||
|
|
@ -823,9 +835,9 @@ def save_listings(
|
|||
is_active = true,
|
||||
price_rub = :price_rub,
|
||||
price_per_m2 = :ppm2,
|
||||
living_area_m2 = :living_area_m2,
|
||||
living_area_m2 = COALESCE(:living_area_m2, living_area_m2),
|
||||
bedrooms_count = :bedrooms_count,
|
||||
balconies_count = :balconies_count,
|
||||
balconies_count = COALESCE(:balconies_count, balconies_count),
|
||||
loggias_count = :loggias_count,
|
||||
description_minhash = :description_minhash,
|
||||
cadastral_number = :cadastral_number,
|
||||
|
|
@ -940,7 +952,7 @@ def save_listings(
|
|||
skip_snapshot = today_row is not None
|
||||
if skip_snapshot:
|
||||
logger.debug(
|
||||
"save_listings:snapshot_skipped (card unchanged) " "source=%s listing_id=%s",
|
||||
"save_listings:snapshot_skipped (card unchanged) source=%s listing_id=%s",
|
||||
lot.source,
|
||||
listing_id,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -18,11 +18,15 @@ release ВСЕГДА в finally — lease не должен течь, даже
|
|||
|
||||
Бан площадки (#2600 п.1): если исключение, поднятое ИЗНУТРИ `with curl_proxy_url(...) as
|
||||
url:`, — `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
|
||||
`DomClickBlockedError` и т.п. — см. `proxy_errors.ProxyBanError`), это НЕ просто
|
||||
`DomClickBlockedError` и т.п. — см. `proxy_errors.ProxyBanError`), это НЕ
|
||||
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
|
||||
`mark_banned` — узел сразу снимается с выдачи ЭТОМУ провайдеру (per-source бан, #2600 п.2;
|
||||
для остальных источников остаётся в строю), кроме случая когда это последний узел,
|
||||
достижимый для провайдера (защита в `app.services.proxy_pool.mark_banned`).
|
||||
`mark_health(ok=False)` на бане НЕ зовётся (#3310): узел исправен, его отбила площадка, а
|
||||
глобальный счётчик после трёх банов выводил его из выдачи ВСЕМ источникам — в том числе
|
||||
последний узел, который защита только что пообещала оставить. Тот же выбор, что у
|
||||
`BrowserFetcher.report_platform_ban` (#3288).
|
||||
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
|
||||
ИЗНУТРИ блока, получает сигнал бесплатно — этот модуль намеренно НЕ импортирует
|
||||
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные
|
||||
|
|
@ -126,13 +130,15 @@ def curl_proxy_url(
|
|||
raise
|
||||
finally:
|
||||
# mark_health/mark_banned/release — best-effort: проблема пула не должна
|
||||
# ронять сбор. mark_banned ПЕРЕД mark_health(ok=False) — оба независимы
|
||||
# (разные поля), но бан — более специфичный/сильный сигнал.
|
||||
# ронять сбор. Бан ВМЕСТО mark_health(ok=False), а не вдобавок (#3310): отказ
|
||||
# площадки — свойство пары «узел × источник», глобальный счётчик здоровья он
|
||||
# копить не должен (см. докстринг модуля).
|
||||
if banned:
|
||||
try:
|
||||
proxy_provider.mark_banned(lease, source=provider)
|
||||
except Exception:
|
||||
logger.warning("proxy_pool: mark_banned failed for %s", provider, exc_info=True)
|
||||
else:
|
||||
try:
|
||||
proxy_provider.mark_health(lease, ok)
|
||||
except Exception:
|
||||
|
|
@ -175,10 +181,10 @@ async def acurl_proxy_url(
|
|||
ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена
|
||||
приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им
|
||||
аренду — замер в тестах: `wait_for(timeout=0.05)` вернулся через ≈0.3 с (всё время
|
||||
acquire). Худший случай `/estimate` — три источника × до 30 с checkout'а коннекта из
|
||||
исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток не прервать.
|
||||
Закрывается со стороны БД — `statement_timeout`/`pool_timeout` короче бюджета,
|
||||
follow-up #3408 («пул коннектов SQLAlchemy 15 против пикового спроса до 28»).
|
||||
acquire). Худший случай `/estimate` — три источника × до `pool_timeout` checkout'а
|
||||
коннекта из исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток
|
||||
не прервать. Закрыто со стороны БД (#3408, PR #3444): `pool_timeout` = 5 с — короче
|
||||
самого короткого бюджета источника (8 с), см. app/core/db.py.
|
||||
"""
|
||||
cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url)
|
||||
enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__))
|
||||
|
|
|
|||
|
|
@ -32,7 +32,7 @@ from scraper_kit.offer_price_history import (
|
|||
validate_diff_percent,
|
||||
)
|
||||
from scraper_kit.providers._base import build_curl_cffi_session
|
||||
from scraper_kit.providers._proxy import curl_proxy_url
|
||||
from scraper_kit.providers._proxy import acurl_proxy_url
|
||||
from scraper_kit.proxy_errors import caused_by_no_proxy
|
||||
from scraper_kit.repair_state_normalizer import (
|
||||
infer_repair_state_from_text,
|
||||
|
|
@ -142,7 +142,7 @@ def _is_captcha_title(title: str) -> bool:
|
|||
"""Заголовок — капча Циана? Один предикат на оба места детекта.
|
||||
|
||||
Детект живёт в двух точках намеренно (#3402 follow-up): общий parse-путь ниже —
|
||||
ЗА границей `with curl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула
|
||||
ЗА границей `async with acurl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула
|
||||
уже не доходит (узел успел получить `mark_health(ok=True)` и остаётся в выдаче
|
||||
Циану — ровно дефект #2700, только на HTTP 200). Поэтому curl-путь спрашивает про
|
||||
капчу ВНУТРИ блока, а не после него.
|
||||
|
|
@ -228,9 +228,13 @@ async def fetch_detail(
|
|||
# curl_cffi own-session path (legacy, back-compat).
|
||||
# proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian.
|
||||
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
|
||||
# Пусто → прямое подключение (dev/no-op). curl_proxy_url: mark_health + release на выходе.
|
||||
# Пусто → прямое подключение (dev/no-op). На выходе mark_health/mark_banned + release.
|
||||
# acurl_proxy_url (#3408 п.3): операции пула в потоке — синхронный вход держал event
|
||||
# loop публичного tradein-backend (админ-ручка истории цен Циана, пара на листинг).
|
||||
_env = config.cian_proxy_url if config is not None else None
|
||||
with curl_proxy_url(config, proxy_provider, "cian", env_fallback_url=_env) as _proxy_url:
|
||||
async with acurl_proxy_url(
|
||||
config, proxy_provider, "cian", env_fallback_url=_env
|
||||
) as _proxy_url:
|
||||
own_session = build_curl_cffi_session(
|
||||
proxy_url=_proxy_url,
|
||||
timeout=25.0,
|
||||
|
|
@ -242,7 +246,7 @@ async def fetch_detail(
|
|||
try:
|
||||
resp = await own_session.get(offer_url, allow_redirects=True)
|
||||
if resp.status_code != 200:
|
||||
# ВНУТРИ curl_proxy_url: поднятый отсюда ProxyBanError доходит до
|
||||
# ВНУТРИ acurl_proxy_url: поднятый отсюда ProxyBanError доходит до
|
||||
# пула (mark_banned на пару «узел × cian», #2600 п.2). Раньше здесь
|
||||
# был `return None` — узел получал mark_health(ok=True) и оставался
|
||||
# в выдаче Циану (#2700, 15 суток по 50 отказов в сутки).
|
||||
|
|
@ -252,7 +256,7 @@ async def fetch_detail(
|
|||
html = resp.text
|
||||
# Капча приходит с HTTP 200, то есть мимо `_raise_if_blocked`. Спрашиваем
|
||||
# ЗДЕСЬ, пока lease жив: общий parse-путь ниже поднимет то же исключение
|
||||
# уже после выхода из `with curl_proxy_url`, где узлу проставлен
|
||||
# уже после выхода из `async with acurl_proxy_url`, где узлу проставлен
|
||||
# mark_health(ok=True) и бана пары «узел × cian» не будет (#3402).
|
||||
title = _page_title(html)
|
||||
if _is_captcha_title(title):
|
||||
|
|
|
|||
|
|
@ -27,7 +27,8 @@ _extract_ssr_state делает balanced-brace scan (с пропуском ск
|
|||
Форма SSR-стейта подтверждена на живой карточке 2075729321 (2026-06-27):
|
||||
productCard.objectInfo.{renovation,livingArea,kitchenArea}, productCard.priceInfo.
|
||||
priceHistory ({date ISO8601+tz, price, diff, state}), productCard.egrnData
|
||||
(snake_case owners_count/collateral/collateral_sber), productCard.legalOptions.
|
||||
(snake_case; area/floor/owners_count — обёртка {status, value}, collateral/
|
||||
collateral_sber — bool; перепроверено 17.09, #3252), productCard.legalOptions.
|
||||
saleType, productCard.viewsCount/callsCount, houseInfo.info.{wallType,floorType},
|
||||
и ТОП-УРОВНЕМ pricePrediction (DomClick AVM → raw_extra.avm).
|
||||
|
||||
|
|
@ -283,6 +284,29 @@ def _to_int(value: Any) -> int | None:
|
|||
return None
|
||||
|
||||
|
||||
def _pos_float(value: Any) -> float | None:
|
||||
"""float > 0, иначе None (#3252).
|
||||
|
||||
Домклик отдаёт ``livingArea: 0`` / ``kitchenArea: 0`` как «не указано» (живая
|
||||
карточка 2078257603, 17.09). Без этого в колонку ложился 0.00, а COALESCE в
|
||||
UPDATE затирал им ранее известную площадь: на проде 303 кухни и 84 жилые
|
||||
площади = 0 — только у domklik и только у обогащённых после 26.08.
|
||||
"""
|
||||
f = _to_float(value)
|
||||
return f if f is not None and f > 0 else None
|
||||
|
||||
|
||||
def _egrn_value(value: Any) -> Any:
|
||||
"""Узел egrnData → значение (#3252).
|
||||
|
||||
Поля ЕГРН приходят обёрткой ``{"status": "success"|"warning"|"error", "value": X}``
|
||||
(живые карточки 17.09: area, floor, owners_count). ``status`` — сверка с тем, что
|
||||
указал продавец, а ``value`` — сам факт из ЕГРН, поэтому берём value при любом
|
||||
статусе. Скаляр пропускаем как есть.
|
||||
"""
|
||||
return value.get("value") if isinstance(value, dict) else value
|
||||
|
||||
|
||||
def _parse_change_time(value: Any) -> datetime | None:
|
||||
"""Нормализует дату изменения цены → timezone-aware datetime.
|
||||
|
||||
|
|
@ -436,8 +460,8 @@ def parse_detail_html(html: str, source_url: str) -> DomClickDetailEnrichment:
|
|||
repair_state = _map_repair_state(repair_type)
|
||||
|
||||
# площади
|
||||
living_area_m2 = _to_float(oi.get("livingArea"))
|
||||
kitchen_area_m2 = _to_float(oi.get("kitchenArea"))
|
||||
living_area_m2 = _pos_float(oi.get("livingArea"))
|
||||
kitchen_area_m2 = _pos_float(oi.get("kitchenArea"))
|
||||
|
||||
# балконы: int → count+флаг; иное (строка/др.) → raw_extra, counts None
|
||||
balconies = oi.get("balconies")
|
||||
|
|
@ -454,10 +478,11 @@ def parse_detail_html(html: str, source_url: str) -> DomClickDetailEnrichment:
|
|||
|
||||
sale_type = (pc.get("legalOptions") or {}).get("saleType")
|
||||
|
||||
# owners_count — защищаемся от двух написаний ключа
|
||||
owners_count = _to_int(egrn.get("owners_count"))
|
||||
# owners_count — обёртка {status, value} (см. _egrn_value); до #3252 int()
|
||||
# от словаря молча давал None: 0 из ~2970 карточек, обогащённых с 29.08.
|
||||
owners_count = _to_int(_egrn_value(egrn.get("owners_count")))
|
||||
if owners_count is None:
|
||||
owners_count = _to_int(egrn.get("ownersCount"))
|
||||
owners_count = _to_int(_egrn_value(egrn.get("ownersCount")))
|
||||
|
||||
# encumbrances_clean (ИНВЕРСИЯ): collateral/collateral_sber обозначают
|
||||
# НАЛИЧИЕ обременения. Любой truthy → clean=False; оба отсутствуют/false
|
||||
|
|
@ -523,8 +548,10 @@ def parse_detail_html(html: str, source_url: str) -> DomClickDetailEnrichment:
|
|||
|
||||
# raw_extra → merge в listings.raw_payload. wall/floor type ОСТАЮТСЯ тут,
|
||||
# НЕ пишем в listings.house_type (кросс-источниковый словарь грязный).
|
||||
# egrn area key — точное написание не подтверждено; пробуем несколько.
|
||||
egrn_area = egrn.get("area") or egrn.get("rosreestrArea") or egrn.get("object_area")
|
||||
# egrnData.area — ключ подтверждён на живых карточках 17.09 (#3252). В raw_payload
|
||||
# кладём обёртку {status, value} целиком: status говорит, сошлась ли площадь
|
||||
# ЕГРН с площадью объявления, и уже лежит в таком виде у 2156 строк.
|
||||
egrn_area = egrn.get("area")
|
||||
demand = _compact(
|
||||
{
|
||||
"calls": pc.get("callsCount"),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue