From 5f9dc5d5123a02c32bc27004f487fc401969be5b Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 30 Aug 2026 15:34:36 +0300 Subject: [PATCH] =?UTF-8?q?feat(tradein/cian):=20=D1=83=20=D0=A6=D0=B8?= =?UTF-8?q?=D0=B0=D0=BD=D0=B0=20=D0=BD=D0=B5=20=D0=B1=D1=8B=D0=BB=D0=BE=20?= =?UTF-8?q?=D0=B4=D0=BE=D0=B1=D0=BE=D1=80=D0=B0=20=D0=BA=D0=B0=D1=80=D1=82?= =?UTF-8?q?=D0=BE=D1=87=D0=B5=D0=BA=20=E2=80=94=20=D1=82=D0=BE=D0=BB=D1=8C?= =?UTF-8?q?=D0=BA=D0=BE=20=D0=BF=D0=BE=D0=B1=D0=BE=D1=87=D0=BD=D1=8B=D0=B9?= =?UTF-8?q?=20=D1=8D=D1=84=D1=84=D0=B5=D0=BA=D1=82=20=D0=B7=D0=B0=D0=B4?= =?UTF-8?q?=D0=B0=D1=87=D0=B8=20=D0=BF=D1=80=D0=BE=20=D0=B8=D1=81=D1=82?= =?UTF-8?q?=D0=BE=D1=80=D0=B8=D1=8E=20(#3284)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Карточки Циана доставались побочным эффектом cian_history_backfill, а её выборка ключуется по offer_price_history. Следствия на 30.08: карточка есть у 4494 из 25222 объявлений (17.8% — последнее место при втором месте по объёму), 19046 без истории при квоте 100/сутки (190 дней на остаток, тогда как очередь растёт вдвадцатеро быстрее), и 1697 объявлений с историей и без карточки, которые исторической выборке недостижимы в принципе. Фетчер при этом исправен: прогоны 5154/5240/5328 дали 100/100, 99/100, 100/100. Чинить нечего — не выдана мощность. Добавлен второй режим выборки (listings_pending="detail", по detail_enriched_at, свежие первыми) и второе расписание поверх ТОГО ЖЕ тела: машинерия работает, дублировать её новым модулем незачем. Историческая выборка оставлена побайтово — по ней живёт суточный прогон. batch_size=400 не на глаз: замеренный темп ~28с на объявление, порог reap_zombies 6ч по heartbeat, бюджетного сторожа у задачи нет — 400×28с≈3.1ч проходит, 800 как у Яндекса (≈6.2ч) убивало бы жнецом. Расписание засеяно enabled=false, как domclick_detail_backfill в миграции 175: это третий круглосуточный добор на общий пул из четырёх узлов, влияние на соседей надо посмотреть, а не предположить. Тесты: 13 проверок, ключ выборки / порядок / неизменность прежнего режима / проводка параметров через посредника / регистрация обоих source. Проверено мутациями: снятие ORDER BY, молчаливый дефолт вместо ValueError и потеря listings_pending в посреднике роняют по 2-3 теста каждая. Набор целиком — 5196 passed, 37 skipped. --- .../backend/app/services/product_handlers.py | 9 + tradein-mvp/backend/app/services/scheduler.py | 11 +- .../app/tasks/cian_history_backfill.py | 67 +++-- ...pe_schedules_seed_cian_detail_backfill.sql | 70 ++++++ .../test_3284_cian_detail_backfill_queue.py | 235 ++++++++++++++++++ 5 files changed, 374 insertions(+), 18 deletions(-) create mode 100644 tradein-mvp/backend/data/sql/282_scrape_schedules_seed_cian_detail_backfill.sql create mode 100644 tradein-mvp/backend/tests/test_3284_cian_detail_backfill_queue.py diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index b7a74b60..fb879a3f 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -706,6 +706,15 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]: "cian_history_backfill", pre_claim=_cian_pre_claim, ), + # #3284: то же тело, другая очередь. Разные source нужны именно как РАЗНЫЕ + # строки расписания — у них свои окна, свой next_run_at и свой счётчик + # прогонов; гонять оба режима под одним source нельзя, планировщик держит + # на source ровно один активный прогон. + "cian_detail_backfill": Handler( + _job_cian_history_backfill, + "cian_detail_backfill", + pre_claim=_cian_pre_claim, + ), "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( diff --git a/tradein-mvp/backend/app/services/scheduler.py b/tradein-mvp/backend/app/services/scheduler.py index 523dbf6c..b5cc645c 100644 --- a/tradein-mvp/backend/app/services/scheduler.py +++ b/tradein-mvp/backend/app/services/scheduler.py @@ -99,10 +99,18 @@ async def _execute_cian_backfill( Params (from default_params jsonb): batch_size: int — rows per run (listings + houses counted separately). + listings_pending: str — "history" (дефолт) | "detail", см. #3284. + do_houses: bool — дефолт true; у cian_detail_backfill выключен. """ from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history batch_size = int(params.get("batch_size", 100)) + # #3284: одно тело обслуживает ДВА расписания. cian_history_backfill идёт с + # дефолтами (история + дома), cian_detail_backfill — с listings_pending="detail" + # и do_houses=false: дома у него уже разбирает суточный сосед, а гонять их + # круглосуточно незачем. + listings_pending = str(params.get("listings_pending", "history")) + do_houses = bool(params.get("do_houses", True)) def _counters(result: CianBackfillResult) -> dict[str, int]: return { @@ -143,9 +151,10 @@ async def _execute_cian_backfill( db, batch_size=batch_size, do_listings=True, - do_houses=True, + do_houses=do_houses, do_valuations=False, on_progress=_heartbeat, + listings_pending=listings_pending, ) counters = {**_counters(result), "duration_sec": int(result.duration_sec)} diff --git a/tradein-mvp/backend/app/tasks/cian_history_backfill.py b/tradein-mvp/backend/app/tasks/cian_history_backfill.py index 3715f0a5..94762ae0 100644 --- a/tradein-mvp/backend/app/tasks/cian_history_backfill.py +++ b/tradein-mvp/backend/app/tasks/cian_history_backfill.py @@ -110,6 +110,41 @@ def _note_refusal(result: CianBackfillResult, status: int | None) -> str | None: return kind +# ── Выборки «что ещё не добрано» ───────────────────────────────────────────── +# Две РАЗНЫЕ цели, которые до #3284 были склеены в одну. Историческая выборка +# ключуется по offer_price_history, поэтому объявление, у которого история уже +# есть, а карточки нет, не вернётся ей НИКОГДА (на 30.08 таких 1697). Вторая +# выборка закрывает ровно этот пробел и заодно даёт Циану то, что у avito / +# domclick / yandex есть давно, — добор по признаку «нет карточки». +_LISTINGS_PENDING_SQL: dict[str, str] = { + # Прежнее поведение, побайтово. Менять его правкой про добор карточек нельзя: + # по этой выборке живёт суточный cian_history_backfill. + "history": """ + SELECT l.id, l.source_url + FROM listings l + LEFT JOIN offer_price_history oph ON oph.listing_id = l.id + WHERE l.source = 'cian' + AND l.source_url IS NOT NULL + AND oph.listing_id IS NULL + LIMIT :lim + """, + # #3284. ORDER BY здесь есть, а в "history" нет, и это намеренно: очередь + # карточек (на 30.08 — 20 728 объявлений) заведомо длиннее любого батча, и + # порядок решает, что мы успеем добрать. Свежие важнее: по ним считается + # оценка. У "history" очередь того же порядка, но её сортировку трогать — + # отдельное решение с отдельной проверкой, не побочный эффект этой правки. + "detail": """ + SELECT l.id, l.source_url + FROM listings l + WHERE l.source = 'cian' + AND l.source_url IS NOT NULL + AND l.detail_enriched_at IS NULL + ORDER BY l.last_seen_at DESC NULLS LAST + LIMIT :lim + """, +} + + async def backfill_cian_history( db: Session, *, @@ -119,6 +154,7 @@ async def backfill_cian_history( do_valuations: bool = False, dry_run: bool = False, on_progress: Callable[[CianBackfillResult], None] | None = None, + listings_pending: str = "history", ) -> CianBackfillResult: """Iterate Cian listings + houses with missing history, fetch+save. @@ -133,6 +169,11 @@ async def backfill_cian_history( do_valuations: process Cian Valuation Calculator batch (external_valuations backfill). Default False — opt-in because each call hits Cian auth-gated API. dry_run: skip all fetch+save; only count and log pending rows. + listings_pending: какую очередь разбирает блок listings (#3284) -- + "history" (дефолт, прежнее поведение): нет строки в offer_price_history; + "detail": нет карточки (detail_enriched_at IS NULL), свежие первыми. + Неизвестное значение -- ValueError, а не молчаливый дефолт: пустой + батч из-за опечатки в default_params выглядел бы как «всё добрано». on_progress: колбэк живости (#2725) — вызывается на каждой сущности ЛЮБОГО из трёх блоков, до её обработки, с текущим (мутируемым) result. Caller пишет heartbeat; исключения колбэка — на его совести (планировщик глушит их сам), @@ -145,24 +186,16 @@ async def backfill_cian_history( t0 = time.time() delay = get_scraper_delay("cian") # seconds; default 5.0 - # ── 1. Listings: missing offer_price_history ────────────────────────────── + # ── 1. Listings: pending-выборка (см. _LISTINGS_PENDING_SQL) ────────────── if do_listings: - rows = ( - db.execute( - text(""" - SELECT l.id, l.source_url - FROM listings l - LEFT JOIN offer_price_history oph ON oph.listing_id = l.id - WHERE l.source = 'cian' - AND l.source_url IS NOT NULL - AND oph.listing_id IS NULL - LIMIT :lim - """), - {"lim": batch_size}, - ) - .mappings() - .all() - ) + try: + pending_sql = _LISTINGS_PENDING_SQL[listings_pending] + except KeyError: + raise ValueError( + f"listings_pending={listings_pending!r} неизвестен; " + f"допустимы {sorted(_LISTINGS_PENDING_SQL)}" + ) from None + rows = db.execute(text(pending_sql), {"lim": batch_size}).mappings().all() result.listings_total = len(rows) if dry_run: diff --git a/tradein-mvp/backend/data/sql/282_scrape_schedules_seed_cian_detail_backfill.sql b/tradein-mvp/backend/data/sql/282_scrape_schedules_seed_cian_detail_backfill.sql new file mode 100644 index 00000000..28845bd2 --- /dev/null +++ b/tradein-mvp/backend/data/sql/282_scrape_schedules_seed_cian_detail_backfill.sql @@ -0,0 +1,70 @@ +-- 282_scrape_schedules_seed_cian_detail_backfill.sql +-- Отдельный добор карточек Циана (issue #3284). +-- +-- ЧТО БЫЛО НЕ ТАК +-- У avito, domclick и yandex есть выделенный *_detail_backfill. У Циана его не было +-- вообще: карточки доставались побочным эффектом cian_history_backfill — суточной +-- задачи про историю цен, — и её выборка ключуется по offer_price_history, а не по +-- наличию карточки. Следствия на 30.08.2026: +-- * 25 222 объявления Циана, карточка есть у 4 494 (17.8%) — последнее место среди +-- четырёх источников при втором месте по объёму (domklik 76%, yandex 50%, +-- avito 24%); +-- * 19 046 объявлений без истории цен при квоте batch_size=100 в сутки — это 190 +-- дней на текущий остаток, притом что Циан приносит ~4 900 объявлений за двое +-- суток, то есть очередь растёт примерно в двадцать раз быстрее, чем разбирается; +-- * 1 697 объявлений имеют историю цен и НЕ имеют карточки — исторической выборке +-- они уже «обработаны» и не вернутся к ней никогда. +-- При этом сам фетчер исправен: прогоны 5154/5240/5328 дали 100/100, 99/100, 100/100. +-- Чинить нечего — не выдана мощность. +-- +-- ЧТО ДЕЛАЕТ ЭТА СТРОКА +-- Второе расписание поверх ТОГО ЖЕ тела (app.services.scheduler._execute_cian_backfill), +-- отличающееся двумя параметрами: +-- listings_pending = "detail" — выборка по detail_enriched_at IS NULL, свежие первыми +-- (ORDER BY last_seen_at DESC): очередь заведомо длиннее батча, и +-- порядок решает, что успеем добрать; по свежим считается оценка. +-- do_houses = false — дома разбирает суточный cian_history_backfill, дублировать их +-- круглосуточно незачем. +-- +-- ОКНО 0–23 +-- Ровно как у avito_detail_backfill и yandex_detail_backfill. Трёхчасовое окно, как у +-- истории, здесь бессмысленно: прогон 5328 уложился в 54 минуты из трёх часов, то есть +-- две трети окна простаивали, а очередь при этом росла. +-- +-- BATCH_SIZE 400 — арифметика, а не вкус +-- Замеренный темп ~28 с на объявление (задержка 5 с, остальное сеть). 400 × 28 с ≈ 3.1 ч. +-- Порог reap_zombies — 6 ч по heartbeat_at, и прогон обязан заканчиваться заметно раньше: +-- у backfill_cian_history НЕТ бюджетного сторожа, единственный ограничитель — batch_size. +-- Отсюда же не 800, как у yandex: 800 × 28 с ≈ 6.2 ч — прогон убивало бы зомби-жнецом. +-- +-- ENABLED=false +-- Как и у domclick_detail_backfill (миграция 175): включение — отдельный осознанный шаг +-- после деплоя и дымовой пробы. Причина не в куках (гейт _cian_pre_claim общий с историей +-- и проходит), а в нагрузке: это ТРЕТИЙ круглосуточный добор на общий пул из четырёх +-- узлов, и его влияние на соседей надо посмотреть, а не предположить. +-- +-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)). +-- Идемпотентно: ON CONFLICT (source) DO NOTHING. + +BEGIN; + +INSERT INTO scrape_schedules ( + source, + enabled, + window_start_hour, + window_end_hour, + next_run_at, + default_params +) +VALUES +( + 'cian_detail_backfill', + false, + 0, + 23, + ((CURRENT_DATE + INTERVAL '1 day')) AT TIME ZONE 'UTC', + '{"batch_size": 400, "listings_pending": "detail", "do_houses": false}'::jsonb +) +ON CONFLICT (source) DO NOTHING; + +COMMIT; diff --git a/tradein-mvp/backend/tests/test_3284_cian_detail_backfill_queue.py b/tradein-mvp/backend/tests/test_3284_cian_detail_backfill_queue.py new file mode 100644 index 00000000..8a5e27d4 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3284_cian_detail_backfill_queue.py @@ -0,0 +1,235 @@ +"""#3284: у Циана не было добора карточек — только побочный эффект задачи про историю цен. + +Историческая выборка ключуется по offer_price_history. Объявление, у которого история +уже есть, а карточки нет, для неё «обработано» и не вернётся никогда — на 30.08.2026 +таких 1697. Плюс квота 100/сутки против очереди в 19 046 — это 190 дней, при том что +Циан приносит ~4900 объявлений за двое суток. + +Правка добавляет второй режим выборки (`listings_pending="detail"`) и второе расписание +поверх того же тела. Тесты ниже закрепляют ровно то, что делает режим полезным: +выборку по detail_enriched_at, свежие первыми, неизменность прежнего режима и отказ +на опечатке в default_params (иначе пустой батч читался бы как «всё добрано»). +""" + +from __future__ import annotations + +import os +import re +from types import SimpleNamespace +from typing import Any +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from app.tasks import cian_history_backfill +from app.tasks.cian_history_backfill import _LISTINGS_PENDING_SQL + + +def _norm(sql: str) -> str: + """Схлопнуть пробелы — сравниваем смысл запроса, а не его отступы.""" + return re.sub(r"\s+", " ", sql).strip() + + +class _FakeBrowserFetcher: + def __init__(self, **kwargs: Any) -> None: + self.last_response_status = None + + async def __aenter__(self) -> _FakeBrowserFetcher: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + +def _enrichment() -> SimpleNamespace: + return SimpleNamespace(price_changes=[]) + + +def _db_returning(rows: list[dict[str, Any]]) -> MagicMock: + db = MagicMock() + db.execute.return_value.mappings.return_value.all.return_value = rows + return db + + +async def _run(db: MagicMock, **kwargs: Any) -> Any: + with ( + patch.object(cian_history_backfill, "BrowserFetcher", _FakeBrowserFetcher), + patch.object(cian_history_backfill, "fetch_detail", AsyncMock(return_value=_enrichment())), + patch.object(cian_history_backfill, "save_detail_enrichment", MagicMock()), + patch.object(cian_history_backfill, "RealMatcherAdapter", MagicMock()), + patch("asyncio.sleep", new_callable=AsyncMock), + ): + return await cian_history_backfill.backfill_cian_history( + db, do_houses=False, do_valuations=False, **kwargs + ) + + +# ── выборка "detail" ───────────────────────────────────────────────────────── + + +def test_detail_queue_selects_by_missing_card_not_by_history() -> None: + """Ключ выборки — detail_enriched_at, а не offer_price_history. + + Это вся суть #3284: пока признаком служит история, 1697 объявлений с историей + и без карточки недостижимы. + """ + sql = _norm(_LISTINGS_PENDING_SQL["detail"]) + assert "l.detail_enriched_at IS NULL" in sql + assert "offer_price_history" not in sql + + +def test_detail_queue_takes_freshest_first() -> None: + """Очередь длиннее батча, поэтому порядок решает, что мы успеем добрать.""" + sql = _norm(_LISTINGS_PENDING_SQL["detail"]) + assert "ORDER BY l.last_seen_at DESC NULLS LAST" in sql + assert sql.index("ORDER BY") < sql.index("LIMIT") + + +def test_detail_queue_stays_scoped_to_cian() -> None: + """Источник обязан быть в WHERE: иначе добор Циана заберёт чужие объявления.""" + assert "l.source = 'cian'" in _norm(_LISTINGS_PENDING_SQL["detail"]) + assert "l.source_url IS NOT NULL" in _norm(_LISTINGS_PENDING_SQL["detail"]) + + +# ── прежний режим не тронут ────────────────────────────────────────────────── + + +def test_history_queue_unchanged() -> None: + """По этой выборке живёт суточный cian_history_backfill — она обязана остаться прежней.""" + sql = _norm(_LISTINGS_PENDING_SQL["history"]) + assert "LEFT JOIN offer_price_history oph ON oph.listing_id = l.id" in sql + assert "oph.listing_id IS NULL" in sql + assert "detail_enriched_at" not in sql + # Сортировки в исторической выборке не было и не появилось: её добавление — + # отдельное решение с отдельной проверкой, а не побочный эффект #3284. + assert "ORDER BY" not in sql + + +def test_default_mode_is_history() -> None: + """Дефолт обязан оставаться прежним: расписание истории параметра не передаёт.""" + db = _db_returning([]) + import asyncio + + asyncio.run(_run(db)) + used = _norm(str(db.execute.call_args_list[0].args[0])) + assert "offer_price_history" in used + + +# ── режим доезжает до запроса ──────────────────────────────────────────────── + + +async def test_detail_mode_reaches_the_query() -> None: + """Параметр не должен потеряться по дороге — проверяем сам исполненный SQL.""" + db = _db_returning([]) + await _run(db, listings_pending="detail") + used = _norm(str(db.execute.call_args_list[0].args[0])) + assert "detail_enriched_at IS NULL" in used + assert "offer_price_history" not in used + + +async def test_detail_mode_processes_rows_normally() -> None: + """Смена выборки не меняет обработку: строки те же, счётчики те же.""" + db = _db_returning([{"id": 7, "source_url": "https://cian.ru/7"}]) + result = await _run(db, listings_pending="detail", batch_size=10) + assert result.listings_processed == 1 + assert result.listings_succeeded == 1 + + +async def test_batch_size_reaches_the_query() -> None: + """batch_size обязан доезжать как :lim — иначе квота из расписания ничего не значит.""" + db = _db_returning([]) + await _run(db, listings_pending="detail", batch_size=400) + assert db.execute.call_args_list[0].args[1] == {"lim": 400} + + +# ── опечатка не должна читаться как «всё добрано» ──────────────────────────── + + +async def test_unknown_mode_raises_instead_of_silently_defaulting() -> None: + """Опечатка в default_params обязана падать громко. + + Молчаливый откат на "history" дал бы прогон с нулём добранных карточек, который + выглядит как штатный: очередь якобы пуста. Такую ошибку ищут днями. + """ + db = _db_returning([]) + with pytest.raises(ValueError, match="listings_pending"): + await _run(db, listings_pending="detali") + + +async def test_unknown_mode_names_the_allowed_values() -> None: + """Сообщение должно называть допустимые значения — иначе оно не помогает.""" + db = _db_returning([]) + with pytest.raises(ValueError) as e: + await _run(db, listings_pending="") + assert "detail" in str(e.value) and "history" in str(e.value) + + +# ── параметры доезжают от расписания до задачи ─────────────────────────────── +# Три предыдущих теста проверяют саму задачу. Эти — путь от строки расписания: +# default_params → _execute_cian_backfill → backfill_cian_history. Без них правка +# в теле-посреднике молча вернула бы cian_detail_backfill к разбору истории. + + +async def test_scheduler_passes_detail_mode_from_params() -> None: + """listings_pending из default_params обязан доехать до задачи.""" + from app.services import scheduler as scheduler_mod + + captured: dict[str, Any] = {} + + async def _fake_backfill(db: Any, **kwargs: Any) -> Any: + captured.update(kwargs) + # Настоящий результат, а не SimpleNamespace: посредник читает у него + # поля, которых в самодельной заглушке легко недосчитаться. + return cian_history_backfill.CianBackfillResult() + + with ( + patch("app.tasks.cian_history_backfill.backfill_cian_history", _fake_backfill), + patch.object(scheduler_mod.runs_mod, "update_heartbeat", MagicMock()), + patch.object(scheduler_mod.runs_mod, "mark_done", MagicMock()), + ): + await scheduler_mod._execute_cian_backfill( + MagicMock(), + run_id=1, + params={"batch_size": 400, "listings_pending": "detail", "do_houses": False}, + ) + + assert captured["listings_pending"] == "detail" + assert captured["do_houses"] is False + assert captured["batch_size"] == 400 + + +async def test_scheduler_defaults_stay_history_and_houses() -> None: + """Без параметров поведение прежнее — по нему живёт суточный cian_history_backfill.""" + from app.services import scheduler as scheduler_mod + + captured: dict[str, Any] = {} + + async def _fake_backfill(db: Any, **kwargs: Any) -> Any: + captured.update(kwargs) + # Настоящий результат, а не SimpleNamespace: посредник читает у него + # поля, которых в самодельной заглушке легко недосчитаться. + return cian_history_backfill.CianBackfillResult() + + with ( + patch("app.tasks.cian_history_backfill.backfill_cian_history", _fake_backfill), + patch.object(scheduler_mod.runs_mod, "update_heartbeat", MagicMock()), + patch.object(scheduler_mod.runs_mod, "mark_done", MagicMock()), + ): + await scheduler_mod._execute_cian_backfill(MagicMock(), run_id=1, params={}) + + assert captured["listings_pending"] == "history" + assert captured["do_houses"] is True + + +def test_both_cian_sources_are_registered() -> None: + """Оба source обязаны быть в реестре: планировщик держит на source один прогон, + поэтому два режима не могут делить одно имя.""" + from app.services.product_handlers import build_product_handlers + + handlers = build_product_handlers(MagicMock()) + assert "cian_history_backfill" in handlers + assert "cian_detail_backfill" in handlers + # Гейт кук общий: он читает schedule_row["source"], а не хардкодит имя. + assert handlers["cian_detail_backfill"].pre_claim is handlers["cian_history_backfill"].pre_claim