diff --git a/tradein-mvp/backend/app/services/proxy_egress.py b/tradein-mvp/backend/app/services/proxy_egress.py index 12880ea7..a7503ad0 100644 --- a/tradein-mvp/backend/app/services/proxy_egress.py +++ b/tradein-mvp/backend/app/services/proxy_egress.py @@ -16,7 +16,8 @@ reap_stale_leases) — она рассчитана на долгоживущие гарантированного `release` на каждом пути выхода; занимать под них lease значило бы дырявить пул фантомно занятыми узлами при малейшей утечке release. Резолвер ниже — ЧИСТО READ, той же таблицы `scrape_proxies` + `scrape_proxy_source_bans`, без блокировок -и без мутаций. +и без мутаций пула. Единственная запись — атрибуция прогона (`scrape_runs.proxy_id`, #3404), +своей короткой сессией, см. `resolve_proxy_url`. ПРАВИЛО ВЫБОРА: enabled=true, consecutive_fails < proxy_pool.MAX_CONSECUTIVE_FAILS (тот же карантинный порог, что у acquire), нет АКТИВНОЙ строки (banned_until > now()) @@ -69,7 +70,7 @@ from sqlalchemy.orm import Session from app.core.config import settings as _settings from app.core.db import SessionLocal as _SessionLocal -from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS +from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS, attribute_run_proxy logger = logging.getLogger(__name__) @@ -277,6 +278,21 @@ def resolve_proxy_url(db: Session, source: str) -> str | None: _safe_label(candidate.id, candidate.label, candidate.url), candidate.ban_count, ) + # #3404: прогоны на этом пути (yandex_detail_backfill, yandex_address_backfill, + # curl-ветка avito_detail_backfill) lease не берут, а атрибуцию раньше писал только + # acquire() — у yandex_detail_backfill proxy_id был NULL в 56 прогонах из 56. + # Своя сессия, не db вызывающего: attribute_run_proxy коммитит, а на сбое + # откатывает, db же — долгоживущая сессия прогона посреди работы (см. + # _pick_candidate). Сбой атрибуции проглатывается внутри — выдачу не роняет. + from scraper_kit.orchestration.run_context import current_run_id + + run_id = current_run_id.get() + if run_id is not None: + attr_db = _SessionLocal() + try: + attribute_run_proxy(attr_db, run_id, candidate.id) + finally: + attr_db.close() return candidate.url diag = _diagnose_no_candidate(db, source) diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index 437e937c..66c3153d 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() ) ) @@ -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() ) ) diff --git a/tradein-mvp/backend/data/sql/320_listings_domclick_zero_area_to_null.sql b/tradein-mvp/backend/data/sql/320_listings_domclick_zero_area_to_null.sql new file mode 100644 index 00000000..e2147b23 --- /dev/null +++ b/tradein-mvp/backend/data/sql/320_listings_domclick_zero_area_to_null.sql @@ -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; diff --git a/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py b/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py index 68f606d4..8bcccd4e 100644 --- a/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py +++ b/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py @@ -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 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/skip_allowlist.txt b/tradein-mvp/backend/tests/skip_allowlist.txt index d46a4876..01f3b26b 100644 --- a/tradein-mvp/backend/tests/skip_allowlist.txt +++ b/tradein-mvp/backend/tests/skip_allowlist.txt @@ -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. diff --git a/tradein-mvp/backend/tests/test_2700_cian_detail_403_node.py b/tradein-mvp/backend/tests/test_2700_cian_detail_403_node.py index d6637486..c4da8e96 100644 --- a/tradein-mvp/backend/tests/test_2700_cian_detail_403_node.py +++ b/tradein-mvp/backend/tests/test_2700_cian_detail_403_node.py @@ -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 не течёт даже на бане diff --git a/tradein-mvp/backend/tests/test_2830_pool_bypass_tails.py b/tradein-mvp/backend/tests/test_2830_pool_bypass_tails.py index 3c559fc8..27aa1d92 100644 --- a/tradein-mvp/backend/tests/test_2830_pool_bypass_tails.py +++ b/tradein-mvp/backend/tests/test_2830_pool_bypass_tails.py @@ -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] diff --git a/tradein-mvp/backend/tests/test_3252_domclick_card_fields.py b/tradein-mvp/backend/tests/test_3252_domclick_card_fields.py new file mode 100644 index 00000000..d5af85af --- /dev/null +++ b/tradein-mvp/backend/tests/test_3252_domclick_card_fields.py @@ -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"" + 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) 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..70b378ec --- /dev/null +++ b/tradein-mvp/backend/tests/test_3299_fallback_counts_any_nodes.py @@ -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 diff --git a/tradein-mvp/backend/tests/test_3310_curl_ban_keeps_last_node.py b/tradein-mvp/backend/tests/test_3310_curl_ban_keeps_last_node.py new file mode 100644 index 00000000..40a75a4d --- /dev/null +++ b/tradein-mvp/backend/tests/test_3310_curl_ban_keeps_last_node.py @@ -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 diff --git a/tradein-mvp/backend/tests/test_3402_cian_captcha_http200.py b/tradein-mvp/backend/tests/test_3402_cian_captcha_http200.py index 78cd03b9..f7b77441 100644 --- a/tradein-mvp/backend/tests/test_3402_cian_captcha_http200.py +++ b/tradein-mvp/backend/tests/test_3402_cian_captcha_http200.py @@ -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 не течёт diff --git a/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py b/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py new file mode 100644 index 00000000..552bdfd6 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3404_egress_run_attribution.py @@ -0,0 +1,148 @@ +"""Прогон на egress-пути знает свой узел (#3404, хвост PR #3405). + +PR #3405 писал `scrape_runs.proxy_id` только из `proxy_pool.acquire()`. Прогоны, которые +берут прокси через `proxy_egress.resolve_proxy_url` (без аренды: yandex_detail_backfill, +yandex_address_backfill, curl-ветка avito_detail_backfill), атрибуцию не получали — +прод 17.09: у yandex_detail_backfill proxy_id NULL в 56 прогонах из 56. + +Живой Postgres, настоящие сессии (атрибуция идёт своей сессией, поэтому строки +коммитятся и удаляются в finally). Без БД — skip. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +from collections.abc import Iterator +from typing import Any + +import pytest +from scraper_kit.orchestration.run_context import current_run_id +from sqlalchemy import create_engine, text +from sqlalchemy.orm import Session, sessionmaker + +from app.services import proxy_egress + + +def _live_engine() -> Any | None: + dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "") + if not dsn or "localhost:5432/test" in dsn: + return None + try: + engine = create_engine(dsn, future=True) + with engine.connect() as conn: + conn.execute(text("SELECT proxy_id FROM scrape_runs LIMIT 1")) + return engine + except Exception: + return None + + +_ENGINE = _live_engine() +pytestmark = pytest.mark.skipif(_ENGINE is None, reason="no reachable Postgres test DB") + + +@pytest.fixture +def env(monkeypatch: pytest.MonkeyPatch) -> Iterator[dict[str, Any]]: + """Узел пула + два прогона: `run` — текущий, `other` — строка, которую вызывающий + держит изменённой без коммита.""" + assert _ENGINE is not None + factory = sessionmaker(bind=_ENGINE, future=True) + monkeypatch.setattr(proxy_egress, "_SessionLocal", factory) + setup = factory() + proxy_id = setup.execute( + text( + "INSERT INTO scrape_proxies (url, provider_affinity) " + "VALUES ('http://t3404-' || gen_random_uuid(), 'any') RETURNING id" + ) + ).scalar_one() + run_ids = [ + setup.execute( + text( + "INSERT INTO scrape_runs (source, status) " + "VALUES ('test_3404', 'running') RETURNING id" + ) + ).scalar_one() + for _ in range(2) + ] + setup.commit() + token = current_run_id.set(None) + try: + yield {"proxy_id": int(proxy_id), "run": int(run_ids[0]), "other": int(run_ids[1])} + finally: + current_run_id.reset(token) + setup.rollback() + setup.execute(text("DELETE FROM scrape_runs WHERE id = ANY(:ids)"), {"ids": run_ids}) + setup.execute(text("DELETE FROM scrape_proxies WHERE id = :id"), {"id": proxy_id}) + setup.commit() + setup.close() + + +def _run_row(run_id: int) -> Any: + assert _ENGINE is not None + with _ENGINE.connect() as conn: + return conn.execute( + text("SELECT proxy_id, counters FROM scrape_runs WHERE id = :id"), {"id": run_id} + ).one() + + +def _node_id_of(url: str) -> int: + assert _ENGINE is not None + with _ENGINE.connect() as conn: + return int( + conn.execute( + text("SELECT id FROM scrape_proxies WHERE url = :u"), {"u": url} + ).scalar_one() + ) + + +def test_resolve_within_run_writes_node_to_scrape_runs(env: dict[str, Any]) -> None: + """Главный случай: прогон идёт, egress выбран — прогон знает узел. + + На старом коде proxy_id остаётся NULL.""" + current_run_id.set(env["run"]) + caller = Session(bind=_ENGINE, future=True) + try: + url = proxy_egress.resolve_proxy_url(caller, "yandex") + finally: + caller.close() + + assert url is not None + node = _node_id_of(url) + row = _run_row(env["run"]) + assert row.proxy_id == node + counters = row.counters if isinstance(row.counters, dict) else json.loads(row.counters) + assert counters.get("proxy_ids") == [node] + + +def test_resolve_outside_run_writes_nothing(env: dict[str, Any]) -> None: + """Без прогона (админка, проверка кук) — scrape_runs не трогаем.""" + caller = Session(bind=_ENGINE, future=True) + try: + assert proxy_egress.resolve_proxy_url(caller, "yandex") is not None + finally: + caller.close() + + assert _run_row(env["run"]).proxy_id is None + + +def test_attribution_does_not_commit_callers_transaction(env: dict[str, Any]) -> None: + """db вызывающего — долгоживущая сессия прогона посреди работы. Атрибуция не имеет + права закоммитить его незавершённые изменения (или откатить их на сбое).""" + current_run_id.set(env["run"]) + caller = Session(bind=_ENGINE, future=True) + try: + caller.execute( + text("UPDATE scrape_runs SET counters = '{\"t3404\": 1}'::jsonb WHERE id = :id"), + {"id": env["other"]}, + ) + proxy_egress.resolve_proxy_url(caller, "yandex") + caller.rollback() + finally: + caller.close() + + other = _run_row(env["other"]).counters or {} + assert "t3404" not in other, "атрибуция закоммитила чужую транзакцию" + assert _run_row(env["run"]).proxy_id is not None diff --git a/tradein-mvp/backend/tests/test_3408_cian_detail_pool_ops_off_loop.py b/tradein-mvp/backend/tests/test_3408_cian_detail_pool_ops_off_loop.py new file mode 100644 index 00000000..e1107e0f --- /dev/null +++ b/tradein-mvp/backend/tests/test_3408_cian_detail_pool_ops_off_loop.py @@ -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}" diff --git a/tradein-mvp/backend/tests/test_3466_corridor_tier_a.py b/tradein-mvp/backend/tests/test_3466_corridor_tier_a.py index 4350f621..72542446 100644 --- a/tradein-mvp/backend/tests/test_3466_corridor_tier_a.py +++ b/tradein-mvp/backend/tests/test_3466_corridor_tier_a.py @@ -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( diff --git a/tradein-mvp/backend/tests/test_estimator_radius_floor.py b/tradein-mvp/backend/tests/test_estimator_radius_floor.py index c532ed95..0f3e8485 100644 --- a/tradein-mvp/backend/tests/test_estimator_radius_floor.py +++ b/tradein-mvp/backend/tests/test_estimator_radius_floor.py @@ -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']} сделок " diff --git a/tradein-mvp/backend/tests/test_proxy_pool_curl_paths.py b/tradein-mvp/backend/tests/test_proxy_pool_curl_paths.py index fe94b232..5e06ee9a 100644 --- a/tradein-mvp/backend/tests/test_proxy_pool_curl_paths.py +++ b/tradein-mvp/backend/tests/test_proxy_pool_curl_paths.py @@ -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 всё равно освобождён diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/base.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/base.py index d8676724..f470f51e 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/base.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/base.py @@ -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, ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py index a6f3e539..ce861335 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py @@ -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,17 +130,19 @@ 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) - try: - proxy_provider.mark_health(lease, ok) - except Exception: - logger.warning("proxy_pool: mark_health failed for %s", provider, exc_info=True) + else: + try: + proxy_provider.mark_health(lease, ok) + except Exception: + logger.warning("proxy_pool: mark_health failed for %s", provider, exc_info=True) try: proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт 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__)) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py index bd34a25d..dd4ae8f2 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py @@ -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): diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py index 66912922..2a845426 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py @@ -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"),