From 1ba1a557707c13892ff6c4b96190c135598599e1 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 09:44:23 +0000 Subject: [PATCH 1/3] =?UTF-8?q?fix(tradein/browser):=20=D0=B3=D0=BE=D0=BD?= =?UTF-8?q?=D0=BA=D0=B0=20=C2=ABexecution=20context=20destroyed=C2=BB=20?= =?UTF-8?q?=E2=80=94=20=D0=B2=D0=BE=D1=81=D1=81=D1=82=D0=B0=D0=BD=D0=BE?= =?UTF-8?q?=D0=B2=D0=B8=D0=BC=D0=B0=D1=8F,=20=D1=80=D0=B5=D1=82=D1=80?= =?UTF-8?q?=D0=B0=D0=B9=20=D0=B1=D0=B5=D0=B7=20relaunch=20(#2676)=20(#2716?= =?UTF-8?q?)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tradein-mvp/browser/server.py | 60 ++++++ tradein-mvp/browser/test_server_fetch_json.py | 176 ++++++++++++++++++ 2 files changed, 236 insertions(+) diff --git a/tradein-mvp/browser/server.py b/tradein-mvp/browser/server.py index 46c9deb0..b2836430 100644 --- a/tradein-mvp/browser/server.py +++ b/tradein-mvp/browser/server.py @@ -125,6 +125,11 @@ FETCH_JSON_SETTLE_MS: int = int(os.environ.get("FETCH_JSON_SETTLE_MS", "1200")) FETCH_JSON_INPAGE_RETRIES: int = int(os.environ.get("FETCH_JSON_INPAGE_RETRIES", "1")) # Пауза между in-page попытками fetch(), мс. FETCH_JSON_RETRY_DELAY_MS: int = int(os.environ.get("FETCH_JSON_RETRY_DELAY_MS", "800")) +# Сколько ждать события `load` на ПОВТОРЕ после гонки «execution context destroyed» +# (#2676). Только на повторе: happy-path остаётся на дешёвом FETCH_JSON_SETTLE_MS, +# иначе бесконечно дозагружающаяся страница удлиняла бы КАЖДЫЙ запрос. Ожидание +# best-effort — по таймауту всё равно пробуем evaluate. +FETCH_JSON_LOAD_WAIT_MS: int = int(os.environ.get("FETCH_JSON_LOAD_WAIT_MS", "15000")) # Известные поставщики. "generic" — фолбэк для всех прочих хостов (один общий # инстанс на неузнанные домены). Порядок задаёт детерминированный health-вывод. @@ -1000,6 +1005,27 @@ async def _do_fetch_json( return await _fetch_json_once( provider, url, method=method, headers=headers, body=body, origin=origin ) + if _is_page_context_lost(exc): + # #2676: браузер жив, умерла ОДНА страница — relaunch не нужен (стоил бы + # ~10-20с и тёплые cookies инстанса). Повторяем на свежей странице, но с + # ожиданием `load`: без него повтор попадает в то же окно клиентской + # навигации, и «транзиентная» ошибка воспроизводится детерминированно. + logger.warning( + "tradein-browser[%s]: страница ушла в навигацию (%s), retry fetch-json " + "с ожиданием load: %s", + provider, + type(exc).__name__, + url, + ) + return await _fetch_json_once( + provider, + url, + method=method, + headers=headers, + body=body, + origin=origin, + wait_for_load=True, + ) raise @@ -1011,6 +1037,7 @@ async def _fetch_json_once( headers: dict, body: object, origin: str, + wait_for_load: bool = False, ) -> dict: """Переходит на origin (same-origin якорь) и выполняет in-page fetch(url). @@ -1029,6 +1056,22 @@ async def _fetch_json_once( await _apply_resource_block(page) await _pace_provider(provider) await page.goto(origin, timeout=BROWSER_NAV_TIMEOUT_MS, wait_until="domcontentloaded") # type: ignore[attr-defined] + if wait_for_load: + # Только ретрай после #2676: даём клиентской навигации доиграть до `load`, + # иначе повтор попадает ровно в то же окно и падает так же. + try: + await page.wait_for_load_state( # type: ignore[attr-defined] + "load", timeout=FETCH_JSON_LOAD_WAIT_MS + ) + except Exception as exc: + # Best-effort: страница может дозагружаться бесконечно (реклама/трекеры). + # Не пробрасываем — evaluate ниже сам скажет, готова страница или нет; + # молчать нельзя, поэтому пишем в лог. + logger.info( + "tradein-browser[%s]: load не дождались (%s), пробуем evaluate как есть", + provider, + type(exc).__name__, + ) # БЕЗ полного BROWSER_WAIT_MS: нам нужен лишь origin-контекст (cookies + # same-origin scope для fetch), а не отрендеренные listings. Settle-паузы # (#1917, FETCH_JSON_SETTLE_MS) хватает, чтобы страница инициализировалась @@ -1106,6 +1149,23 @@ def _is_browser_crash(exc: BaseException) -> bool: ) +# Литерал playwright, а не строка из нашего лога: driver 1.60.0 (в образе сайдкара) +# бросает ровно «Execution context was destroyed» / «... , most likely because of a +# navigation.» — обе формы начинаются одинаково, поэтому хватает одного маркера. +# Проверено grep'ом по playwright/driver/package/lib/coreBundle.js в живом контейнере. +_PAGE_CONTEXT_LOST_MARKER = "execution context was destroyed" + + +def _is_page_context_lost(exc: BaseException) -> bool: + """Страница потеряла JS-контекст (ушла в навигацию между goto и evaluate), #2676. + + НЕ краш браузера: инстанс жив, потеряна одна страница. Поэтому обрабатывается + отдельно от _is_browser_crash — relaunch здесь стоил бы ~10-20с и тёплый профиль + (cookies/фингерпринт инстанса) ради браузера, с которым всё в порядке. + """ + return _PAGE_CONTEXT_LOST_MARKER in str(exc).lower() + + # ── login handler ────────────────────────────────────────────────────────────── diff --git a/tradein-mvp/browser/test_server_fetch_json.py b/tradein-mvp/browser/test_server_fetch_json.py index 0af86622..91466ad0 100644 --- a/tradein-mvp/browser/test_server_fetch_json.py +++ b/tradein-mvp/browser/test_server_fetch_json.py @@ -67,6 +67,7 @@ class _FakePage: def __init__(self, evaluate_result: dict[str, Any]) -> None: self.goto_urls: list[str] = [] self.waits: list[int] = [] # записанные wait_for_timeout(ms) — settle-проверка #1917 + self.load_waits: list[int] = [] # wait_for_load_state("load", timeout=) — #2676 self.closed = 0 # evaluate — AsyncMock, чтобы проверять как сам результат, так и аргументы. self.evaluate = AsyncMock(return_value=evaluate_result) @@ -80,6 +81,9 @@ class _FakePage: async def wait_for_timeout(self, ms: int) -> None: self.waits.append(ms) + async def wait_for_load_state(self, state: str, timeout: int = 0) -> None: + self.load_waits.append(timeout) + async def close(self) -> None: self.closed += 1 @@ -365,3 +369,175 @@ def test_do_fetch_json_relaunch_on_browser_crash(monkeypatch: pytest.MonkeyPatch healthy_page.evaluate.assert_awaited_once() assert crashing_page.closed == 1 assert healthy_page.closed == 1 + + +# ── гонка «execution context was destroyed» (#2676) ─────────────────────────────── + + +class _FakeSequenceBrowser: + """Отдаёт страницы по очереди: первая попытка ≠ вторая (retry на СВЕЖЕЙ странице).""" + + def __init__(self, pages: list[_FakePage]) -> None: + self._pages = list(pages) + self.opened = 0 + + async def new_page(self) -> _FakePage: + page = self._pages[min(self.opened, len(self._pages) - 1)] + self.opened += 1 + return page + + +def _no_relaunch(monkeypatch: pytest.MonkeyPatch) -> list[str]: + """Подменяет _relaunch_browser счётчиком — тест падает, если его всё-таки позвали.""" + calls: list[str] = [] + + async def _fake(provider: str) -> None: + calls.append(provider) + + monkeypatch.setattr(server, "_relaunch_browser", _fake) + return calls + + +# Обе формы, которые бросает playwright 1.60 (driver coreBundle.js) — короткая и полная. +@pytest.mark.parametrize( + "message", + [ + "Page.evaluate: Execution context was destroyed, most likely because of a navigation.", + "Execution context was destroyed", + ], +) +def test_do_fetch_json_retries_on_destroyed_context( + monkeypatch: pytest.MonkeyPatch, message: str +) -> None: + """#2676: страница ушла в навигацию → повтор на свежей странице, БЕЗ relaunch. + + До правки такая ошибка не попадала ни в одну ветку восстановления (_is_browser_crash + матчит только закрытие цели/браузера/соединения) и уезжала наверх как 500. + """ + racing_page = _FakePage({"status": 0, "body": ""}) + racing_page.evaluate = AsyncMock(side_effect=RuntimeError(message)) + settled_page = _FakePage({"status": 200, "body": '{"recovered": true}'}) + + server._browsers["avito"] = _FakeSequenceBrowser([racing_page, settled_page]) + monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000) + monkeypatch.setattr(server, "FETCH_JSON_LOAD_WAIT_MS", 4242) + relaunched = _no_relaunch(monkeypatch) + + result = asyncio.run( + server._do_fetch_json( + "avito", + "https://www.avito.ru/api/x", + method="GET", + headers={}, + body=None, + origin="https://www.avito.ru/", + ) + ) + + assert result == {"status": 200, "body": '{"recovered": true}'} + # Браузер живой — перезапускать его нельзя (тёплые cookies + ~10-20с). + assert relaunched == [] + # Первая попытка НЕ ждала load (happy-path не удлиняется), повтор — ждал. + assert racing_page.load_waits == [] + assert settled_page.load_waits == [4242] + assert racing_page.closed == 1 and settled_page.closed == 1 + + +def test_do_fetch_json_retry_survives_load_timeout(monkeypatch: pytest.MonkeyPatch) -> None: + """Ожидание load на повторе — best-effort: таймаут не отменяет саму попытку.""" + racing_page = _FakePage({"status": 0, "body": ""}) + racing_page.evaluate = AsyncMock(side_effect=RuntimeError("Execution context was destroyed")) + settled_page = _FakePage({"status": 200, "body": "ok"}) + settled_page.wait_for_load_state = AsyncMock( # type: ignore[method-assign] + side_effect=TimeoutError("Timeout 15000ms exceeded") + ) + + server._browsers["avito"] = _FakeSequenceBrowser([racing_page, settled_page]) + monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000) + _no_relaunch(monkeypatch) + + result = asyncio.run( + server._do_fetch_json( + "avito", + "https://www.avito.ru/api/x", + method="GET", + headers={}, + body=None, + origin="https://www.avito.ru/", + ) + ) + assert result == {"status": 200, "body": "ok"} + + +def test_do_fetch_json_gives_up_after_one_context_retry(monkeypatch: pytest.MonkeyPatch) -> None: + """Повтор ровно один: вторая та же ошибка уезжает наверх, а не крутит цикл.""" + message = "Page.evaluate: Execution context was destroyed" + first = _FakePage({"status": 0, "body": ""}) + first.evaluate = AsyncMock(side_effect=RuntimeError(message)) + second = _FakePage({"status": 0, "body": ""}) + second.evaluate = AsyncMock(side_effect=RuntimeError(message)) + + server._browsers["avito"] = _FakeSequenceBrowser([first, second]) + monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000) + _no_relaunch(monkeypatch) + + with pytest.raises(RuntimeError, match="Execution context was destroyed"): + asyncio.run( + server._do_fetch_json( + "avito", + "https://www.avito.ru/api/x", + method="GET", + headers={}, + body=None, + origin="https://www.avito.ru/", + ) + ) + first.evaluate.assert_awaited_once() + second.evaluate.assert_awaited_once() + + +def test_do_fetch_json_does_not_retry_unrelated_error(monkeypatch: pytest.MonkeyPatch) -> None: + """Чужая ошибка НЕ ретраится — повтор невосстановимого жжёт бюджет прогона.""" + page = _FakePage({"status": 0, "body": ""}) + page.evaluate = AsyncMock(side_effect=RuntimeError("net::ERR_PROXY_CONNECTION_FAILED")) + server._browsers["avito"] = _FakeSequenceBrowser([page]) + monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000) + _no_relaunch(monkeypatch) + + with pytest.raises(RuntimeError, match="ERR_PROXY_CONNECTION_FAILED"): + asyncio.run( + server._do_fetch_json( + "avito", + "https://www.avito.ru/api/x", + method="GET", + headers={}, + body=None, + origin="https://www.avito.ru/", + ) + ) + page.evaluate.assert_awaited_once() + + +def test_fetch_json_handler_500_carries_reason_in_body(monkeypatch: pytest.MonkeyPatch) -> None: + """Причина отказа остаётся в теле 500 — её читает _raise_for_sidecar_status (#2708).""" + message = "Execution context was destroyed" + first = _FakePage({"status": 0, "body": ""}) + first.evaluate = AsyncMock(side_effect=RuntimeError(message)) + second = _FakePage({"status": 0, "body": ""}) + second.evaluate = AsyncMock(side_effect=RuntimeError(message)) + server._browsers["avito"] = _FakeSequenceBrowser([first, second]) + monkeypatch.setattr(server, "BROWSER_RECYCLE_PAGES", 10_000) + _no_relaunch(monkeypatch) + + async def _ensure(provider: str, proxy_override: str | None = None) -> bool: + return True + + monkeypatch.setattr(server, "_ensure_browser", _ensure) + + response = asyncio.run( + server.fetch_json_handler( + _make_request({"url": "https://www.avito.ru/api/x", "source": "avito"}) + ) + ) + assert response.status == 500 + assert "Execution context was destroyed" in _json_body(response)["error"] From 663a831775f60b2a307ae6aa79f7fd23ce50eab4 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 09:55:50 +0000 Subject: [PATCH 2/3] =?UTF-8?q?fix(tradein/scraper):=20=D0=BE=D1=82=D0=BC?= =?UTF-8?q?=D0=B5=D1=82=D0=BA=D0=B8=20=D0=B2=D1=80=D0=B5=D0=BC=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=D0=B0=20=D0=BF?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D1=81=D1=82=D0=B0=D1=8E=D1=82=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=BC=D0=B5=D1=80=D0=B7=D0=B0=D1=82=D1=8C=20=D0=B2=20=D0=B5?= =?UTF-8?q?=D0=B3=D0=BE=20=D0=B6=D0=B5=20=D1=82=D1=80=D0=B0=D0=BD=D0=B7?= =?UTF-8?q?=D0=B0=D0=BA=D1=86=D0=B8=D0=B8=20(#2702)=20(#2718)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../backend/app/services/scrape_runs.py | 51 +++- .../223_scrape_runs_time_columns_meaning.sql | 71 ++++++ .../backend/tests/test_2702_run_timestamps.py | 237 ++++++++++++++++++ .../tests/test_scrape_skip_visibility.py | 4 +- .../src/scraper_kit/orchestration/runs.py | 43 +++- .../scraper_kit/orchestration/scheduler.py | 16 +- 6 files changed, 396 insertions(+), 26 deletions(-) create mode 100644 tradein-mvp/backend/data/sql/223_scrape_runs_time_columns_meaning.sql create mode 100644 tradein-mvp/backend/tests/test_2702_run_timestamps.py diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index a2868d65..a7ac714a 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -2,6 +2,31 @@ Таблица scrape_runs создана в 015_scrape_runs.sql. Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled. + +ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL — +синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции и не двигается, +сколько бы та ни жила. Финализаторы (mark_done/mark_failed/mark_banned) выполняются +ТОЙ ЖЕ сессией, что и работа задачи, — и если рабочая транзакция всё это время +оставалась открытой (задача ничего не коммитила: нечего было сохранять, батч читающий, +сохранение шло чужой сессией), их UPDATE попадал ВНУТРЬ неё, и `finished_at` получал +время НАЧАЛА работы, а не её конца. + +Замер на проде 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик +counters.duration_sec): у 153 заявленная длительность превышала собственное окно +finished_at − started_at более чем в 1.5 раза, у 133 окно было меньше секунды при +работе дольше 10 с. 126 из этих 133 окон лежат в диапазоне 9-64 мс — это не разброс, +а подпись механизма: столько проходит от коммита claim'а до первого запроса рабочей +транзакции. Крайний случай — прогон 346 (cian_history_backfill): 18230 с работы, +окно 32 мс. + +Дефект был не сплошной ровно потому, что зависел от того, коммитила ли задача перед +финалом: cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят +поштучно, у них окно совпадало с работой; yandex_address_backfill (45 из 50 прогонов), +newbuilding_enrich, cian_history_backfill — нет. + +Побочно это чинит и `heartbeat_at`: он писался тем же `now()` и по той же причине +отставал от реальности на возраст открытой транзакции, а на нём стоит поиск зависших +прогонов (reap_zombies, порог 6 ч). """ from __future__ import annotations @@ -282,7 +307,10 @@ def _alert_on_run_id( def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: - """INSERT scrape_runs(source, status='running', params, started_at=NOW()). + """INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()). + + started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей + транзакции задачи его уже не достаёт (#2702). Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …); отдельной колонки run_type больше нет — она 3244 прогона подряд молчала @@ -293,7 +321,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: text( """ INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at) - VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW()) + VALUES ( + :source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp() + ) RETURNING id """ ), @@ -305,7 +335,7 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None: - """UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки. + """UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки. total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926). @@ -316,7 +346,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None text( """ UPDATE scrape_runs - SET heartbeat_at = NOW(), + SET heartbeat_at = clock_timestamp(), counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -364,7 +394,7 @@ def is_cancelled(db: Session, run_id: int) -> bool: def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: - """Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count. + """Финализация run: status='done', finished_at + counters + total_seen/new_count. total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся в выделенные колонки — иначе admin/observability показывает 0 (audit #1926). @@ -374,7 +404,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: text( """ UPDATE scrape_runs - SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'done', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -413,7 +444,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) text( """ UPDATE scrape_runs - SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'failed', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), error = :error, counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -472,7 +504,8 @@ def mark_banned( text( """ UPDATE scrape_runs - SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'banned', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), error = :error, counters = CAST(:counters AS jsonb), ban_kind = :ban_kind, total_seen = COALESCE(CAST(:total_seen AS int), total_seen), @@ -583,7 +616,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool: text( """ UPDATE scrape_runs - SET status = 'cancelled', finished_at = NOW() + SET status = 'cancelled', finished_at = clock_timestamp() WHERE id = :run_id AND status = 'running' RETURNING id """ diff --git a/tradein-mvp/backend/data/sql/223_scrape_runs_time_columns_meaning.sql b/tradein-mvp/backend/data/sql/223_scrape_runs_time_columns_meaning.sql new file mode 100644 index 00000000..b1759289 --- /dev/null +++ b/tradein-mvp/backend/data/sql/223_scrape_runs_time_columns_meaning.sql @@ -0,0 +1,71 @@ +-- 223_scrape_runs_time_columns_meaning.sql +-- Purpose (#2702): зафиксировать в схеме, что отметки времени прогона до этой +-- правки не охватывали его работу, и куда смотреть аналитике вместо разности. +-- +-- Dependencies: 015_scrape_runs.sql (создала started_at/finished_at/heartbeat_at), +-- 051_scrape_runs_extend.sql (finished_at/counters). +-- Apply after: 221_backfill_house_suggestions_image_link.sql +-- Идемпотентно: только COMMENT ON COLUMN (перезаписывает сам себя), данных не трогает. +-- +-- ── ЧТО БЫЛО СЛОМАНО ───────────────────────────────────────────────────────── +-- Финализаторы (mark_done / mark_failed / mark_banned) писали finished_at и +-- heartbeat_at через now(). В PostgreSQL now() == transaction_timestamp(): он +-- замерзает на СТАРТЕ транзакции. Финализатор выполняется той же сессией, что и +-- работа задачи; если рабочая транзакция всё это время оставалась открытой (задаче +-- нечего было коммитить — читающий батч, ноль сохранений, сохранение чужой сессией), +-- UPDATE финализатора попадал ВНУТРЬ неё и получал время НАЧАЛА работы. +-- +-- Прод-замер 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик +-- counters.duration_sec): +-- 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза; +-- 133 — окно меньше секунды при работе дольше 10 с. +-- 126 из этих 133 окон лежат в 9-64 мс: это не разброс, а подпись механизма — +-- столько проходит от коммита claim'а до первого запроса рабочей транзакции. +-- Крайние: прогон 346 (cian_history_backfill) — 18 230 с работы при окне 32 мс; +-- 497 (newbuilding_enrich) — 6 124 с при 21 мс; 341 (yandex_address_backfill) — +-- 1 460 с при 19 мс. +-- +-- Дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом. +-- Средние окно/duration_sec по источникам на том же замере: +-- cian_history_backfill 2554 / 4222 ← окно короче работы +-- newbuilding_enrich 1302 / 2078 ← короче +-- yandex_address_backfill 107 / 1022 ← короче в 10 раз +-- avito_detail_backfill 2829 / 2016 ← длиннее (норма) +-- house_imv_backfill 1173 / 1173 ← совпадает +-- cadastral_geo_match 11 / 10 ← совпадает +-- +-- ── ПОЧЕМУ ИСТОРИЮ НЕ ЧИНИМ ────────────────────────────────────────────────── +-- Восстановить настоящий finished_at по строке нельзя: реальное время конца нигде +-- не сохранилось. Но counters.duration_sec измерялся монотонными часами процесса +-- (time.monotonic / time.time в самих задачах) и транзакцией не затронут — он у +-- этих строк верный. Поэтому история не переписывается, а помечается: аналитика +-- обязана брать длительность из счётчика, а не из разности отметок. +-- +-- Правка кода (clock_timestamp() вместо now() во всех финализаторах и в heartbeat) +-- живёт в app/services/scrape_runs.py + packages/scraper-kit/.../orchestration/runs.py. + +COMMENT ON COLUMN scrape_runs.started_at IS + 'Момент claim''а прогона. Пишется create_run своей транзакцией (commit сразу ' + 'после INSERT), поэтому откат рабочей транзакции задачи его не затрагивает. ' + 'С #2702 — clock_timestamp(); до него now() (= старт транзакции тика планировщика), ' + 'что давало сдвиг в пределах тика.'; + +COMMENT ON COLUMN scrape_runs.finished_at IS + 'Момент финализации прогона. ВНИМАНИЕ: у строк ДО #2702 (2026-08-06) значение ' + 'недостоверно — писалось now() (= transaction_timestamp) внутри рабочей транзакции ' + 'задачи, поэтому у прогонов, ничего не коммитивших по ходу работы, равно времени ' + 'её НАЧАЛА. На проде так вышло у 133 из 487 прогонов со счётчиком длительности ' + '(окно < 1 с при работе > 10 с). Длительность таких прогонов брать из ' + 'counters->>''duration_sec'' (монотонные часы процесса, транзакцией не затронуты), ' + 'а НЕ из finished_at − started_at.'; + +COMMENT ON COLUMN scrape_runs.heartbeat_at IS + 'Последний признак жизни прогона; на нём стоит поиск зависших (reap_zombies, порог ' + '6 ч). У строк ДО #2702 отставал от реальности на возраст открытой рабочей ' + 'транзакции по той же причине, что finished_at, — то есть критерий «завис» решал ' + 'по замороженной отметке. С #2702 пишется clock_timestamp().'; + +COMMENT ON COLUMN scrape_runs.counters IS + 'Счётчики прогона (jsonb). Ключ duration_sec, где он есть, измерен монотонными ' + 'часами процесса и остаётся единственным достоверным источником длительности для ' + 'строк до #2702 (см. комментарий к finished_at).'; diff --git a/tradein-mvp/backend/tests/test_2702_run_timestamps.py b/tradein-mvp/backend/tests/test_2702_run_timestamps.py new file mode 100644 index 00000000..99e2792f --- /dev/null +++ b/tradein-mvp/backend/tests/test_2702_run_timestamps.py @@ -0,0 +1,237 @@ +"""#2702: отметки времени прогона не охватывают его работу. + +Что было. Финализаторы писали `finished_at`/`heartbeat_at` через `now()`, а `now()` +в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. +Выполняются финализаторы той же сессией, что и работа задачи, поэтому если рабочая +транзакция всё это время оставалась открытой, их UPDATE попадал ВНУТРЬ неё и получал +время НАЧАЛА работы. + +Прод-замер 2026-08-06 (487 прогонов с finished_at и counters.duration_sec): + * 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза; + * 133 — окно меньше секунды при работе дольше 10 с, причём 126 из них укладываются + в 9-64 мс: столько проходит от коммита claim'а до первого запроса рабочей + транзакции — подпись механизма, а не разброс. +Крайний случай, воспроизведённый ниже дословно: прогон 346 (cian_history_backfill) — +18 230 с работы, окно 32 мс. + +Почему дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом. +`cadastral_geo_match` / `house_imv_backfill` / `avito_detail_backfill` коммитят +поштучно — у них окно совпадало с работой; `yandex_address_backfill` (45 из 50 +прогонов), `newbuilding_enrich`, `cian_history_backfill` — нет. + +Фальсификация: на старом коде (`now()`) тесты 1 и 3 дают другой ответ. +""" + +from __future__ import annotations + +import inspect +import os +import re +from typing import Any + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from scraper_kit.orchestration import runs as kit_runs +from scraper_kit.orchestration import scheduler as kit_scheduler + +from app.services import scrape_runs as app_runs + +_MODULES = {"kit": kit_runs, "app": app_runs} + +# Колонки, которые обязаны нести НАСТОЯЩЕЕ время, а не время старта транзакции. +_TS_COLS = ("started_at", "finished_at", "heartbeat_at") + + +def _split_top(items: str) -> list[str]: + """Разбить список SQL-элементов по запятым ВЕРХНЕГО уровня (CAST(:x AS t) — один).""" + out: list[str] = [] + depth = 0 + cur = "" + for ch in items: + if ch == "," and depth == 0: + out.append(cur.strip()) + cur = "" + continue + depth += (ch == "(") - (ch == ")") + cur += ch + out.append(cur.strip()) + return out + + +class _Row: + """Строка ответа: id для create_run, source для alert-хука.""" + + id = 1 + source = "src" + status = "done" + + +class _FakeResult: + def fetchone(self) -> _Row: + return _Row() + + def first(self) -> _Row: + return _Row() + + def fetchall(self) -> list[_Row]: + return [] + + +class _FakePg: + """Мини-модель PostgreSQL на две функции времени и ленивую транзакцию. + + `now()` == `transaction_timestamp()` — замерзает на старте транзакции; + `clock_timestamp()` — настоящие часы. Транзакция открывается лениво на первом + execute (autobegin SQLAlchemy) и закрывается commit/rollback. Больше модель + ничего не умеет — этого достаточно, чтобы отличить одно от другого. + """ + + def __init__(self) -> None: + self.wall: float = 0.0 # «стенные часы» теста + self.tx_start: float | None = None + self.row: dict[str, float] = {} # что осело в scrape_runs + self.calls: list[str] = [] + + def _stamp(self, col: str, func: str) -> None: + assert self.tx_start is not None + self.row[col] = self.tx_start if func.lower() == "now" else self.wall + + def execute(self, stmt: Any, params: Any = None) -> _FakeResult: + if self.tx_start is None: + self.tx_start = self.wall + self.calls.append("execute") + sql = str(stmt) + for col in _TS_COLS: # UPDATE ... SET = () + m = re.search(rf"\b{col}\s*=\s*(now|clock_timestamp)\s*\(\s*\)", sql, re.I) + if m is not None: + self._stamp(col, m.group(1)) + ins = re.search( + r"INSERT INTO scrape_runs\s*\((.*?)\).*?VALUES\s*\((.*)\)", sql, re.S | re.I + ) + if ins is not None: # INSERT — отметки времени позиционные, в VALUES + for col, val in zip(_split_top(ins.group(1)), _split_top(ins.group(2)), strict=False): + f = re.fullmatch(r"(now|clock_timestamp)\s*\(\s*\)", val, re.I) + if col in _TS_COLS and f is not None: + self._stamp(col, f.group(1)) + return _FakeResult() + + def commit(self) -> None: + self.calls.append("commit") + self.tx_start = None + + def rollback(self) -> None: + self.calls.append("rollback") + self.tx_start = None + + +# ── 1. Финал прогона датируется концом работы, а не стартом транзакции ──────── + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_finished_at_covers_the_work_not_the_transaction_start(name: str) -> None: + """Прогон 346 дословно: 18 230 с работы в одной незакоммиченной транзакции. + + На старом коде finished_at = 0.032 (старт рабочей транзакции) → окно 32 мс при + пяти часах работы. Это и есть та строка, ради которой заведена задача. + """ + mod = _MODULES[name] + db = _FakePg() + db.row["started_at"] = 0.0 # create_run уже закоммитил claim + + db.wall = 0.020 + mod.update_heartbeat(db, 1, {"listings_processed": 0}) # heartbeat + commit + + db.wall = 0.032 + db.execute("SELECT id FROM listings WHERE history IS NULL") # рабочая транзакция + + db.wall = 18230.0 # пять часов работы, ни одного коммита + mod.mark_done(db, 1, {"listings_processed": 1200, "duration_sec": 18230}) + + assert db.row["finished_at"] == pytest.approx(18230.0) + assert db.row["heartbeat_at"] == pytest.approx(18230.0) + window = db.row["finished_at"] - db.row["started_at"] + assert window == pytest.approx(18230.0), "окно прогона обязано охватывать его работу" + + +@pytest.mark.parametrize("name", list(_MODULES)) +@pytest.mark.parametrize("finalizer", ["mark_failed", "mark_banned"]) +def test_failed_and_banned_finals_are_wall_clock_too(name: str, finalizer: str) -> None: + """Тот же инвариант для неуспешных финалов. + + Их спасал defensive-rollback в начале (он закрывал рабочую транзакцию), но + полагаться на побочный эффект чужой защиты нельзя — проверяем явно. + """ + mod = _MODULES[name] + db = _FakePg() + db.wall = 0.019 + db.execute("SELECT 1") # рабочая транзакция открыта + db.wall = 1460.0 + getattr(mod, finalizer)(db, 1, "boom", {"checked": 0, "duration_sec": 1460}) + assert db.row["finished_at"] == pytest.approx(1460.0) + + +# ── 2. started_at переживает откат рабочей транзакции ──────────────────────── + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_started_at_survives_rolled_back_work_transaction(name: str) -> None: + """Требование #2702 п.1: отметка старта живёт в СВОЕЙ закоммиченной транзакции. + + create_run коммитит INSERT до возврата run_id, а ни один финализатор прогона + started_at не переписывает — поэтому откат рабочей транзакции его не достаёт. + Финал при этом обязан быть датирован концом работы (это и падает на старом коде). + """ + mod = _MODULES[name] + db = _FakePg() + + run_id = mod.create_run(db, source="cian_history_backfill", params={}) + assert run_id == 1 + assert db.calls == ["execute", "commit"], "INSERT прогона обязан коммититься сразу" + + db.wall = 0.030 + db.execute("UPDATE listings SET address = 'x'") # рабочая транзакция + db.wall = 100.0 + db.rollback() # работа упала и откатилась + + db.wall = 100.5 + db.execute("SELECT count(*) FROM listings") # новая рабочая транзакция + db.wall = 1460.0 + mod.mark_done(db, 1, {"checked": 0, "duration_sec": 1460}) + + assert db.row["started_at"] == pytest.approx(0.0), "started_at не должен сдвигаться" + assert db.row["finished_at"] == pytest.approx(1460.0) + + +# ── 3. Инвариант источника: никакая отметка времени не пишется now() ───────── + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_no_run_timestamp_is_written_with_now(name: str) -> None: + """`now()` в этих модулях не имеет корректного применения — его быть не должно. + + Проверяем весь исходник, а не отдельные запросы: INSERT в create_run пишет + started_at/heartbeat_at позиционно (в VALUES), и построчная проверка его бы + пропустила — ровно так дефект и дожил до 3 300 прогонов. + """ + src = inspect.getsource(_MODULES[name]) + # Регистрозависимо: SQL в этих модулях пишется в верхнем регистре, а строчное + # `now()` встречается в объяснительной прозе docstring'ов — ловим SQL, не текст. + assert re.search(r"\bNOW\s*\(\s*\)", src) is None + assert "clock_timestamp()" in src + + +def test_zombie_criterion_compares_real_clocks() -> None: + """#2702 п.2: поиск зависших сравнивает записанный heartbeat со «сейчас». + + Обе стороны сравнения обязаны быть настоящим временем: на проде у всех 6 + прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и + остался на отметке старта (max advance 0.0 с) — критерий решал по замороженной + отметке, хотя нормальный прогон этого источника длится до 5.06 ч. + """ + src = inspect.getsource(kit_scheduler.reap_zombies) + stmt = re.search(r"UPDATE scrape_runs.*?RETURNING id", src, re.S) + assert stmt is not None + assert re.search(r"\bNOW\s*\(\s*\)", stmt.group(0)) is None + assert stmt.group(0).count("clock_timestamp()") == 2 # finished_at + порог сравнения diff --git a/tradein-mvp/backend/tests/test_scrape_skip_visibility.py b/tradein-mvp/backend/tests/test_scrape_skip_visibility.py index a16b5ae2..0683f683 100644 --- a/tradein-mvp/backend/tests/test_scrape_skip_visibility.py +++ b/tradein-mvp/backend/tests/test_scrape_skip_visibility.py @@ -129,7 +129,9 @@ def test_mark_skipped_collapse_refreshes_detail_and_started_at() -> None: update = db.sql_of("UPDATE scrape_runs r") assert update is not None sql, params = update - assert "started_at = NOW()" in sql + # clock_timestamp(), а не NOW(): отметки времени прогона пишутся настоящими + # часами, иначе внутри долгой открытой транзакции они замерзают (#2702). + assert "started_at = clock_timestamp()" in sql assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска" assert "first_skip_at" in sql, "начало стрика потеряно" assert params["details"] == "37 дн. назад" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index bbe367ca..ecb26a58 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -9,6 +9,15 @@ app-копии: 2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim), app-копии эта функция не нужна. + +ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL — +синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы +выполняются той же сессией, что и работа задачи, и если её транзакция оставалась +открытой всё время работы (нечего было коммитить), их UPDATE попадал ВНУТРЬ неё — +`finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с +окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в +диапазоне 9-64 мс (= от коммита claim'а до первого запроса рабочей транзакции). +Полный разбор — в docstring app-копии `app/services/scrape_runs.py`. """ from __future__ import annotations @@ -297,7 +306,10 @@ def _alert_on_run_id( def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: - """INSERT scrape_runs(source, status='running', params, started_at=NOW()). + """INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()). + + started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей + транзакции задачи его уже не достаёт (#2702). Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …); отдельной колонки run_type больше нет — она 3244 прогона подряд молчала @@ -308,7 +320,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int: text( """ INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at) - VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW()) + VALUES ( + :source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp() + ) RETURNING id """ ), @@ -357,9 +371,9 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = LIMIT 1 ) UPDATE scrape_runs r - SET heartbeat_at = NOW(), - started_at = NOW(), - finished_at = NOW(), + SET heartbeat_at = clock_timestamp(), + started_at = clock_timestamp(), + finished_at = clock_timestamp(), counters = COALESCE(r.counters, '{}'::jsonb) || jsonb_build_object( 'skips', COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1, 'detail', CAST(:details AS text), @@ -385,7 +399,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = source, status, error, counters, started_at, heartbeat_at, finished_at ) VALUES ( - :source, 'skipped', :reason, CAST(:counters AS jsonb), NOW(), NOW(), NOW() + :source, 'skipped', :reason, CAST(:counters AS jsonb), clock_timestamp(), clock_timestamp(), clock_timestamp() ) RETURNING id """ @@ -409,7 +423,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None: - """UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки. + """UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки. total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926). @@ -420,7 +434,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None text( """ UPDATE scrape_runs - SET heartbeat_at = NOW(), + SET heartbeat_at = clock_timestamp(), counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -447,7 +461,7 @@ def is_cancelled(db: Session, run_id: int) -> bool: def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: - """Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count. + """Финализация run: status='done', finished_at + counters + total_seen/new_count. total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся в выделенные колонки — иначе admin/observability показывает 0 (audit #1926). @@ -457,7 +471,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: text( """ UPDATE scrape_runs - SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'done', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -496,7 +511,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) text( """ UPDATE scrape_runs - SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'failed', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), error = :error, counters = CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) @@ -555,7 +571,8 @@ def mark_banned( text( """ UPDATE scrape_runs - SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(), + SET status = 'banned', + finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), error = :error, counters = CAST(:counters AS jsonb), ban_kind = :ban_kind, total_seen = COALESCE(CAST(:total_seen AS int), total_seen), @@ -586,7 +603,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool: text( """ UPDATE scrape_runs - SET status = 'cancelled', finished_at = NOW() + SET status = 'cancelled', finished_at = clock_timestamp() WHERE id = :run_id AND status = 'running' RETURNING id """ diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index e012353d..c670a298 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -235,16 +235,26 @@ def has_running_run(db: Session, source: str) -> bool: def reap_zombies(db: Session) -> int: - """Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours.""" + """Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours. + + clock_timestamp(), а не now() (#2702): критерий сравнивает ЗАПИСАННЫЙ heartbeat со + временем «сейчас», и обе стороны сравнения должны быть настоящим временем. `now()` + замерзает на старте транзакции — у писавшей heartbeat стороны это давало отставание + на весь возраст открытой рабочей транзакции (см. docstring runs.py), у читающей + стороны — на возраст тика. На проде это уже стоило ложных срабатываний: у всех 6 + прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и + остался на отметке старта (max advance 0.0 с) — при том что нормальный прогон этого + источника длится до 5.06 ч (прогон 346) и обязан был двигать heartbeat. + """ zombie_interval = f"{ZOMBIE_THRESHOLD_HOURS} hours" result = db.execute( text( """ UPDATE scrape_runs - SET status = 'zombie', finished_at = NOW() + SET status = 'zombie', finished_at = clock_timestamp() WHERE status = 'running' AND (heartbeat_at IS NULL - OR heartbeat_at < NOW() - CAST(:interval AS interval)) + OR heartbeat_at < clock_timestamp() - CAST(:interval AS interval)) RETURNING id """ ), From 01440928561ccbc4b4cdbdda7843678d571ba9b4 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 09:55:54 +0000 Subject: [PATCH 3/3] =?UTF-8?q?fix(tradein/houses):=20camelCase-=D1=82?= =?UTF-8?q?=D0=B8=D0=BF=D1=8B=20=D0=B4=D0=BE=D0=BC=D0=BE=D0=B2=20=D0=BF?= =?UTF-8?q?=D1=80=D0=B8=D0=B2=D0=BE=D0=B4=D1=8F=D1=82=D1=81=D1=8F=20=D0=BA?= =?UTF-8?q?=20=D0=BA=D0=B0=D0=BD=D0=BE=D0=BD=D1=83=20=D1=83=20=D0=B8=D1=81?= =?UTF-8?q?=D1=82=D0=BE=D1=87=D0=BD=D0=B8=D0=BA=D0=B0=20(#2678)=20(#2719)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../data/sql/224_houses_house_type_canon.sql | 83 ++++++++++++ .../tests/test_2678_house_type_canon.py | 128 ++++++++++++++++++ .../src/scraper_kit/house_type_normalizer.py | 5 + .../src/scraper_kit/providers/avito/houses.py | 19 ++- 4 files changed, 233 insertions(+), 2 deletions(-) create mode 100644 tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql create mode 100644 tradein-mvp/backend/tests/test_2678_house_type_canon.py diff --git a/tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql b/tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql new file mode 100644 index 00000000..b8b58859 --- /dev/null +++ b/tradein-mvp/backend/data/sql/224_houses_house_type_canon.sql @@ -0,0 +1,83 @@ +-- 224_houses_house_type_canon.sql +-- Issue #2678 (хвост #2675/#2674): twin миграции 141 для таблицы ДОМОВ. +-- +-- Миграция 141 привела camelCase-вокабуляр Циана к канону только в listings. +-- В houses он остался — и каждый читатель типа дома чинил его у себя (#2675 +-- починил домовую оценку; поштучный путь и подбор аналогов продолжали сравнивать +-- 'monolithBrick' с 'monolith_brick' и не совпадать). +-- +-- ЗАМЕР ПРОДА 2026-08-06 (SELECT source, house_type, count(*) FROM houses GROUP BY 1,2): +-- канон: brick 388+5+4+1 · panel 325+3+2+1 · monolith 118+37+6+1 · +-- block 90+1+1 · monolith_brick 2+1 · wood 2 +-- camelCase: monolithBrick 48 (derived) + 8 (cian_newbuilding) = 56 · +-- gasSilicateBlock 1 · aerocreteBlock 1 +-- прочее: stalin 3 · other 18 · wireframe 1 +-- NULL: 8440 из 8880 строк (тип дома вообще неизвестен — не наш случай) +-- +-- ЖИВОГО ПИСАТЕЛЯ camelCase В houses НЕТ: у всех 80 неканоничных строк +-- last_scraped_at = 2026-05-24 14:04:20.012209 — одна и та же метка, т.е. +-- единственный прогон backfill'а 063 (промоут типов из listings ДО миграции 141). +-- Единственный живой писатель houses.house_type — avito-каталог домов +-- (providers/avito/houses.py), он пишет русские подписи через свою карту; в том +-- же PR он переведён на общий normalize_house_type, чтобы неизвестное значение +-- шло как NULL, а не как 'other' (его единственный источник неканона). +-- +-- ПРОВЕРКА СМЫСЛА ПЕРЕД СКЛЕЙКОЙ (требование #2674 — не слепить разное): +-- контрольная группа в своих же данных. Для каждой неканоничной строки взяты +-- типы её ЖЕ объявлений (listings.house_id_fk), уже нормализованных 141: +-- monolithBrick 56 домов — monolith_brick присутствует у ВСЕХ 56 → одно и то же +-- stalin 3 дома — brick (совпадает с решением 141: «сталинка» = кирпич) +-- aerocreteBlock 1 дом — block +-- gasSilicateBlock 1 дом — block +-- other 18 домов — разброс monolith/monolith_brick/brick/block, т.е. +-- 'other' = «неизвестно», а не отдельный материал +-- wireframe 1 дом — wireframe и у объявлений (само-согласовано) +-- Вывод: склейка безопасна ТОЛЬКО для четырёх camelCase-токенов + stalin. +-- +-- ЧТО НАМЕРЕННО НЕ ТРОГАЕМ: +-- 'other' (18) и 'wireframe' (1) — честного соответствия в каноне нет +-- (см. #2675: normalize_house_type схлопывает их в None на чтении, и это +-- правильный ответ — NULL нейтрален для soft-penalty эстиматора, а выдуманный +-- материал был бы враньём). Стирать их здесь тоже не будем: это единственный +-- след того, что источник что-то про дом сказал. +-- +-- BACKFILL (счётчики сняты на проде ДО применения, 2026-08-06): +-- monolithBrick -> monolith_brick : 56 строк +-- stalin -> brick : 3 строки +-- aerocreteBlock -> block : 1 строка +-- gasSilicateBlock -> block : 1 строка +-- foamConcreteBlock-> block : 0 строк (в houses не встречается, +-- оставлен для паритета с картой 141) +-- ИТОГО ожидаемо тронуто: 61 строка. +-- +-- Idempotent: WHERE перечисляет только мапимые токены → повторный прогон 0 строк. +-- Маппинг тождественен house_type_normalizer._RAW_TO_CANON и миграции 141 — +-- третьего словаря не заводим. + +BEGIN; + +UPDATE houses + SET house_type = CASE house_type + WHEN 'monolithBrick' THEN 'monolith_brick' + WHEN 'gasSilicateBlock' THEN 'block' + WHEN 'aerocreteBlock' THEN 'block' + WHEN 'foamConcreteBlock' THEN 'block' + WHEN 'stalin' THEN 'brick' + ELSE house_type + END + WHERE house_type IN ( + 'monolithBrick', 'gasSilicateBlock', 'aerocreteBlock', + 'foamConcreteBlock', 'stalin' + ); + +COMMENT ON COLUMN houses.house_type IS + 'Материал/тип дома, канон: panel/brick/monolith/monolith_brick/block/wood ' + '(тот же enum, что listings.house_type и scraper_kit.house_type_normalizer). ' + 'Писать сюда только через normalize_house_type — сырые вокабуляры источников ' + '(cian camelCase monolithBrick/gasSilicateBlock/stalin, yandex SCREAMING ' + 'MONOLIT_BRICK, русские подписи Авито) приводятся ДО записи, миграция 224 ' + 'вычистила исторические. Вне канона осталось намеренно: other (источник сказал ' + '«другое») и wireframe (каркас — материала в каноне нет). Неизвестный тип = ' + 'NULL, а не панель и не other: NULL нейтрален для soft-penalty эстиматора.'; + +COMMIT; diff --git a/tradein-mvp/backend/tests/test_2678_house_type_canon.py b/tradein-mvp/backend/tests/test_2678_house_type_canon.py new file mode 100644 index 00000000..c349a222 --- /dev/null +++ b/tradein-mvp/backend/tests/test_2678_house_type_canon.py @@ -0,0 +1,128 @@ +"""#2678: тип дома приводится к канону У ИСТОЧНИКА, и словарь ровно один. + +Кейсы взяты не из головы, а из фактического замера прода 2026-08-06 +(`SELECT source, house_type, count(*) FROM houses GROUP BY 1,2`): + + monolithBrick 56 · other 18 · stalin 3 · aerocreteBlock 1 · + gasSilicateBlock 1 · wireframe 1 · плюс канон (brick/panel/monolith/ + monolith_brick/block/wood) и 8440 NULL. + +Проверяется три вещи: + 1. каждый фактический вариант → канон (или честный None); + 2. живой писатель houses.house_type (avito-каталог) больше не изобретает + 'other' и ходит через общий нормализатор; + 3. миграция 224 не заводит третий словарь — её CASE совпадает с картой кода. +""" + +from __future__ import annotations + +import re +from pathlib import Path + +import pytest +from scraper_kit.house_type_normalizer import _RAW_TO_CANON, normalize_house_type +from scraper_kit.providers.avito.houses import _normalize_house_type as avito_house_type + +_MIGRATION_224 = ( + Path(__file__).resolve().parents[1] / "data" / "sql" / "224_houses_house_type_canon.sql" +) + +# Фактический словарь houses.house_type на проде 2026-08-06 → чем он обязан стать. +# None = «честно неизвестно» (NULL нейтрален для soft-penalty эстиматора, в отличие +# от выдуманного материала). +_PROD_VALUES: list[tuple[str, str | None]] = [ + ("monolithBrick", "monolith_brick"), # 56 строк + ("other", None), # 18 строк — источник сказал «другое», материала нет + ("stalin", "brick"), # 3 строки — «сталинка» = кирпич (решение миграции 141) + ("aerocreteBlock", "block"), # 1 строка + ("gasSilicateBlock", "block"), # 1 строка + ("wireframe", None), # 1 строка — каркас, в каноне такого материала нет + ("brick", "brick"), + ("panel", "panel"), + ("monolith", "monolith"), + ("monolith_brick", "monolith_brick"), + ("block", "block"), + ("wood", "wood"), +] + + +@pytest.mark.parametrize(("raw", "expected"), _PROD_VALUES) +def test_prod_value_maps_to_canon(raw: str, expected: str | None) -> None: + assert normalize_house_type(raw) == expected + + +def test_canon_survives_uppercase_including_monolith_brick() -> None: + """#2678 п.6: сквозной проброс канона был регистрозависим — кроме monolith_brick. + + Значений в верхнем регистре в базе сегодня ноль; это страховка на новый источник, + который отдаст канон «как в документации». + """ + assert normalize_house_type("MONOLITH_BRICK") == "monolith_brick" + assert normalize_house_type("Monolith_Brick") == "monolith_brick" + # Остальной канон и раньше переживал регистр — фиксируем, что не сломали. + for token in ("BRICK", "Panel", "MONOLITH", "Block", "WOOD"): + assert normalize_house_type(token) == token.lower() + + +# ── живой писатель houses.house_type: avito-каталог домов ──────────────────────── + + +@pytest.mark.parametrize( + ("label", "expected"), + [ + ("Монолитно-кирпичный", "monolith_brick"), + ("Панельный", "panel"), + ("КИРПИЧНЫЙ", "brick"), + (" Блочный ", "block"), + ("Деревянный", "wood"), + ], +) +def test_avito_house_label_maps_to_canon(label: str, expected: str) -> None: + assert avito_house_type(label) == expected + + +def test_avito_unknown_label_is_null_not_other() -> None: + """Незнакомая подпись → NULL. До #2678 здесь появлялось 'other'. + + 'other' всегда != канону, т.е. читатель получал не «неизвестно», а гарантированное + несовпадение: ложный штраф при подборе аналогов и пропуск оценки. + """ + assert avito_house_type("Саманный") is None + assert avito_house_type("") is None + assert avito_house_type(None) is None + + +def test_avito_writer_handles_foreign_vocabulary() -> None: + """Писатель ходит через общий нормализатор, а не только через свою карту.""" + assert avito_house_type("monolithBrick") == "monolith_brick" + assert avito_house_type("MONOLIT_BRICK") == "monolith_brick" + + +# ── миграция 224: тот же словарь, что в коде ───────────────────────────────────── + + +def _migration_case_pairs() -> dict[str, str]: + """WHEN 'x' THEN 'y' из исполняемой части миграции (без `--`-комментариев).""" + code = "\n".join( + line.split("--", 1)[0] for line in _MIGRATION_224.read_text(encoding="utf-8").splitlines() + ) + return dict(re.findall(r"WHEN\s+'([^']+)'\s+THEN\s+'([^']+)'", code)) + + +def test_migration_224_mapping_matches_code() -> None: + """Миграция не заводит третий словарь — каждая пара есть в _RAW_TO_CANON.""" + pairs = _migration_case_pairs() + assert pairs, "в миграции 224 не нашлось ни одного WHEN ... THEN" + for raw, canon in pairs.items(): + assert _RAW_TO_CANON.get(raw) == canon, f"{raw!r} расходится с house_type_normalizer" + + +def test_migration_224_touches_only_mapped_tokens() -> None: + """WHERE ограничен теми же токенами → 'other'/'wireframe'/канон не трогаются.""" + code = "\n".join( + line.split("--", 1)[0] for line in _MIGRATION_224.read_text(encoding="utf-8").splitlines() + ) + where_tokens = set(re.findall(r"'([A-Za-z]+)'", code.split("WHERE", 1)[1].split(";", 1)[0])) + assert where_tokens == set(_migration_case_pairs()) + assert "other" not in where_tokens + assert "wireframe" not in where_tokens diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/house_type_normalizer.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/house_type_normalizer.py index f9ad280d..aacb99ba 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/house_type_normalizer.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/house_type_normalizer.py @@ -47,6 +47,11 @@ _RAW_TO_CANON: dict[str, str] = { "panel": "panel", "block": "block", "wood": "wood", + # #2678 п.6: канон целиком, включая monolith_brick. Pass-through по _CANON + # регистрозависим, и без этого ключа 'MONOLITH_BRICK'/'Monolith_Brick' (форма, + # в которой канон может прийти от нового источника) уезжали бы в None — + # единственный канонический токен без такой страховки. + "monolith_brick": "monolith_brick", "monolithBrick": "monolith_brick", "gasSilicateBlock": "block", "aerocreteBlock": "block", diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/houses.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/houses.py index a62f1755..c91b256b 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/houses.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/houses.py @@ -39,6 +39,7 @@ from sqlalchemy import text from sqlalchemy.orm import Session from scraper_kit.avito_exceptions import AvitoBlockedError, AvitoRateLimitedError +from scraper_kit.house_type_normalizer import normalize_house_type from scraper_kit.providers.avito.serp import _is_firewall_page from scraper_kit.providers.avito.shared import RUS_MONTHS, _unix_to_date @@ -258,10 +259,24 @@ def _strip_price(price_str: str | None) -> int | None: def _normalize_house_type(raw: str | None) -> str | None: - """Нормализует тип дома: "Монолитный" → "monolith".""" + """Нормализует тип дома: "Монолитный" → "monolith". Незнакомое → None. + + Единственный живой писатель houses.house_type (см. save-функцию ниже), поэтому + канон обязан приводиться ЗДЕСЬ, а не у каждого читателя (#2678). + + Русская подпись Авито снимается локальной HOUSE_TYPE_MAP, результат прогоняется + через общий scraper_kit.house_type_normalizer: он знает и канон, и чужие + вокабуляры (cian camelCase, yandex SCREAMING) — на случай, если карточка дома + однажды придёт с чужим токеном. + + #2678: раньше незнакомое значение становилось 'other' — единственный источник + неканоничных значений в houses среди живых писателей. 'other' всегда != канон, + т.е. для soft-penalty эстиматора это ложный штраф, а NULL нейтрален (тот же + довод, что в docstring house_type_normalizer). + """ if not raw: return None - return HOUSE_TYPE_MAP.get(raw.lower(), "other") + return normalize_house_type(HOUSE_TYPE_MAP.get(raw.strip().lower(), raw)) def _normalize_house_class(raw: str | None) -> str | None: