Merge remote-tracking branch 'origin/main' into fix/browser-instances-limit
All checks were successful
CI Trade-In / changes (pull_request) Successful in 29s
CI / changes (pull_request) Successful in 33s
CI Trade-In / backend-tests (pull_request) Successful in 7m23s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Successful in 2m11s
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

This commit is contained in:
bot-backend 2026-09-17 15:50:18 +05:00
commit 92af3491ae
19 changed files with 1034 additions and 76 deletions

View file

@ -16,7 +16,8 @@ reap_stale_leases) — она рассчитана на долгоживущие
гарантированного `release` на каждом пути выхода; занимать под них lease значило бы гарантированного `release` на каждом пути выхода; занимать под них lease значило бы
дырявить пул фантомно занятыми узлами при малейшей утечке release. Резолвер ниже дырявить пул фантомно занятыми узлами при малейшей утечке release. Резолвер ниже
ЧИСТО READ, той же таблицы `scrape_proxies` + `scrape_proxy_source_bans`, без блокировок ЧИСТО READ, той же таблицы `scrape_proxies` + `scrape_proxy_source_bans`, без блокировок
и без мутаций. и без мутаций пула. Единственная запись атрибуция прогона (`scrape_runs.proxy_id`, #3404),
своей короткой сессией, см. `resolve_proxy_url`.
ПРАВИЛО ВЫБОРА: enabled=true, consecutive_fails < proxy_pool.MAX_CONSECUTIVE_FAILS ПРАВИЛО ВЫБОРА: enabled=true, consecutive_fails < proxy_pool.MAX_CONSECUTIVE_FAILS
(тот же карантинный порог, что у acquire), нет АКТИВНОЙ строки (banned_until > now()) (тот же карантинный порог, что у acquire), нет АКТИВНОЙ строки (banned_until > now())
@ -69,7 +70,7 @@ from sqlalchemy.orm import Session
from app.core.config import settings as _settings from app.core.config import settings as _settings
from app.core.db import SessionLocal as _SessionLocal 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__) 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), _safe_label(candidate.id, candidate.label, candidate.url),
candidate.ban_count, 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 return candidate.url
diag = _diagnose_no_candidate(db, source) diag = _diagnose_no_candidate(db, source)

View file

@ -275,12 +275,13 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
чужая только запасной вариант, чтобы источник не голодал при живых свободных узлах чужая только запасной вариант, чтобы источник не голодал при живых свободных узлах
чужой affinity (#2600). чужой affinity (#2600).
Fallback НЕ трогает последний enabled-узел выделенной (не-'any') affinity: если Fallback НЕ трогает последний пригодный узел выделенной (не-'any') affinity: если
fallback заберёт его под чужой источник, «свой» останется без прокси вообще хуже, fallback заберёт его под чужой источник, «свой» останется без прокси вообще хуже,
чем голодание исходного источника, которое фикс призван устранить. Кандидат чем голодание исходного источника, которое фикс призван устранить. Кандидат
участвует в fallback, только если его affinity='any' ИЛИ у этой affinity есть ДРУГОЙ участвует в fallback, только если его affinity='any' ИЛИ у источника этой affinity
enabled-узел (EXISTS-подзапрос) т.е. выдача не обнулит доступность выделенной без него останется ДРУГОЙ кандидат (EXISTS-подзапрос): узел той же affinity или
affinity целиком. 'any', здоровый, с живой арендой порта и не забаненный этим источником (#3299 —
до него 'any'-узлы не считались, и единственный выделенный узел не выдавался никому).
Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql
единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR
@ -401,17 +402,29 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
-- fallback увести последний реально рабочий узел выделенной -- fallback увести последний реально рабочий узел выделенной
-- affinity и обрушить её (два domclick-узла, один забанен -- affinity и обрушить её (два domclick-узла, один забанен
-- domclick'ом → второй уходит под avito → domclick без прокси). -- domclick'ом → второй уходит под avito → domclick без прокси).
--
-- #3299: вопрос «останется ли у источника sp.provider_affinity
-- хоть один кандидат без sp», а не «есть ли ВТОРОЙ узел той же
-- привязки». 'any'-узлы основной запрос выдаёт выделенному
-- источнику наравне, значит и backup'ом они считаются; прежний
-- `= sp.provider_affinity` прятал единственный выделенный узел от
-- всех. Здоровье и срок аренды как в основном запросе.
-- leased_by НЕ проверяем: аренда вернётся через минуты, а
-- отказ из-за чужой аренды прятал бы узел при любом параллельном
-- прогоне (тот же довод, что у защиты в mark_banned).
OR EXISTS ( OR EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxies AS other FROM scrape_proxies AS other
WHERE other.provider_affinity = sp.provider_affinity WHERE other.provider_affinity IN (sp.provider_affinity, 'any')
AND other.enabled AND other.enabled
AND other.consecutive_fails < CAST(:max_fails AS integer)
AND (other.expires_at IS NULL OR other.expires_at > now())
AND other.id <> sp.id AND other.id <> sp.id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxy_source_bans b2 FROM scrape_proxy_source_bans b2
WHERE b2.proxy_id = other.id WHERE b2.proxy_id = other.id
AND b2.source = other.provider_affinity AND b2.source = sp.provider_affinity
AND b2.banned_until > now() AND b2.banned_until > now()
) )
) )
@ -446,9 +459,9 @@ def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLe
) )
db.commit() db.commit()
if run_id is not None and run_id != NON_RUN_LEASE_MARKER: if run_id is not None and run_id != NON_RUN_LEASE_MARKER:
# #3404: одна точка, покрывающая ВСЕ пути выдачи (curl — acquire на каждый # #3404: покрывает все пути АРЕНДЫ (curl — acquire на каждый вызов, браузер —
# вызов, браузер — sticky lease на весь прогон, ре-acquire при ротации узла # sticky lease на весь прогон, ре-acquire при ротации узла mid-run). Путь без
# mid-run) — см. attribute_run_proxy docstring. # аренды (proxy_egress.resolve_proxy_url) пишет атрибуцию сам — см. docstring.
attribute_run_proxy(db, run_id, proxy_id) attribute_run_proxy(db, run_id, proxy_id)
if fallback_used: if fallback_used:
logger.warning( 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: def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None:
"""Записать узел, через который идёт прогон run_id, в scrape_runs (#3404). """Записать узел, через который идёт прогон run_id, в scrape_runs (#3404).
Единственный писатель `acquire()` сразу после выдачи lease'а: покрывает и Два писателя. `acquire()` сразу после выдачи lease'а: покрывает и curl-путь
curl-путь (acquire на каждый вызов), и браузерный sticky lease (один acquire на (acquire на каждый вызов), и браузерный sticky lease (один acquire на весь прогон),
весь прогон), и ре-acquire при ротации узла mid-run (`browser_fetcher`, и ре-acquire при ротации узла mid-run (`browser_fetcher`,
`_LEASE_ROTATE_AFTER_FAILS`) то есть смена узла ЗА прогон фиксируется сама, `_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` ПОСЛЕДНИЙ использованный узел (перезаписывается при `scrape_runs.proxy_id` ПОСЛЕДНИЙ использованный узел (перезаписывается при
каждой новой выдаче); полная цепочка узлов, если она менялась, в каждой новой выдаче); полная цепочка узлов, если она менялась, в
@ -535,8 +549,7 @@ def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None:
db.commit() db.commit()
if row is None: if row is None:
logger.debug( logger.debug(
"proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already " "proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already finalized?)",
"finalized?)",
run_id, run_id,
) )
except Exception: except Exception:
@ -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) WHERE sp.id <> CAST(:proxy_id AS bigint)
AND sp.enabled AND sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer) 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 ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxy_source_bans b 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'ом -- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом
-- не считается (иначе защита сочла бы affinity живой, когда она -- не считается (иначе защита сочла бы affinity живой, когда она
-- уже нет). -- уже нет).
-- #3299: предикат backup'а — буква в букву как в fallback
-- acquire() (там же обоснование): 'any'-узлы в счёт, здоровье
-- и срок аренды проверяются.
OR EXISTS ( OR EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxies other FROM scrape_proxies other
WHERE other.provider_affinity = sp.provider_affinity WHERE other.provider_affinity IN (sp.provider_affinity, 'any')
AND other.enabled AND other.enabled
AND other.consecutive_fails < CAST(:max_fails AS integer)
AND (other.expires_at IS NULL OR other.expires_at > now())
AND other.id <> sp.id AND other.id <> sp.id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM scrape_proxy_source_bans b2 FROM scrape_proxy_source_bans b2
WHERE b2.proxy_id = other.id WHERE b2.proxy_id = other.id
AND b2.source = other.provider_affinity AND b2.source = sp.provider_affinity
AND b2.banned_until > now() AND b2.banned_until > now()
) )
) )

