From 8ffc2cb9422daff54d5e2ff4a8bbd4bbee2cf3a9 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 11:45:04 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein/scraper):=20=D0=B4=D0=BD=D0=B5?= =?UTF-8?q?=D0=B2=D0=BD=D0=BE=D0=B9=20=D1=81=D0=BD=D0=B8=D0=BC=D0=BE=D0=BA?= =?UTF-8?q?=20=D1=83=D0=B7=D0=BD=D0=B0=D1=91=D1=82=20=D1=81=D0=B2=D0=BE?= =?UTF-8?q?=D0=B9=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=20(#2701)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Из девяти вызовов save_listings в pipeline.py три не передавали run_id. У двух — run_avito_city_sweep:1113 и run_avito_newbuilding_sweep:1754 — run_id лежал в той же функции (keyword-параметр, под которым идёт runs.mark_done/heartbeat/is_cancelled) и просто не доезжал до вызова. Передан. run_avito_pipeline:570 — оставлен без run_id осознанно, решение записано в код: у функции нет строки в scrape_runs, продуктовых вызывающих нет (единственные — тесты), привязывать снимок не к чему. Писатели вне pipeline.py проверены и тоже помечены: admin /scrape (ручной скрейп без прогона) и scripts/ingest_domclick_jsonl.py (разовый ingest файла). Второй и есть главная дыра замера — 129 908 снимков domklik без run_id, а вовсе не Авито: по источникам заполнение yandex 100%, cian 98.8%, avito 80.3%, domklik 2.9%. Дыра domklik историческая (ежедневный run_domclick_city_sweep run_id передаёт), а вот у avito пустые строки приходят КАЖДЫЙ день (944 за 2026-08-06, 635 за 2026-08-05 — все без run_id). Бэкфилл 150 515 исторических строк невозможен, проверено: - привязки снимок→прогон в схеме больше нет; - по времени неоднозначно: на всех 36 днях с осиротевшими avito-снимками в сутках было >1 пишущего avito-прогона (до 29); - observed_at не спасает — ON CONFLICT перезаписывает его последним писателем дня, а не тем, кто строку создал. Ложная атрибуция хуже пустоты — бэкфилл не делается. Refs #2701 --- tradein-mvp/backend/app/api/v1/admin.py | 2 ++ .../backend/scripts/ingest_domclick_jsonl.py | 4 +++ .../tests/test_scraper_kit_pipeline_parity.py | 30 +++++++++++++++++++ .../test_scraper_kit_pipeline_parity2.py | 21 +++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 9 ++++++ 5 files changed, 66 insertions(+) diff --git a/tradein-mvp/backend/app/api/v1/admin.py b/tradein-mvp/backend/app/api/v1/admin.py index 6e41c279..50635e45 100644 --- a/tradein-mvp/backend/app/api/v1/admin.py +++ b/tradein-mvp/backend/app/api/v1/admin.py @@ -232,6 +232,8 @@ async def scrape_around( ) else: lots = await scraper.fetch_around(payload.lat, payload.lon, payload.radius_m) + # run_id нет и не будет (#2701): ручной admin-скрейп строки в scrape_runs не + # заводит — снимок пишется вне прогона, поле честно остаётся NULL. inserted, updated = save_listings( db, lots, matcher=matcher, region_code=DEFAULT_REGION_CODE ) diff --git a/tradein-mvp/backend/scripts/ingest_domclick_jsonl.py b/tradein-mvp/backend/scripts/ingest_domclick_jsonl.py index ff0b65ce..08b6b70a 100644 --- a/tradein-mvp/backend/scripts/ingest_domclick_jsonl.py +++ b/tradein-mvp/backend/scripts/ingest_domclick_jsonl.py @@ -190,6 +190,10 @@ def run(jsonl_path: str, limit: int | None, dry_run: bool) -> dict[str, int]: db = SessionLocal() try: if lots: + # run_id нет и не будет (#2701): разовый ingest файла — не прогон скрапера, + # строки в scrape_runs под него не существует. Именно отсюда 129 908 снимков + # domklik без run_id (2.9% заполнения у источника) — исторические, не текущие: + # ежедневный run_domclick_city_sweep run_id передаёт. inserted, updated = save_listings( db, lots, matcher=RealMatcherAdapter(), region_code=DEFAULT_REGION_CODE ) diff --git a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity.py b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity.py index 18b796cb..7da8983b 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity.py @@ -442,3 +442,33 @@ async def test_avito_city_sweep_passes_proxy_provider_to_scraper_constructor() - avito_scraper_cls = capture["avito_scraper_cls"] avito_scraper_cls.assert_called_once() assert avito_scraper_cls.call_args.kwargs.get("proxy_provider") is sentinel + + +# ── #2701: снимок обязан знать свой прогон ──────────────────────────────────── +# +# Замер на проде до правки: listings_snapshots 396 162 строки, run_id заполнен у +# 245 647 (62.0%); по avito 80.3%, и у ВСЕХ дневных city-sweep строк (напр. 944 за +# 2026-08-06, 635 за 2026-08-05) run_id пуст — sweep его просто не передавал, хотя +# держал в своей же сигнатуре и логировал в каждой строке. +# +# Вторая половина цепочки (save_listings прокидывает run_id в upsert_listing_snapshot) +# уже под замком: test_snapshot_writer.py::test_save_listings_snapshot_receives_run_id. + + +@pytest.mark.asyncio +async def test_avito_city_sweep_passes_run_id_to_save_listings() -> None: + """run_avito_city_sweep(run_id=1) → save_listings(..., run_id=1) → снимок с прогоном. + + Falsification: убрать `run_id=run_id` из вызова save_listings в pipeline.py — + kwargs не содержит run_id, assert падает. + """ + scenario = _Scenario( + anchors=_ANCHORS_2, + per_anchor=[("lots", 3, 3, 0), ("lots", 2, 2, 0)], + ) + capture: dict[str, Any] = {} + await _drive(scenario, capture=capture) + save_mock = capture["save_mock"] + assert save_mock.call_count == 2, "оба anchor'а сохраняют — проверяем оба вызова" + for call in save_mock.call_args_list: + assert call.kwargs.get("run_id") == 1, "снимок anchor'а остался бы без прогона" diff --git a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py index c38ffc0d..248ebe81 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py @@ -865,3 +865,24 @@ async def test_full_load_honest_empty_stays_done(source: str) -> None: counters, calls = await _drive_full_load_empty(source=source, attempts=6, failures=0) assert counters["unique_fetched"] == 0 assert calls[-1][0] == "mark_done" + + +# ── #2701: снимок обязан знать свой прогон ──────────────────────────────────── +# +# Замер на проде до правки: run_id пуст у 150 515 из 396 162 снимков (38%). +# Из девяти вызовов save_listings в pipeline.py шесть передавали run_id, три нет; +# у двух из трёх (city sweep + этот novostroyka-обход) run_id лежал в той же функции. + + +@pytest.mark.asyncio +async def test_avito_newbuilding_sweep_passes_run_id_to_save_listings() -> None: + """run_avito_newbuilding_sweep(run_id=1) → save_listings(..., run_id=1). + + Falsification: убрать `run_id=run_id` из вызова save_listings — kwargs пуст, assert падает. + """ + capture: dict[str, Any] = {} + await _drive_nb_sweep(capture=capture) + save_mock = capture["save_mock"] + assert save_mock.call_count > 0 + for call in save_mock.call_args_list: + assert call.kwargs.get("run_id") == 1, "novostroyka-снимок остался бы без прогона" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 57a17fc6..c571f75d 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -567,6 +567,13 @@ async def run_avito_pipeline( # ── Step 2: save listings ─────────────────────────────── if lots: try: + # run_id НЕ передаётся намеренно (#2701): у этой функции нет строки в + # scrape_runs — она одиночный anchor-путь, продуктовых вызовов не имеет + # (единственные вызывающие — тесты; развёртки идут через + # run_avito_city_sweep / run_avito_newbuilding_sweep / run_avito_full_load, + # которые run_id заводят и передают). Привязывать снимок не к чему; + # снимок с run_id несуществующего прогона был бы хуже пустого поля. + # Появится вызывающий с прогоном — добавить параметр run_id, как у sweep'ов. counters.lots_inserted, counters.lots_updated = save_listings( db, lots, matcher=matcher, region_code=region_code ) @@ -1115,6 +1122,7 @@ async def run_avito_city_sweep( anchor_lots, matcher=matcher, region_code=region_code, + run_id=run_id, city=_city_name, city_anchor=_city_anchor_point, city_radius_km=_city_radius_km, @@ -1756,6 +1764,7 @@ async def run_avito_newbuilding_sweep( lots, matcher=matcher, region_code=region_code, + run_id=run_id, city=EKATERINBURG_CITY_NAME, ) counters.lots_inserted += ins