diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index 9030ca16..f6cb9b4a 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -750,6 +750,13 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]: post_claim=reschedule_after_minutes(param="interval_minutes", default=360), ), "rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"), + # Wildcard (#3051 п.3): rosreestr_dkp_import_77 (Москва, миграция 288) и любой + # будущий region_code-суффикс из той же семьи резолвятся сюда через + # resolve_handler по префиксу (тот же механизм, что deactivate_stale_* / + # avito_city_sweep_* — см. scraper_kit.orchestration.scheduler.resolve_handler). + # region_code берётся из default_params строки расписания (import_rosreestr_dkp + # сам валидирует его через app.services.regions.REGIONS), Handler-тело общее. + "rosreestr_dkp_import_*": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import_*"), "listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"), "asking_to_sold_ratio_refresh": Handler( _job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh" diff --git a/tradein-mvp/backend/app/services/regions.py b/tradein-mvp/backend/app/services/regions.py index c7b4e3c6..380c7ea0 100644 --- a/tradein-mvp/backend/app/services/regions.py +++ b/tradein-mvp/backend/app/services/regions.py @@ -47,6 +47,15 @@ class Region: Регион без тира должен деградировать ЯВНО (потребитель спрашивает unsupported_tier_reason и логирует/маркирует), а не молча считать дальше без источника. + canonical_city — #3051: имя города, которым ПЕРЕЗАПИСЫВАЕТСЯ `city` + строк, приходящих из источника без надёжного city-поля + (Росреестр по Москве отдаёт муниципальный округ/поселение + вместо города — «Раменки», «Сосенское» — а не «Москва»). + None — источник несёт свой city как есть, без override + (регион 66: byte-for-byte прежнее поведение). Not-None — + потребитель (import_rosreestr_dkp) подставляет это имя + вместо city источника и НЕ фильтрует по city IS NOT NULL + (иначе на 77 теряется ~10% строк с пустым city). """ code: int @@ -58,6 +67,7 @@ class Region: city_token: str cities: frozenset[str] enrichment_tiers: frozenset[str] + canonical_city: str | None = None def is_within_bbox(lat: float, lon: float, bbox: BBox) -> bool: @@ -144,6 +154,10 @@ REGIONS: dict[int, Region] = { # sber_index покрывают регион 66. Пустое множество здесь — не заглушка, # а ФАКТ, который потребители обязаны озвучивать (см. класс-докстринг). enrichment_tiers=frozenset(), + # #3051: Росреестр по Москве отдаёт в city муниципальный округ/поселение + # ("муниципальный округ Раменки", "поселение Сосенское"), не сам город — + # import_rosreestr_dkp подставляет каноничное имя вместо city источника. + canonical_city="Москва", ), } diff --git a/tradein-mvp/backend/app/services/scheduler.py b/tradein-mvp/backend/app/services/scheduler.py index 0f67306d..4eaa4e1b 100644 --- a/tradein-mvp/backend/app/services/scheduler.py +++ b/tradein-mvp/backend/app/services/scheduler.py @@ -25,6 +25,7 @@ Zombie-reap, advisory-lock claim и tick-loop теперь целиком в from __future__ import annotations +import json import logging from typing import Any @@ -50,6 +51,7 @@ from sqlalchemy.orm import Session from app.core.shutdown import shutdown_requested from app.services import scrape_runs as runs_mod +from app.services.regions import REGIONS __all__ = ["compute_next_run_at", "has_running_run"] @@ -238,6 +240,24 @@ async def _execute_cian_backfill( _DKP_SOURCE = "rosreestr_dkp_import" + +def _dkp_source_for_region(region_code: int) -> str: + """Имя scrape_runs.source для чекпоинта данного региона (#3051 п.3). + + 66 — байт-в-байт прежнее имя ('rosreestr_dkp_import'), под которым годами + писались scrape_runs. Остальные регионы получают суффикс кода — тот же + формат, что и у строки scrape_schedules ('rosreestr_dkp_import_77', + seed — миграция 288), которую резолвит wildcard 'rosreestr_dkp_import_*' + в product_handlers.py. Изоляция чекпоинтов между регионами держится именно + на разных source: _resume_dkp_cursor ищет ПРЕДЫДУЩИЙ прогон с ТЕМ ЖЕ source, + поэтому курсор региона 77 никогда не подхватит last_id региона 66 (и + наоборот) — они просто разные строки в scrape_runs.source. + """ + if region_code == 66: + return _DKP_SOURCE + return f"{_DKP_SOURCE}_{region_code}" + + # Потолок возраста чекпоинта: старше — last_id прошлого прогона не подхватываем, прогон # стартует с id=0 (issue #3168). У предиката `id > last_id` нет протухания в смысле # свипов (он остаётся корректным сколь угодно долго), но апстрим @@ -263,7 +283,9 @@ _DKP_RESUME_CANDIDATE_SQL = text(""" """) -def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]: +def _resume_dkp_cursor( + db: Session, run_id: int, source: str = _DKP_SOURCE +) -> tuple[int, dict[str, Any]]: """Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168). last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat @@ -271,7 +293,14 @@ def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]: обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял пере-сканировать источник с начала. - Кандидат — ПОСЛЕДНИЙ прогон этого source (тот же принцип, что и + `source` (#3051 п.3) — per-region ключ чекпоинта (см. _dkp_source_for_region): + дефолт _DKP_SOURCE сохраняет прежнее поведение вызовов без явного аргумента + (регион 66). Кандидат ищется СТРОГО по этому source — прогон региона 77 + (source='rosreestr_dkp_import_77') никогда не видит last_id региона 66 + (source='rosreestr_dkp_import') и наоборот: разные регионы физически не + матчат друг друга в WHERE source = :source ниже. + + Кандидат — ПОСЛЕДНИЙ прогон ЭТОГО source (тот же принцип, что и scraper_kit.orchestration.scheduler._pick_resume, локальная копия ладдера — контракт другой: нет params/interval_days, курсор числовой, а не bucket-set): - 'running' / 'zombie' — прогон, которого не завершили штатно. @@ -285,7 +314,7 @@ def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]: кодом (kit_runs.update_heartbeat — merge, не замена), чтобы решение было видно в scrape_runs, а не только в логе. """ - row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": _DKP_SOURCE, "rid": run_id}).fetchone() + row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": source, "rid": run_id}).fetchone() verdict: dict[str, Any] = {"resume_from": None} if row is None: @@ -322,22 +351,35 @@ def import_rosreestr_dkp( ) -> None: """Import ДКП-сделок из gendesign rosreestr_deals через postgres_fdw. - Python-порт import-rosreestr.sh (Variant C из #563). + Python-порт import-rosreestr.sh (Variant C из #563). #3051 п.3: параметризовано + по региону (params["region_code"], реестр — app.services.regions.REGIONS) — + было хардкод region_code=66. - Источник: foreign table gendesign_rosreestr_deals (создана в migration 072). + Источник: foreign table gendesign_rosreestr_deals (создана в migration 072, + okato/quarter_cad_number/district добавлены миграцией 288). SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql. USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader). - Область покрытия: вся Свердловская область (region_code=66), не только Екатеринбург — - прежний ILIKE-фильтр по подстроке города (ограничивавший импорт одним Екатеринбургом) - снят (Mera trade-in расширяется на весь регион, unlocks +47183 сделок вне ЕКБ уже - сидящих в source foreign table). address и deals.city строятся из реального city - источника (не хардкод "Екатеринбург"), deals.region_code заполняется из строки - источника (= 66 при текущем фильтре). + Область покрытия: region_code из params (default 66 — вся Свердловская область, + не только Екатеринбург). Неизвестный код региона (нет в REGIONS) — ValueError, + прогон падает явно, а не молча импортирует мусор с чужим region_code. + + region_code=66 (регион БЕЗ canonical_city в реестре) — поведение байт-в-байт + прежнее: city/address строятся из city источника, обязателен фильтр + city IS NOT NULL AND trim(city) != ''. + + Регион С canonical_city (77 — Москва): Росреестр отдаёт в city муниципальный + округ/поселение ("муниципальный округ Раменки", "поселение Сосенское"), НЕ + город — city/address подставляют region.canonical_city, а не city источника; + фильтр city IS NOT NULL НЕ применяется (иначе теряется ~10% строк с пустым + city источника). Исходные city/okato/quarter_cad_number/district уходят в + raw_payload (jsonb) — единственная ветка SQL решает это через bind-параметр + :canonical_city (CASE WHEN ... IS NOT NULL), а не отдельный Python if/else на + конкретный код региона. Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24): - - region_code = 66 (вся Свердловская область, все города) - - city IS NOT NULL AND trim(city) != '' (непустой город → корректный address) + - region_code = :region_code (параметризовано, было хардкод 66) + - city IS NOT NULL AND trim(city) != '' — ТОЛЬКО если у региона нет canonical_city - realestate_type_code = '002001003000' (квартира) - area BETWEEN 18 AND 200 - deal_price BETWEEN 1000000 AND 100000000 @@ -354,17 +396,33 @@ def import_rosreestr_dkp( (WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint), мержем (kit_runs.update_heartbeat), а не заменой. На старте _resume_dkp_cursor решает продолжить с last_id прошлого прогона или начать с 0 — чекпоинт переживает рестарт - процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168). + процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168). Курсор — ПЕР + РЕГИОН (#3051 п.3): _resume_dkp_cursor вызывается с source=_dkp_source_for_region + (region_code), поэтому last_id региона 77 никогда не подхватывает last_id региона + 66 — они разные scrape_runs.source ('rosreestr_dkp_import' vs + 'rosreestr_dkp_import_77'), см. докстринг _dkp_source_for_region. SAVEPOINT per row — один сбойный row не откатывает батч. Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals). TODO (follow-up): запустить geocode backfill после import. - Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549). + Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549, + region-agnostic — синтетические строки существовали только для ЕКБ). """ since: str = str(params.get("since", "2024-01-01")) batch_size: int = int(params.get("batch_size", 2000)) + region_code: int = int(params.get("region_code", 66)) + region = REGIONS.get(region_code) + if region is None: + raise ValueError( + f"rosreestr_dkp_import: region_code={region_code} не найден в " + f"app.services.regions.REGIONS (известны: {sorted(REGIONS)}) — " + "прогон остановлен, чтобы не импортировать сделки с неизвестным " + "региональным контекстом (city/address-правила для него не определены)" + ) + dkp_source = _dkp_source_for_region(region_code) + counters: dict[str, int] = { "rows_fetched": 0, "rows_inserted": 0, @@ -400,7 +458,7 @@ def import_rosreestr_dkp( ) db.rollback() - last_id, resume_verdict = _resume_dkp_cursor(db, run_id) + last_id, resume_verdict = _resume_dkp_cursor(db, run_id, source=dkp_source) total_batches = 0 kit_runs.update_heartbeat(db, run_id, resume_verdict) logger.info( @@ -450,9 +508,20 @@ def import_rosreestr_dkp( id, id AS source_id_src, 'ros:dkp:' || CAST(id AS text) AS dedup_hash, - trim(city) || ', ' || trim(street) AS address, + -- #3051: регион с canonical_city (Москва) подставляет его вместо + -- city источника (округ/поселение, не город) — CASE на bind-параметре, + -- не Python if/else на код региона. + CASE + WHEN CAST(:canonical_city AS text) IS NOT NULL + THEN CAST(:canonical_city AS text) || ', ' || trim(street) + ELSE trim(city) || ', ' || trim(street) + END AS address, region_code, - trim(city) AS city, + CASE + WHEN CAST(:canonical_city AS text) IS NOT NULL + THEN CAST(:canonical_city AS text) + ELSE trim(city) + END AS city, CASE WHEN area < 30 THEN 0 WHEN area < 44 THEN 1 @@ -472,10 +541,27 @@ def import_rosreestr_dkp( year_build AS year_built, round(deal_price)::bigint AS price_rub, round(price_per_sqm)::int AS price_per_m2, - period_start_date AS deal_date + period_start_date AS deal_date, + doc_type, + -- Исходный city/okato/quarter_cad_number/district — ТОЛЬКО когда + -- city перезаписан canonical_city выше (иначе NULL, регион 66 + -- byte-for-byte прежний: raw_payload не заполнялся и не заполняется). + CASE + WHEN CAST(:canonical_city AS text) IS NOT NULL THEN + jsonb_build_object( + 'src_city', city, + 'okato', okato, + 'quarter_cad_number', quarter_cad_number, + 'district', district + ) + ELSE NULL + END AS raw_payload FROM gendesign_rosreestr_deals - WHERE region_code = 66 - AND city IS NOT NULL AND trim(city) <> '' + WHERE region_code = CAST(:region_code AS int) + AND ( + CAST(:canonical_city AS text) IS NOT NULL + OR (city IS NOT NULL AND trim(city) <> '') + ) AND realestate_type_code = '002001003000' AND area BETWEEN 18 AND 200 AND deal_price BETWEEN 1000000 AND 100000000 @@ -486,7 +572,13 @@ def import_rosreestr_dkp( ORDER BY id LIMIT CAST(:batch_size AS int) """), - {"since": since, "last_id": last_id, "batch_size": batch_size}, + { + "since": since, + "last_id": last_id, + "batch_size": batch_size, + "region_code": region_code, + "canonical_city": region.canonical_city, + }, ) .mappings() .all() @@ -524,7 +616,7 @@ def import_rosreestr_dkp( INSERT INTO deals ( source, dedup_hash, source_id, address, region_code, city, rooms, area_m2, floor, year_built, price_rub, price_per_m2, - deal_date + deal_date, doc_type, raw_payload ) VALUES ( 'rosreestr', @@ -539,7 +631,9 @@ def import_rosreestr_dkp( CAST(:year_built AS int), CAST(:price_rub AS bigint), CAST(:price_per_m2 AS int), - CAST(:deal_date AS date) + CAST(:deal_date AS date), + CAST(:doc_type AS text), + CAST(:raw_payload AS jsonb) ) ON CONFLICT (dedup_hash) DO UPDATE SET address = EXCLUDED.address, @@ -551,7 +645,9 @@ def import_rosreestr_dkp( year_built = EXCLUDED.year_built, price_rub = EXCLUDED.price_rub, price_per_m2 = EXCLUDED.price_per_m2, - deal_date = EXCLUDED.deal_date + deal_date = EXCLUDED.deal_date, + doc_type = EXCLUDED.doc_type, + raw_payload = EXCLUDED.raw_payload WHERE deals.address IS DISTINCT FROM EXCLUDED.address OR deals.region_code IS DISTINCT FROM EXCLUDED.region_code OR deals.city IS DISTINCT FROM EXCLUDED.city @@ -562,6 +658,8 @@ def import_rosreestr_dkp( OR deals.price_rub IS DISTINCT FROM EXCLUDED.price_rub OR deals.price_per_m2 IS DISTINCT FROM EXCLUDED.price_per_m2 OR deals.deal_date IS DISTINCT FROM EXCLUDED.deal_date + OR deals.doc_type IS DISTINCT FROM EXCLUDED.doc_type + OR deals.raw_payload IS DISTINCT FROM EXCLUDED.raw_payload RETURNING (xmax = 0) AS was_inserted """), { @@ -577,6 +675,12 @@ def import_rosreestr_dkp( "price_rub": row["price_rub"], "price_per_m2": row["price_per_m2"], "deal_date": row["deal_date"], + "doc_type": row["doc_type"], + "raw_payload": ( + json.dumps(row["raw_payload"], ensure_ascii=False) + if row["raw_payload"] is not None + else None + ), }, ).fetchone() if result is None: diff --git a/tradein-mvp/backend/data/sql/288_deals_doc_type.sql b/tradein-mvp/backend/data/sql/288_deals_doc_type.sql new file mode 100644 index 00000000..948b2a09 --- /dev/null +++ b/tradein-mvp/backend/data/sql/288_deals_doc_type.sql @@ -0,0 +1,87 @@ +-- 288_deals_doc_type.sql +-- deals.doc_type (#3051 п.3) + foreign table gendesign_rosreestr_deals: okato/ +-- quarter_cad_number/district + disabled seed schedule rosreestr_dkp_import_77. +-- +-- Dependencies: 002_core_tables.sql (deals), 072_scrape_schedules_seed_cian_rosreestr.sql +-- (gendesign_rosreestr_deals, scrape_schedules). +-- Apply after: 287_proxy_run_attribution.sql +-- +-- WHY: +-- Трек 2 подготовки Mera к Москве — импорт сделок Росреестра параметризуется по +-- региону (66 Свердловская обл. / 77 Москва, code-часть в scheduler.py). deals +-- до сих пор не различал ДКП (вторичка) от ДДУ (застройщик) в самой строке — +-- различие жило только в WHERE-фильтре импортёра. Явная колонка нужна для +-- будущего ДДУ-импорта (#3051 п.3) и для аналитики, которая иначе не может +-- отличить типы сделок в одной таблице. +-- +-- Foreign table расширена тремя колонками источника (okato, quarter_cad_number, +-- district) — они нужны import_rosreestr_dkp для raw_payload по регионам, где +-- city источника не используется как есть (Москва: city = муниципальный +-- округ/поселение, не город). Существование колонок в public.rosreestr_deals +-- на gendesign-стороне проверено live (2026-09-08, prod psql). +-- +-- Seed-строка rosreestr_dkp_import_77 — ВЫКЛЮЧЕНА (enabled=false): миграция +-- только заводит расписание, включение и первый прогон по Москве — отдельное +-- решение main-сессии после ревью кода-части. +-- +-- ИДЕМПОТЕНТНОСТЬ: ADD COLUMN IF NOT EXISTS × 4, COMMENT ON COLUMN (безусловны, +-- но перезаписывают тот же текст), ON CONFLICT (source) DO NOTHING для seed. +-- Backfill (doc_type='ДКП' WHERE source='rosreestr' AND doc_type IS NULL) — +-- повторный прогон no-op (второй раз IS NULL уже не матчит ни одну строку). + +BEGIN; + +SET LOCAL lock_timeout = '5s'; + +-- ── deals.doc_type ──────────────────────────────────────────────────────────── + +ALTER TABLE deals + ADD COLUMN IF NOT EXISTS doc_type text; + +-- Backfill: весь текущий rosreestr-импорт в deals уже отфильтрован по doc_type='ДКП' +-- на стороне import_rosreestr_dkp (WHERE doc_type = 'ДКП') — строки, попавшие в deals +-- ДО этой миграции, были все ДКП, ДДУ-импорта ещё нет ни одной строки. +UPDATE deals + SET doc_type = 'ДКП' + WHERE source = 'rosreestr' + AND doc_type IS NULL; + +COMMENT ON COLUMN deals.doc_type IS + 'Тип документа сделки Росреестра: ДКП (договор купли-продажи, вторичка) или ' + 'ДДУ (договор долевого участия, застройщик) — #3051 п.3. NULL для строк не из ' + 'rosreestr (etazhi/domklik_history) — у них своя типизация сделки, либо не ' + 'применимо. Заполняется import_rosreestr_dkp из значения источника (WHERE ' + 'уже фильтрует doc_type=''ДКП'', колонка носит его явно, а не только в фильтре).'; + +-- ── gendesign_rosreestr_deals: колонки под региональный raw_payload (#3051) ──── +-- ALTER FOREIGN TABLE ADD COLUMN — только локальные метаданные (не трогает +-- реальную remote-таблицу), безопасно как ALTER TABLE ADD COLUMN без DEFAULT. + +ALTER FOREIGN TABLE gendesign_rosreestr_deals + ADD COLUMN IF NOT EXISTS okato text; + +ALTER FOREIGN TABLE gendesign_rosreestr_deals + ADD COLUMN IF NOT EXISTS quarter_cad_number text; + +ALTER FOREIGN TABLE gendesign_rosreestr_deals + ADD COLUMN IF NOT EXISTS district text; + +-- ── Seed: rosreestr_dkp_import_77 (выключено) ────────────────────────────────── + +INSERT INTO scrape_schedules ( + source, + enabled, + window_start_hour, + window_end_hour, + default_params +) +VALUES ( + 'rosreestr_dkp_import_77', + false, + 4, + 6, + '{"region_code": 77, "since": "2024-01-01", "batch_size": 2000}'::jsonb +) +ON CONFLICT (source) DO NOTHING; + +COMMIT; diff --git a/tradein-mvp/backend/tests/test_3051_rosreestr_region_param.py b/tradein-mvp/backend/tests/test_3051_rosreestr_region_param.py new file mode 100644 index 00000000..e35c72e1 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3051_rosreestr_region_param.py @@ -0,0 +1,247 @@ +"""#3051 п.3: параметризация import_rosreestr_dkp по региону + deals.doc_type. + +Трек 2 подготовки Mera к Москве (region_code=77). Чисто-юнит: SQL-текст +(inspect.getsource), _FakeDb-двойники для чекпоинта, миграция 288 и реестр +регионов — без живого FDW/Postgres (тот же стиль, что test_rosreestr_dedup_key.py +и test_3168_backfill_cursor_resume.py). + +Покрывает пункты задачи: + (a) region_code больше не литерал 66 в SQL — bind-параметр (см. также + test_rosreestr_dedup_key.py::test_live_import_region_code_is_bind_param). + (b) маппинг 77: city='Москва', address с префиксом, raw_payload с src_city/okato/ + quarter_cad_number/district. + (c) маппинг 66 не изменился (canonical_city=None → старые SQL-выражения нетронуты). + (d) чекпоинт per-region: _resume_dkp_cursor(source=...) изолирует регионы. + (e) реестр хендлеров резолвит rosreestr_dkp_import_77 через wildcard. +""" + +from __future__ import annotations + +import inspect +import os +from pathlib import Path +from typing import Any +from unittest.mock import MagicMock + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.services import scheduler as sched +from app.services.regions import REGIONS + +_IMPORT_SRC = inspect.getsource(sched.import_rosreestr_dkp) +_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql" +_MIGRATION_288 = _SQL_DIR / "288_deals_doc_type.sql" + + +# ── регион 66/77 в реестре (canonical_city) ────────────────────────────────── + + +def test_region_66_has_no_canonical_city_override() -> None: + """66 — источник несёт свой city как есть, поведение byte-for-byte прежнее.""" + assert REGIONS[66].canonical_city is None + + +def test_region_77_canonical_city_is_moskva() -> None: + """77 — Росреестр отдаёт округ/поселение вместо города, нужен override.""" + assert REGIONS[77].canonical_city == "Москва" + + +# ── _dkp_source_for_region — имя чекпоинта per-region ──────────────────────── + + +def test_dkp_source_for_region_66_keeps_legacy_name() -> None: + """66 — byte-for-byte прежнее имя, под которым годами писались scrape_runs.""" + assert sched._dkp_source_for_region(66) == "rosreestr_dkp_import" + + +def test_dkp_source_for_region_77_gets_suffix() -> None: + """77 — суффикс кода, тот же формат, что у строки scrape_schedules (миграция 288).""" + assert sched._dkp_source_for_region(77) == "rosreestr_dkp_import_77" + + +# ── валидация неизвестного региона — падает ДО любого обращения к БД ──────── + + +def test_import_rejects_unknown_region_code_before_touching_db() -> None: + """Неизвестный region_code — явный ValueError, БД не трогается вовсе.""" + db = MagicMock() + try: + sched.import_rosreestr_dkp(db, run_id=1, params={"region_code": 404}) + raised = False + except ValueError as exc: + raised = True + assert "404" in str(exc) + assert raised, "expected ValueError for unknown region_code" + db.execute.assert_not_called() + + +# ── SQL: region_code — bind-параметр, не литерал (см. также test_rosreestr_dedup_key) ─ + + +def test_sql_selects_region_code_via_bind_param() -> None: + assert "WHERE region_code = CAST(:region_code AS int)" in _IMPORT_SRC + + +def test_no_python_branch_on_region_77() -> None: + """Маппинг city/address/raw_payload идёт ОДНОЙ SQL-веткой на bind-параметре + :canonical_city (CASE WHEN), а не Python if/else на конкретный код региона.""" + assert "== 77" not in _IMPORT_SRC + assert "region_code == 77" not in _IMPORT_SRC + + +# ── SQL: маппинг city/address через canonical_city (регион 77) ────────────── + + +def test_sql_address_uses_canonical_city_case() -> None: + assert "CAST(:canonical_city AS text) || ', ' || trim(street)" in _IMPORT_SRC, ( + "address для canonical_city-региона обязан быть 'Москва, '" + ) + assert "trim(city) || ', ' || trim(street)" in _IMPORT_SRC, ( + "ELSE-ветка (регион 66) обязана остаться прежней" + ) + + +def test_sql_city_uses_canonical_city_case() -> None: + # WHEN CAST(:canonical_city AS text) IS NOT NULL THEN CAST(:canonical_city AS text) + assert "THEN CAST(:canonical_city AS text)\n" in _IMPORT_SRC + assert "ELSE trim(city)\n" in _IMPORT_SRC + + +def test_sql_city_not_null_filter_skipped_when_canonical_city_set() -> None: + """Фильтр city IS NOT NULL применяется ТОЛЬКО когда canonical_city не задан — + иначе на 77 теряется ~10% строк с пустым city источника.""" + assert ( + "CAST(:canonical_city AS text) IS NOT NULL\n" + " OR (city IS NOT NULL AND trim(city) <> '')" in _IMPORT_SRC + ) + + +def test_sql_raw_payload_carries_source_columns() -> None: + """raw_payload (только когда canonical_city задан) несёт src_city/okato/ + quarter_cad_number/district — исходные значения источника, не потерянные.""" + for key in ( + "'src_city', city", + "'okato', okato", + "'quarter_cad_number', quarter_cad_number", + "'district', district", + ): + assert key in _IMPORT_SRC, f"raw_payload missing key expression: {key!r}" + + +def test_params_dict_binds_region_code_and_canonical_city() -> None: + assert '"region_code": region_code' in _IMPORT_SRC + assert '"canonical_city": region.canonical_city' in _IMPORT_SRC + + +def test_insert_carries_doc_type_and_raw_payload() -> None: + assert "doc_type, raw_payload" in _IMPORT_SRC + assert "doc_type = EXCLUDED.doc_type" in _IMPORT_SRC + assert "raw_payload = EXCLUDED.raw_payload" in _IMPORT_SRC + + +# ── чекпоинт per-region: 77 не подхватывает last_id региона 66 ────────────── + + +class _FakeRow: + def __init__(self, **kw: Any) -> None: + self.__dict__.update(kw) + + +class _SourceKeyedFakeDb: + """Двойник сессии: отдаёт кандидата ТОЛЬКО для своего source, иначе None. + + Эмулирует реальный `WHERE source = :source` — единственный способ честно + проверить изоляцию чекпоинтов между регионами без живого Postgres. + """ + + def __init__(self, rows_by_source: dict[str, Any]) -> None: + self.rows_by_source = rows_by_source + self.seen_sources: list[str] = [] + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params is not None and "counters" in params: + return MagicMock() + source = (params or {}).get("source") + self.seen_sources.append(source) + return MagicMock(fetchone=lambda: self.rows_by_source.get(source)) + + def commit(self) -> None: + pass + + +def test_resume_cursor_region_77_does_not_see_region_66_checkpoint() -> None: + """Прогон source='rosreestr_dkp_import' (66) не резюмится под source региона 77.""" + row_66 = _FakeRow( + prev_id=777, + prev_status="zombie", + prev_counters={"last_id": 999999}, + age_h=1.0, + ) + db = _SourceKeyedFakeDb({"rosreestr_dkp_import": row_66}) + + last_id_77, verdict_77 = sched._resume_dkp_cursor( + db, run_id=1, source="rosreestr_dkp_import_77" + ) + assert last_id_77 == 0 + assert verdict_77["resume_reason"] == "no_prev_run" + + # Регион 66 по-прежнему видит СВОЙ чекпоинт. + last_id_66, verdict_66 = sched._resume_dkp_cursor(db, run_id=2, source="rosreestr_dkp_import") + assert last_id_66 == 999999 + assert verdict_66["resume_reason"] == "ok" + + +def test_resume_cursor_default_source_is_backward_compatible() -> None: + """Вызов без явного source (как в старых тестах/коде) — прежнее поведение (66).""" + row = _FakeRow(prev_id=1, prev_status="zombie", prev_counters={"last_id": 42}, age_h=1.0) + db = _SourceKeyedFakeDb({"rosreestr_dkp_import": row}) + last_id, _verdict = sched._resume_dkp_cursor(db, run_id=1) + assert last_id == 42 + assert db.seen_sources == ["rosreestr_dkp_import"] + + +# ── реестр хендлеров: rosreestr_dkp_import_77 резолвится через wildcard ────── + + +def test_handler_registry_region_77_uses_same_job_as_region_66() -> None: + from scraper_kit.orchestration.scheduler import build_registry, resolve_handler + + from app.services.product_handlers import _job_rosreestr_dkp, build_product_handlers + + registry = build_registry(build_product_handlers(ctx=None)) # type: ignore[arg-type] + h66 = resolve_handler("rosreestr_dkp_import", registry) + h77 = resolve_handler("rosreestr_dkp_import_77", registry) + assert h66 is not None and h77 is not None + assert h66.job is _job_rosreestr_dkp + assert h77.job is _job_rosreestr_dkp + + +# ── миграция 288 ────────────────────────────────────────────────────────────── + + +def test_migration_288_exists() -> None: + assert _MIGRATION_288.is_file(), f"missing migration: {_MIGRATION_288}" + + +def test_migration_288_adds_doc_type_idempotently() -> None: + sql = _MIGRATION_288.read_text("utf-8") + assert "ADD COLUMN IF NOT EXISTS doc_type text" in sql + assert "SET doc_type = 'ДКП'" in sql + assert "WHERE source = 'rosreestr'" in sql + assert "AND doc_type IS NULL" in sql + + +def test_migration_288_seeds_disabled_region_77_schedule() -> None: + sql = _MIGRATION_288.read_text("utf-8") + assert "'rosreestr_dkp_import_77'" in sql + assert "false" in sql + assert '"region_code": 77' in sql + assert "ON CONFLICT (source) DO NOTHING" in sql + + +def test_migration_288_has_lock_timeout_before_alter_table() -> None: + sql = _MIGRATION_288.read_text("utf-8") + lt_pos = sql.find("SET LOCAL lock_timeout") + alter_pos = sql.find("ALTER TABLE deals") + assert lt_pos != -1 and alter_pos != -1 + assert lt_pos < alter_pos diff --git a/tradein-mvp/backend/tests/test_rosreestr_dedup_key.py b/tradein-mvp/backend/tests/test_rosreestr_dedup_key.py index 90176d3c..b73ea88a 100644 --- a/tradein-mvp/backend/tests/test_rosreestr_dedup_key.py +++ b/tradein-mvp/backend/tests/test_rosreestr_dedup_key.py @@ -151,15 +151,21 @@ def test_migration_077_converts_md5_to_plain_key() -> None: def test_migration_077_shared_filter_matches_live_import() -> None: """Дедуп-релевантные фильтры src CTE миграции 077 совпадают с живым импортом. - Исключение — city ILIKE (см. test_live_import_dropped_ekb_city_filter ниже): - 077 backfill'ил ЕКБ-строки под EKB-only scope того времени; живой импорт расширен - на всю Свердловскую область (region_code=66, все города), поэтому city-фильтр из - живого импорта СНЯТ намеренно. Остальные клозы обязаны совпадать байт-в-байт, - иначе 077 конвертировал бы не тот набор строк. + Исключения: + - city ILIKE (см. test_live_import_dropped_ekb_city_filter ниже): 077 backfill'ил + ЕКБ-строки под EKB-only scope того времени; живой импорт расширен на всю + Свердловскую область (region_code=66, все города), поэтому city-фильтр из + живого импорта СНЯТ намеренно. + - region_code (#3051 п.3): 077 — историческая миграция (уже применена на прод, + хардкод region_code=66 трогать нельзя и не нужно — она навсегда про 66). Живой + импорт с #3051 параметризован по региону (region_code = CAST(:region_code AS + int)) — литерала 66 в его SQL больше нет, см. + test_live_import_region_code_is_bind_param ниже. Остальные клозы обязаны + совпадать байт-в-байт, иначе 077 конвертировал бы не тот набор строк. """ sql = _MIGRATION_077.read_text("utf-8") + assert "region_code = 66" in sql, "migration 077 must keep its historical literal" for clause in ( - "region_code = 66", "realestate_type_code = '002001003000'", "area BETWEEN 18 AND 200", "deal_price BETWEEN 1000000 AND 100000000", @@ -170,6 +176,18 @@ def test_migration_077_shared_filter_matches_live_import() -> None: assert clause in _IMPORT_SRC, f"missing filter clause in import: {clause!r}" +def test_live_import_region_code_is_bind_param() -> None: + """#3051 п.3: живой импорт параметризован по региону, литерала 66 в SQL нет. + + До #3051 живой импорт хардкодил `WHERE region_code = 66` — единственный регион + покрытия. С параметризацией region_code приходит из params (default 66 — обратная + совместимость), в SQL идёт bind-параметром через CAST, не литералом. + """ + assert "region_code = CAST(:region_code AS int)" in _IMPORT_SRC + assert "region_code = 66" not in _IMPORT_BODY + assert 'params.get("region_code", 66)' in _IMPORT_SRC + + def test_live_import_dropped_ekb_city_filter() -> None: """Живой импорт БОЛЬШЕ не фильтрует по городу — Mera расширена на всю обл. 66. diff --git a/tradein-mvp/backend/tests/test_scraper_kit_scheduler_parity.py b/tradein-mvp/backend/tests/test_scraper_kit_scheduler_parity.py index e651265d..efe3585a 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_scheduler_parity.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_scheduler_parity.py @@ -57,6 +57,10 @@ from scraper_kit.orchestration.scheduler import ( _PRODUCT_SOURCES: set[str] = { "cian_history_backfill", "rosreestr_dkp_import", + # #3051 п.3: member-source семейства "rosreestr_dkp_import_*" (Москва, миграция 288 + # — то же раскрытие wildcard в конкретный member, что deactivate_stale_avito/yandex/ + # cian ниже, а не сам wildcard "rosreestr_dkp_import_*"). + "rosreestr_dkp_import_77", "listing_source_snapshot", "asking_to_sold_ratio_refresh", "deal_city_price_bands_refresh", diff --git a/tradein-mvp/deploy/import-rosreestr.sh b/tradein-mvp/deploy/import-rosreestr.sh index 460b076f..be828e8f 100755 --- a/tradein-mvp/deploy/import-rosreestr.sh +++ b/tradein-mvp/deploy/import-rosreestr.sh @@ -2,12 +2,22 @@ # Импорт реальных сделок Росреестра из gendesign-БД в tradein deals. # # Источник: gendesign-postgres-1 / rosreestr_deals (6.8М строк, партиц.). -# Берём ЕКБ квартиры (realestate_type_code=002001003000) за последние ~18 мес. -# rooms выводим из площади (Росреестр не отдаёт кол-во комнат). +# Берём квартиры региона REGION_CODE (realestate_type_code=002001003000) за +# последние ~18 мес. rooms выводим из площади (Росреестр не отдаёт кол-во комнат). # Координаты — NULL, проставляются отдельно (geocode по street). # # Только ДКП (вторичка) — ДДУ застройщиков скёюят median, отделены сознательно. См. PR-A. # +# #3051 п.3: REGION_CODE параметризован (default 66 — Свердловская обл., byte-for-byte +# прежнее поведение). Этот bash-путь — НЕ region-generic: для region_code=77 (Москва) он +# НЕ подставляет canonical_city вместо city источника (Росреестр по Москве отдаёт +# муниципальный округ/поселение, не сам город) и не пишет raw_payload с +# okato/quarter_cad_number/district — эту логику несёт только Python-путь +# (app/services/scheduler.py::import_rosreestr_dkp, боевой планировщик). Если этот +# скрипт когда-нибудь запустят вручную с REGION_CODE=77 — city/address будут +# "муниципальный округ Раменки, ..." как есть из источника, НЕ "Москва, ...". Держать +# паритет фильтров (area/price/doc_type) обязательно, паритет city-override — нет. +# # Запуск на прод-хосте: ./import-rosreestr.sh # Повторяемо: дедуп по dedup_hash, новый запуск подтянет свежие кварталы. @@ -18,8 +28,9 @@ DST_PG="${DST_PG:-tradein-postgres}" SRC_DB="${SRC_DB:-gendesign}" SRC_USER="${SRC_USER:-gendesign}" SINCE="${SINCE:-2024-01-01}" +REGION_CODE="${REGION_CODE:-66}" -echo "[$(date -u +%H:%M:%S)] import-rosreestr: ЕКБ квартиры с $SINCE" +echo "[$(date -u +%H:%M:%S)] import-rosreestr: регион $REGION_CODE, квартиры с $SINCE" # 1. Staging-таблица в tradein. docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c " @@ -53,7 +64,7 @@ docker exec "$SRC_PG" psql -U "$SRC_USER" -d "$SRC_DB" -v ON_ERROR_STOP=on -c " round(price_per_sqm)::int AS price_per_m2, period_start_date AS deal_date FROM rosreestr_deals - WHERE region_code = 66 + WHERE region_code = $REGION_CODE AND city IS NOT NULL AND trim(city) <> '' AND realestate_type_code = '002001003000' AND area BETWEEN 18 AND 200