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