View file

@ -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.0814.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;

View file

@ -83,9 +83,9 @@ _SSR_LITERAL = """{
}, },
"legalOptions": {"saleType": "Свободная продажа"}, "legalOptions": {"saleType": "Свободная продажа"},
"egrnData": { "egrnData": {
"area": 38.2, "area": {"status": "success", "value": 38.2},
"floor": 5, "floor": {"status": "success", "value": 5},
"owners_count": 1, "owners_count": {"status": "success", "value": 1},
"collateral": true, "collateral": true,
"collateral_sber": false "collateral_sber": false
}, },
@ -264,7 +264,7 @@ def test_parse_detail_html_raw_extra() -> None:
# wall/floor type живут ТОЛЬКО в raw_extra, НЕ в listings.house_type. # wall/floor type живут ТОЛЬКО в raw_extra, НЕ в listings.house_type.
assert "house_type" not in raw assert "house_type" not in raw
assert raw["domclick_building_guid"] == "abc-guid-123" 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"]["calls"] == 5
assert raw["demand"]["favorites"] == 12 assert raw["demand"]["favorites"] == 12
# AVM (Layer C, top-level pricePrediction) → raw_extra.avm # AVM (Layer C, top-level pricePrediction) → raw_extra.avm

View file

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

View file

@ -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_statement_timeout_overrides_session_ceiling
tests/test_3463_db_timeouts.py::test_set_local_is_scoped_to_its_transaction 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-остановкой исполняет только # Миграция 310 (#3385): настоящий SQL-файл с DELETE и RAISE-остановкой исполняет только
# Postgres. В ci-tradein.yml бегут по-настоящему (postgres-сервис, #2745); краснеют от # Postgres. В ci-tradein.yml бегут по-настоящему (postgres-сервис, #2745); краснеют от
# снятия подписи батча, сужения окна до ×10 и снятия порога — проверено вручную 17.09. # снятия подписи батча, сужения окна до ×10 и снятия порога — проверено вручную 17.09.

View file

@ -103,7 +103,7 @@ async def test_403_bans_the_node_for_cian_only() -> None:
with pytest.raises(CianBlockedError): with pytest.raises(CianBlockedError):
await _fetch(403, spy) await _fetch(403, spy)
assert spy.mark_banned_calls == [(1, "cian")] 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 не течёт даже на бане assert spy.release_calls == [1] # lease не течёт даже на бане

View file

@ -179,7 +179,7 @@ async def test_price_history_403_bans_the_node_for_cian() -> None:
pool = _SpyPool() pool = _SpyPool()
result = await _run_price_history(pool, status_code=403) result = await _run_price_history(pool, status_code=403)
assert pool.mark_banned_calls == [(9, "cian")] 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 # прогон честен: отказ посчитан assert result.errors == 1 # прогон честен: отказ посчитан
@ -335,7 +335,7 @@ async def test_zhk_resolve_403_reaches_the_pool() -> None:
with pytest.raises(CianBlockedError): with pytest.raises(CianBlockedError):
await _resolve(403, spy) await _resolve(403, spy)
assert spy.mark_banned_calls == [(9, "cian")] 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] assert spy.release_calls == [9]

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

View 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

View 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

View file

@ -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 spy.mark_banned_calls == [(1, "cian")]
assert (1, True) not in spy.mark_health_calls, "узел с капчей записан здоровым" 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 не течёт assert spy.release_calls == [1] # lease не течёт

View file

@ -0,0 +1,148 @@
"""Прогон на egress-пути знает свой узел (#3404, хвост PR #3405).
PR #3405 писал `scrape_runs.proxy_id` только из `proxy_pool.acquire()`. Прогоны, которые
берут прокси через `proxy_egress.resolve_proxy_url` (без аренды: yandex_detail_backfill,
yandex_address_backfill, curl-ветка avito_detail_backfill), атрибуцию не получали
прод 17.09: у yandex_detail_backfill proxy_id NULL в 56 прогонах из 56.
Живой Postgres, настоящие сессии (атрибуция идёт своей сессией, поэтому строки
коммитятся и удаляются в finally). Без БД skip.
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
import json
from collections.abc import Iterator
from typing import Any
import pytest
from scraper_kit.orchestration.run_context import current_run_id
from sqlalchemy import create_engine, text
from sqlalchemy.orm import Session, sessionmaker
from app.services import proxy_egress
def _live_engine() -> Any | None:
dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "")
if not dsn or "localhost:5432/test" in dsn:
return None
try:
engine = create_engine(dsn, future=True)
with engine.connect() as conn:
conn.execute(text("SELECT proxy_id FROM scrape_runs LIMIT 1"))
return engine
except Exception:
return None
_ENGINE = _live_engine()
pytestmark = pytest.mark.skipif(_ENGINE is None, reason="no reachable Postgres test DB")
@pytest.fixture
def env(monkeypatch: pytest.MonkeyPatch) -> Iterator[dict[str, Any]]:
"""Узел пула + два прогона: `run` — текущий, `other` — строка, которую вызывающий
держит изменённой без коммита."""
assert _ENGINE is not None
factory = sessionmaker(bind=_ENGINE, future=True)
monkeypatch.setattr(proxy_egress, "_SessionLocal", factory)
setup = factory()
proxy_id = setup.execute(
text(
"INSERT INTO scrape_proxies (url, provider_affinity) "
"VALUES ('http://t3404-' || gen_random_uuid(), 'any') RETURNING id"
)
).scalar_one()
run_ids = [
setup.execute(
text(
"INSERT INTO scrape_runs (source, status) "
"VALUES ('test_3404', 'running') RETURNING id"
)
).scalar_one()
for _ in range(2)
]
setup.commit()
token = current_run_id.set(None)
try:
yield {"proxy_id": int(proxy_id), "run": int(run_ids[0]), "other": int(run_ids[1])}
finally:
current_run_id.reset(token)
setup.rollback()
setup.execute(text("DELETE FROM scrape_runs WHERE id = ANY(:ids)"), {"ids": run_ids})
setup.execute(text("DELETE FROM scrape_proxies WHERE id = :id"), {"id": proxy_id})
setup.commit()
setup.close()
def _run_row(run_id: int) -> Any:
assert _ENGINE is not None
with _ENGINE.connect() as conn:
return conn.execute(
text("SELECT proxy_id, counters FROM scrape_runs WHERE id = :id"), {"id": run_id}
).one()
def _node_id_of(url: str) -> int:
assert _ENGINE is not None
with _ENGINE.connect() as conn:
return int(
conn.execute(
text("SELECT id FROM scrape_proxies WHERE url = :u"), {"u": url}
).scalar_one()
)
def test_resolve_within_run_writes_node_to_scrape_runs(env: dict[str, Any]) -> None:
"""Главный случай: прогон идёт, egress выбран — прогон знает узел.
На старом коде proxy_id остаётся NULL."""
current_run_id.set(env["run"])
caller = Session(bind=_ENGINE, future=True)
try:
url = proxy_egress.resolve_proxy_url(caller, "yandex")
finally:
caller.close()
assert url is not None
node = _node_id_of(url)
row = _run_row(env["run"])
assert row.proxy_id == node
counters = row.counters if isinstance(row.counters, dict) else json.loads(row.counters)
assert counters.get("proxy_ids") == [node]
def test_resolve_outside_run_writes_nothing(env: dict[str, Any]) -> None:
"""Без прогона (админка, проверка кук) — scrape_runs не трогаем."""
caller = Session(bind=_ENGINE, future=True)
try:
assert proxy_egress.resolve_proxy_url(caller, "yandex") is not None
finally:
caller.close()
assert _run_row(env["run"]).proxy_id is None
def test_attribution_does_not_commit_callers_transaction(env: dict[str, Any]) -> None:
"""db вызывающего — долгоживущая сессия прогона посреди работы. Атрибуция не имеет
права закоммитить его незавершённые изменения (или откатить их на сбое)."""
current_run_id.set(env["run"])
caller = Session(bind=_ENGINE, future=True)
try:
caller.execute(
text("UPDATE scrape_runs SET counters = '{\"t3404\": 1}'::jsonb WHERE id = :id"),
{"id": env["other"]},
)
proxy_egress.resolve_proxy_url(caller, "yandex")
caller.rollback()
finally:
caller.close()
other = _run_row(env["other"]).counters or {}
assert "t3404" not in other, "атрибуция закоммитила чужую транзакцию"
assert _run_row(env["run"]).proxy_id is not None

View file

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

View file

@ -127,10 +127,10 @@ def test_flag_on_exception_marks_fail_and_still_releases() -> None:
assert spy.release_calls == [7] 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) внутри блока — """Исключение — подкласс ProxyBanError (напр. AvitoBlockedError) внутри блока —
вызывает mark_banned(lease, source=provider) В ДОПОЛНЕНИЕ к mark_health(ok=False) вызывает mark_banned(lease, source=provider) ВМЕСТО mark_health(ok=False)
(#2600 п.1 curl-путь). Zero изменений в вызывающем коде — сигнал детектируется (#2600 п.1 curl-путь, #3310). Zero изменений в вызывающем коде — сигнал детектируется
по ТИПУ исключения, а не явным вызовом.""" по ТИПУ исключения, а не явным вызовом."""
class _FakeBlockedError(ProxyBanError): class _FakeBlockedError(ProxyBanError):
@ -143,7 +143,8 @@ def test_ban_exception_calls_mark_banned_in_addition_to_mark_health() -> None:
assert url == _LEASE.url assert url == _LEASE.url
raise _FakeBlockedError("firewall page detected") raise _FakeBlockedError("firewall page detected")
assert spy.mark_banned_calls == [(7, "avito")] 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 всё равно освобождён assert spy.release_calls == [7] # lease всё равно освобождён

View file

@ -587,9 +587,18 @@ def save_listings(
price_rub = EXCLUDED.price_rub, price_rub = EXCLUDED.price_rub,
price_per_m2 = EXCLUDED.price_per_m2, price_per_m2 = EXCLUDED.price_per_m2,
-- Cian-specific: обновляем при каждом re-scrape -- 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, bedrooms_count = EXCLUDED.bedrooms_count,
balconies_count = EXCLUDED.balconies_count, balconies_count = COALESCE(
EXCLUDED.balconies_count, listings.balconies_count
),
loggias_count = EXCLUDED.loggias_count, loggias_count = EXCLUDED.loggias_count,
description_minhash = EXCLUDED.description_minhash, description_minhash = EXCLUDED.description_minhash,
cadastral_number = EXCLUDED.cadastral_number, cadastral_number = EXCLUDED.cadastral_number,
@ -742,8 +751,11 @@ def save_listings(
listings.price_previous_rub, listings.newbuilding_id, listings.price_previous_rub, listings.newbuilding_id,
listings.newbuilding_url, listings.card_hash, listings.is_active listings.newbuilding_url, listings.card_hash, listings.is_active
) IS DISTINCT FROM ( ) IS DISTINCT FROM (
EXCLUDED.price_rub, EXCLUDED.price_per_m2, EXCLUDED.living_area_m2, EXCLUDED.price_rub, EXCLUDED.price_per_m2,
EXCLUDED.bedrooms_count, EXCLUDED.balconies_count, EXCLUDED.loggias_count, 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.description_minhash, EXCLUDED.cadastral_number,
EXCLUDED.building_cadastral_number, EXCLUDED.building_cadastral_number,
-- #3063: те же COALESCE, что в SET выше. Правая часть гейта ОБЯЗАНА -- #3063: те же COALESCE, что в SET выше. Правая часть гейта ОБЯЗАНА
@ -823,9 +835,9 @@ def save_listings(
is_active = true, is_active = true,
price_rub = :price_rub, price_rub = :price_rub,
price_per_m2 = :ppm2, price_per_m2 = :ppm2,
living_area_m2 = :living_area_m2, living_area_m2 = COALESCE(:living_area_m2, living_area_m2),
bedrooms_count = :bedrooms_count, bedrooms_count = :bedrooms_count,
balconies_count = :balconies_count, balconies_count = COALESCE(:balconies_count, balconies_count),
loggias_count = :loggias_count, loggias_count = :loggias_count,
description_minhash = :description_minhash, description_minhash = :description_minhash,
cadastral_number = :cadastral_number, cadastral_number = :cadastral_number,
@ -940,7 +952,7 @@ def save_listings(
skip_snapshot = today_row is not None skip_snapshot = today_row is not None
if skip_snapshot: if skip_snapshot:
logger.debug( 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, lot.source,
listing_id, listing_id,
) )

View file

@ -18,11 +18,15 @@ release ВСЕГДА в finally — lease не должен течь, даже
Бан площадки (#2600 п.1): если исключение, поднятое ИЗНУТРИ `with curl_proxy_url(...) as Бан площадки (#2600 п.1): если исключение, поднятое ИЗНУТРИ `with curl_proxy_url(...) as
url:`, `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/ url:`, `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
`DomClickBlockedError` и т.п. см. `proxy_errors.ProxyBanError`), это НЕ просто `DomClickBlockedError` и т.п. см. `proxy_errors.ProxyBanError`), это НЕ
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный `mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
`mark_banned` узел сразу снимается с выдачи ЭТОМУ провайдеру (per-source бан, #2600 п.2; `mark_banned` узел сразу снимается с выдачи ЭТОМУ провайдеру (per-source бан, #2600 п.2;
для остальных источников остаётся в строю), кроме случая когда это последний узел, для остальных источников остаётся в строю), кроме случая когда это последний узел,
достижимый для провайдера (защита в `app.services.proxy_pool.mark_banned`). достижимый для провайдера (защита в `app.services.proxy_pool.mark_banned`).
`mark_health(ok=False)` на бане НЕ зовётся (#3310): узел исправен, его отбила площадка, а
глобальный счётчик после трёх банов выводил его из выдачи ВСЕМ источникам в том числе
последний узел, который защита только что пообещала оставить. Тот же выбор, что у
`BrowserFetcher.report_platform_ban` (#3288).
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
ИЗНУТРИ блока, получает сигнал бесплатно этот модуль намеренно НЕ импортирует ИЗНУТРИ блока, получает сигнал бесплатно этот модуль намеренно НЕ импортирует
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные
@ -126,17 +130,19 @@ def curl_proxy_url(
raise raise
finally: finally:
# mark_health/mark_banned/release — best-effort: проблема пула не должна # mark_health/mark_banned/release — best-effort: проблема пула не должна
# ронять сбор. mark_banned ПЕРЕД mark_health(ok=False) — оба независимы # ронять сбор. Бан ВМЕСТО mark_health(ok=False), а не вдобавок (#3310): отказ
# (разные поля), но бан — более специфичный/сильный сигнал. # площадки — свойство пары «узел × источник», глобальный счётчик здоровья он
# копить не должен (см. докстринг модуля).
if banned: if banned:
try: try:
proxy_provider.mark_banned(lease, source=provider) proxy_provider.mark_banned(lease, source=provider)
except Exception: except Exception:
logger.warning("proxy_pool: mark_banned failed for %s", provider, exc_info=True) logger.warning("proxy_pool: mark_banned failed for %s", provider, exc_info=True)
try: else:
proxy_provider.mark_health(lease, ok) try:
except Exception: proxy_provider.mark_health(lease, ok)
logger.warning("proxy_pool: mark_health failed for %s", provider, exc_info=True) except Exception:
logger.warning("proxy_pool: mark_health failed for %s", provider, exc_info=True)
try: try:
proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт
except Exception: except Exception:
@ -175,10 +181,10 @@ async def acurl_proxy_url(
ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена
приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им
аренду замер в тестах: `wait_for(timeout=0.05)` вернулся через 0.3 с (всё время аренду замер в тестах: `wait_for(timeout=0.05)` вернулся через 0.3 с (всё время
acquire). Худший случай `/estimate` три источника × до 30 с checkout'а коннекта из acquire). Худший случай `/estimate` три источника × до `pool_timeout` checkout'а
исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток не прервать. коннекта из исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток
Закрывается со стороны БД `statement_timeout`/`pool_timeout` короче бюджета, не прервать. Закрыто со стороны БД (#3408, PR #3444): `pool_timeout` = 5 с — короче
follow-up #3408 («пул коннектов SQLAlchemy 15 против пикового спроса до 28»). самого короткого бюджета источника (8 с), см. app/core/db.py.
""" """
cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url) cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url)
enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__)) enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__))

View file

@ -32,7 +32,7 @@ from scraper_kit.offer_price_history import (
validate_diff_percent, validate_diff_percent,
) )
from scraper_kit.providers._base import build_curl_cffi_session 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.proxy_errors import caused_by_no_proxy
from scraper_kit.repair_state_normalizer import ( from scraper_kit.repair_state_normalizer import (
infer_repair_state_from_text, infer_repair_state_from_text,
@ -142,7 +142,7 @@ def _is_captcha_title(title: str) -> bool:
"""Заголовок — капча Циана? Один предикат на оба места детекта. """Заголовок — капча Циана? Один предикат на оба места детекта.
Детект живёт в двух точках намеренно (#3402 follow-up): общий parse-путь ниже — Детект живёт в двух точках намеренно (#3402 follow-up): общий parse-путь ниже —
ЗА границей `with curl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула ЗА границей `async with acurl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула
уже не доходит (узел успел получить `mark_health(ok=True)` и остаётся в выдаче уже не доходит (узел успел получить `mark_health(ok=True)` и остаётся в выдаче
Циану ровно дефект #2700, только на HTTP 200). Поэтому curl-путь спрашивает про Циану ровно дефект #2700, только на HTTP 200). Поэтому curl-путь спрашивает про
капчу ВНУТРИ блока, а не после него. капчу ВНУТРИ блока, а не после него.
@ -228,9 +228,13 @@ async def fetch_detail(
# curl_cffi own-session path (legacy, back-compat). # curl_cffi own-session path (legacy, back-compat).
# proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian. # proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian.
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url. # Прокси: пул за флагом 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 _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( own_session = build_curl_cffi_session(
proxy_url=_proxy_url, proxy_url=_proxy_url,
timeout=25.0, timeout=25.0,
@ -242,7 +246,7 @@ async def fetch_detail(
try: try:
resp = await own_session.get(offer_url, allow_redirects=True) resp = await own_session.get(offer_url, allow_redirects=True)
if resp.status_code != 200: if resp.status_code != 200:
# ВНУТРИ curl_proxy_url: поднятый отсюда ProxyBanError доходит до # ВНУТРИ acurl_proxy_url: поднятый отсюда ProxyBanError доходит до
# пула (mark_banned на пару «узел × cian», #2600 п.2). Раньше здесь # пула (mark_banned на пару «узел × cian», #2600 п.2). Раньше здесь
# был `return None` — узел получал mark_health(ok=True) и оставался # был `return None` — узел получал mark_health(ok=True) и оставался
# в выдаче Циану (#2700, 15 суток по 50 отказов в сутки). # в выдаче Циану (#2700, 15 суток по 50 отказов в сутки).
@ -252,7 +256,7 @@ async def fetch_detail(
html = resp.text html = resp.text
# Капча приходит с HTTP 200, то есть мимо `_raise_if_blocked`. Спрашиваем # Капча приходит с HTTP 200, то есть мимо `_raise_if_blocked`. Спрашиваем
# ЗДЕСЬ, пока lease жив: общий parse-путь ниже поднимет то же исключение # ЗДЕСЬ, пока lease жив: общий parse-путь ниже поднимет то же исключение
# уже после выхода из `with curl_proxy_url`, где узлу проставлен # уже после выхода из `async with acurl_proxy_url`, где узлу проставлен
# mark_health(ok=True) и бана пары «узел × cian» не будет (#3402). # mark_health(ok=True) и бана пары «узел × cian» не будет (#3402).
title = _page_title(html) title = _page_title(html)
if _is_captcha_title(title): if _is_captcha_title(title):

View file

@ -27,7 +27,8 @@ _extract_ssr_state делает balanced-brace scan (с пропуском ск
Форма SSR-стейта подтверждена на живой карточке 2075729321 (2026-06-27): Форма SSR-стейта подтверждена на живой карточке 2075729321 (2026-06-27):
productCard.objectInfo.{renovation,livingArea,kitchenArea}, productCard.priceInfo. productCard.objectInfo.{renovation,livingArea,kitchenArea}, productCard.priceInfo.
priceHistory ({date ISO8601+tz, price, diff, state}), productCard.egrnData 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}, saleType, productCard.viewsCount/callsCount, houseInfo.info.{wallType,floorType},
и ТОП-УРОВНЕМ pricePrediction (DomClick AVM raw_extra.avm). и ТОП-УРОВНЕМ pricePrediction (DomClick AVM raw_extra.avm).
@ -283,6 +284,29 @@ def _to_int(value: Any) -> int | None:
return 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: def _parse_change_time(value: Any) -> datetime | None:
"""Нормализует дату изменения цены → timezone-aware datetime. """Нормализует дату изменения цены → timezone-aware datetime.
@ -436,8 +460,8 @@ def parse_detail_html(html: str, source_url: str) -> DomClickDetailEnrichment:
repair_state = _map_repair_state(repair_type) repair_state = _map_repair_state(repair_type)
# площади # площади
living_area_m2 = _to_float(oi.get("livingArea")) living_area_m2 = _pos_float(oi.get("livingArea"))
kitchen_area_m2 = _to_float(oi.get("kitchenArea")) kitchen_area_m2 = _pos_float(oi.get("kitchenArea"))
# балконы: int → count+флаг; иное (строка/др.) → raw_extra, counts None # балконы: int → count+флаг; иное (строка/др.) → raw_extra, counts None
balconies = oi.get("balconies") 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") sale_type = (pc.get("legalOptions") or {}).get("saleType")
# owners_count — защищаемся от двух написаний ключа # owners_count — обёртка {status, value} (см. _egrn_value); до #3252 int()
owners_count = _to_int(egrn.get("owners_count")) # от словаря молча давал None: 0 из ~2970 карточек, обогащённых с 29.08.
owners_count = _to_int(_egrn_value(egrn.get("owners_count")))
if owners_count is None: 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 обозначают # encumbrances_clean (ИНВЕРСИЯ): collateral/collateral_sber обозначают
# НАЛИЧИЕ обременения. Любой truthy → clean=False; оба отсутствуют/false # НАЛИЧИЕ обременения. Любой 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 ОСТАЮТСЯ тут, # raw_extra → merge в listings.raw_payload. wall/floor type ОСТАЮТСЯ тут,
# НЕ пишем в listings.house_type (кросс-источниковый словарь грязный). # НЕ пишем в listings.house_type (кросс-источниковый словарь грязный).
# egrn area key — точное написание не подтверждено; пробуем несколько. # egrnData.area — ключ подтверждён на живых карточках 17.09 (#3252). В raw_payload
egrn_area = egrn.get("area") or egrn.get("rosreestrArea") or egrn.get("object_area") # кладём обёртку {status, value} целиком: status говорит, сошлась ли площадь
# ЕГРН с площадью объявления, и уже лежит в таком виде у 2156 строк.
egrn_area = egrn.get("area")
demand = _compact( demand = _compact(
{ {
"calls": pc.get("callsCount"), "calls": pc.get("callsCount"),