From 3a0b94e0521335450e6e93a1d658882a288e9766 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 00:02:21 +0500 Subject: [PATCH] =?UTF-8?q?fix(scheduler):=20=D0=BC=D0=B5=D1=80=D0=B8?= =?UTF-8?q?=D1=82=D1=8C=20=D1=81=D0=B2=D0=B5=D0=B6=D0=B5=D1=81=D1=82=D1=8C?= =?UTF-8?q?=20=D0=B8=D1=81=D1=82=D0=BE=D1=87=D0=BD=D0=B8=D0=BA=D0=B0=20?= =?UTF-8?q?=D0=B4=D0=B0=D0=BD=D0=BD=D1=8B=D0=BC=D0=B8,=20=D0=B0=20=D0=BD?= =?UTF-8?q?=D0=B5=20=D1=81=D1=82=D0=B0=D1=82=D1=83=D1=81=D0=BE=D0=BC=20?= =?UTF-8?q?=D0=BF=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `_STALE_SOURCES_SQL` считала свежесть возрастом последнего прогона со статусом 'done'. Прогон с блоком честно финализируется как 'banned' (#2657) и при этом вставляет строки: у domclick_city_sweep 886 строк 25.08 и 128 строк 23.08 — оба 'banned'. Источник, регулярно ловящий блок и столь же регулярно приносящий данные, числился мёртвым навсегда, отсюда ложный P1 #3118 «домклик не собирается с 5 августа». Запрос отдаёт завершённые прогоны, решение «прогон дал данные» принимает `run_brought_data` ТЕМ ЖЕ результатным словарём, которым уже судит сторож нулевого результата (runs._RESULT_COUNTER_KEYS, #2703) — одна мера на оба механизма. Не 'lots_inserted': это новизна, а не наличие данных (здоровый дедуплицированный sweep вставляет ноль). Результат не измерен (28 источников без результатного ключа) → судим прежней мерой, статусом: «не измерено» ≠ «ноль». never_ok считается той же мерой, иначе соврал бы в другую сторону. Closes #3172 --- .../tests/test_2670_stale_source_digest.py | 28 ++- .../tests/test_freshness_by_data_3172.py | 181 ++++++++++++++++++ .../scraper_kit/orchestration/scheduler.py | 96 ++++++++-- 3 files changed, 290 insertions(+), 15 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_freshness_by_data_3172.py diff --git a/tradein-mvp/backend/tests/test_2670_stale_source_digest.py b/tradein-mvp/backend/tests/test_2670_stale_source_digest.py index baf2ec4a..3af2e3a4 100644 --- a/tradein-mvp/backend/tests/test_2670_stale_source_digest.py +++ b/tradein-mvp/backend/tests/test_2670_stale_source_digest.py @@ -83,6 +83,26 @@ PROD_STALE = [ ] +def _run_row(row: Any) -> Any: + """Та же правда в кодировке #3172: строка запроса — ПРОГОН, а не готовая свёртка. + + Свежесть теперь меряется последним прогоном, ПРИНЁСШИМ ДАННЫЕ (`freshness_rows`), + поэтому запрос отдаёт сырые поля прогона; порог, который проверяет этот файл, от + правки не зависит — он считает те же `since`/`never_ok`, просто свёрнутые в Python. + """ + return SimpleNamespace( + source=row.source, + interval_days=row.interval_days, + created_at=NOW - timedelta(days=400), + finished_at=row.since, + status="done", + counters={"lots_fetched": 1}, + ) + + +PROD_RUN_ROWS = [_run_row(r) for r in PROD_ROWS] + + @pytest.fixture(autouse=True) def _reset_digest_clock() -> Any: """Выпуск сводки помнится в памяти модуля — сбрасываем между тестами.""" @@ -131,7 +151,7 @@ def test_never_successful_source_is_reported_with_a_flag() -> None: def test_digest_emits_one_event_listing_all_stale_sources() -> None: sentry = MagicMock() with patch.object(sched, "sentry_sdk", sentry): - stale = sched.emit_stale_digest(_db(PROD_ROWS), now=NOW) + stale = sched.emit_stale_digest(_db(PROD_RUN_ROWS), now=NOW) assert [s.source for s in stale] == PROD_STALE sentry.capture_message.assert_called_once() msg = sentry.capture_message.call_args[0][0] @@ -145,7 +165,7 @@ def test_digest_covers_the_two_sources_the_ladder_cannot_reach() -> None: """Главное свойство: стрик 0 не мешает сводке — она меряет календарь, а не серию.""" sentry = MagicMock() with patch.object(sched, "sentry_sdk", sentry): - stale = sched.emit_stale_digest(_db(PROD_ROWS), now=NOW) + stale = sched.emit_stale_digest(_db(PROD_RUN_ROWS), now=NOW) zero_streak = {"avito_full_load_exhaustive", "cian_history_backfill"} assert zero_streak <= {s.source for s in stale} @@ -154,14 +174,14 @@ def test_digest_is_quiet_when_everything_is_fresh() -> None: sentry = MagicMock() fresh = [_row("avito_city_sweep", None, 1.1), _row("sber_index_pull", 7, 4.1)] with patch.object(sched, "sentry_sdk", sentry): - assert sched.emit_stale_digest(_db(fresh), now=NOW) == [] + assert sched.emit_stale_digest(_db([_run_row(r) for r in fresh]), now=NOW) == [] sentry.capture_message.assert_not_called() def test_digest_is_daily_not_per_tick() -> None: """Планировщик тикает раз в минуту; сводка обязана выходить раз в сутки.""" sentry = MagicMock() - db = _db(PROD_ROWS) + db = _db(PROD_RUN_ROWS) with patch.object(sched, "sentry_sdk", sentry): sched.emit_stale_digest(db, now=NOW) sched.emit_stale_digest(db, now=NOW + timedelta(minutes=1)) diff --git a/tradein-mvp/backend/tests/test_freshness_by_data_3172.py b/tradein-mvp/backend/tests/test_freshness_by_data_3172.py new file mode 100644 index 00000000..1ea859c2 --- /dev/null +++ b/tradein-mvp/backend/tests/test_freshness_by_data_3172.py @@ -0,0 +1,181 @@ +"""Свежесть источника меряется данными, а не статусом прогона (#3172). + +Дефект: `_STALE_SOURCES_SQL` считала свежесть возрастом последнего прогона со статусом +'done'. Прогон, поймавший блок, честно финализируется как 'banned' (#2657) и ПРИ ЭТОМ +вставляет строки — у domclick_city_sweep 886 строк 25.08 и 128 строк 23.08, оба +'banned'. Источник, регулярно ловящий блок и столь же регулярно приносящий данные, для +сводки был мёртв навсегда: отсюда ложный P1 #3118 «домклик не собирается с 5 августа». + +Мок-строки несут ОДНУ И ТУ ЖЕ правду БД в двух кодировках: сырые поля прогона +(finished_at/status/counters), которые читает новая свёртка, и since/never_ok, которые +СТАРЫЙ код считал бы по правилу `status='done'` (`_old_rule`, а не руками). Поэтому +`git apply -R` даёт красное ПО ЗНАЧЕНИЮ, а не ImportError: + - (а) домклик снова попадает в просроченные (старая мера смотрит на 05.08); + - (б) источник с ежедневным 'done' и нулевой выдачей ПЕРЕСТАЁТ быть просроченным + (старая мера засчитывает пустой 'done' как свежесть); + - (в) источник с одними 'banned'-прогонами с данными снова «не собирал ни разу». +""" + +from __future__ import annotations + +import os +from datetime import UTC, datetime +from types import SimpleNamespace +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from scraper_kit.orchestration import scheduler as sched + +NOW = datetime(2026, 8, 26, 8, 0, tzinfo=UTC) +CREATED = datetime(2026, 5, 1, 0, 0, tzinfo=UTC) + + +def _run(day: int, status: str, **counters: Any) -> dict[str, Any]: + return { + "finished_at": datetime(2026, 8, day, 6, 0, tzinfo=UTC), + "status": status, + "counters": counters, + } + + +def _old_rule(runs: list[dict[str, Any]]) -> tuple[datetime, bool]: + """Свежесть ДО правки: max(finished_at) по прогонам со статусом 'done'.""" + done = [r["finished_at"] for r in runs if r["status"] == "done"] + return (max(done) if done else CREATED), not done + + +def _rows(source: str, runs: list[dict[str, Any]], interval_days: Any = None) -> list[Any]: + """Строки запроса: прогон на строку + since/never_ok в кодировке старого правила.""" + since, never_ok = _old_rule(runs) + return [ + SimpleNamespace( + source=source, + interval_days=interval_days, + created_at=CREATED, + since=since, + never_ok=never_ok, + **run, + ) + for run in runs + ] or [ + SimpleNamespace( + source=source, + interval_days=interval_days, + created_at=CREATED, + since=since, + never_ok=never_ok, + finished_at=None, + status=None, + counters=None, + ) + ] + + +def _db(rows: list[Any]) -> MagicMock: + db = MagicMock() + db.execute.return_value.fetchall.return_value = rows + return db + + +def _stale(rows: list[Any]) -> dict[str, sched.StaleSource]: + sched._last_stale_digest_at = None # сводка выходит раз в сутки — снимаем анти-спам + with patch.object(sched, "sentry_sdk", MagicMock()): + return {s.source: s for s in sched.emit_stale_digest(_db(rows), now=NOW)} + + +@pytest.fixture(autouse=True) +def _reset_digest_clock() -> Any: + sched._last_stale_digest_at = None + yield + sched._last_stale_digest_at = None + + +# ── (а) блок с данными — это сбор, а не смерть источника ───────────────────── + +DOMCLICK = [ + _run(5, "done", lots_fetched=1200, lots_inserted=430), + _run(23, "banned", lots_fetched=310, lots_inserted=128), + _run(25, "banned", lots_fetched=940, lots_inserted=886), +] + + +def test_banned_run_with_data_keeps_source_fresh() -> None: + """Прод-случай #3118: последние прогоны 'banned', но строки собраны → НЕ просрочен.""" + assert "domclick_city_sweep" not in _stale(_rows("domclick_city_sweep", DOMCLICK)) + + +def test_freshness_is_the_last_run_that_brought_data() -> None: + """since = 25.08 (banned, 940 в выдаче), а не 05.08 (последний 'done').""" + (row,) = sched.freshness_rows(_rows("domclick_city_sweep", DOMCLICK)) + assert row.since == datetime(2026, 8, 25, 6, 0, tzinfo=UTC) + assert row.never_ok is False + + +# ── (б) пустой 'done' свежести не даёт ─────────────────────────────────────── + + +def test_daily_done_runs_with_zero_results_are_stale() -> None: + """Ежедневный 'done' с пустой выдачей: последние данные 01.08 → 25 суток просрочки.""" + runs = [_run(1, "done", lots_fetched=800, lots_inserted=210)] + runs += [_run(day, "done", lots_fetched=0, lots_inserted=0) for day in range(2, 26)] + + stale = _stale(_rows("yandex_city_sweep", runs)) + + assert "yandex_city_sweep" in stale + assert stale["yandex_city_sweep"].age_days == pytest.approx(25.08, abs=0.1) + + +def test_source_without_result_metric_is_judged_by_status_as_before() -> None: + """КОНТРОЛЬ (зелёный до и после): 'не измерено' ≠ 'ноль' — судим прежней мерой. + + refresh_search_matview и 27 других источников не пишут результатного ключа вовсе; + строгое «данные > 0» разом объявило бы их всех просроченными навсегда. + """ + fresh = _rows("refresh_search_matview", [_run(25, "done")]) + old = _rows("refresh_search_matview", [_run(1, "done")]) + assert "refresh_search_matview" not in _stale(fresh) + assert "refresh_search_matview" in _stale(old) + + +# ── (в) never_ok согласован с той же мерой ─────────────────────────────────── + + +def test_never_ok_counts_banned_runs_that_brought_data() -> None: + """Источник, собиравший только сквозь баны, СОБИРАЛ — «ни разу» о нём соврало бы. + + Через сводку, а не только через `freshness_rows`: на откате правки это красное ПО + ЗНАЧЕНИЮ (старая мера не видит ни одного 'done' → since=created_at, never_ok=True, + просрочка 117 суток), а не «функции нет». + """ + rows = _rows("avito_full_load", [_run(24, "banned", lots_fetched=17, lots_inserted=17)]) + + assert "avito_full_load" not in _stale(rows) + assert sched.freshness_rows(rows)[0].never_ok is False + + +def test_never_ok_stays_true_when_no_run_ever_brought_data() -> None: + """Обратная сторона: прогоны есть, данных не было ни в одном → «ни разу», от created_at.""" + runs = [_run(2, "banned", lots_fetched=0), _run(9, "failed", lots_fetched=0)] + + (row,) = sched.freshness_rows(_rows("cian_history_backfill", runs)) + + assert row.never_ok is True + assert row.since == CREATED + + +def test_schedule_without_runs_is_never_ok() -> None: + (row,) = sched.freshness_rows(_rows("brand_new_source", [])) + assert (row.never_ok, row.since) == (True, CREATED) + + +# ── дополнение к логическим тестам: сам запрос больше не фильтрует по статусу ─ + + +def test_query_no_longer_selects_by_run_status() -> None: + sql = str(sched._STALE_SOURCES_SQL) + assert "status = 'done'" not in sql + assert "LEFT JOIN scrape_runs" in sql 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 7454d657..711ddd0e 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 @@ -108,20 +108,30 @@ SKIP_UNKNOWN_SOURCE = "unknown_source" STALE_DIGEST_INTERVAL_FACTOR = 3 STALE_DIGEST_PERIOD_H = 24 -# Возраст последнего УСПЕШНОГО ('done') прогона на каждое включённое расписание. -# COALESCE(last_ok, created_at): у расписания без единого успеха отсчёт идёт от его -# создания — иначе «никогда не собирал» выглядело бы как «нет данных, судить нечем». +# Завершённые прогоны каждого включённого расписания — «принёс ли прогон данные» +# решается в Python (`run_brought_data`), а не статусом в WHERE (#3172). +# +# Раньше здесь стоял `status = 'done'`, и свежесть источника равнялась возрасту +# последнего прогона с этим статусом. Но прогон, поймавший блок, честно финализируется +# как 'banned' (#2657) — И ПРИ ЭТОМ ВСТАВЛЯЕТ СТРОКИ: у domclick_city_sweep 886 строк +# 25.08 и 128 строк 23.08, оба прогона 'banned'. Источник, который регулярно ловит блок +# и столь же регулярно приносит данные, числился мёртвым навсегда → ложный P1 #3118 +# («домклик не собирается с 5 августа», хотя сбор шёл). +# +# ponytail: полный проход по завершённым прогонам включённых источников (порядка 10k +# строк на проде) — раз в STALE_DIGEST_PERIOD_H часов. Понадобится дешевле — оконный +# фильтр по finished_at, но тогда never_ok перестанет быть честным («не собирал НИ +# РАЗУ» превратится в «не собирал в окне»). _STALE_SOURCES_SQL = text(""" SELECT sch.source, sch.default_params->>'interval_days' AS interval_days, - COALESCE( - (SELECT max(r.finished_at) FROM scrape_runs r - WHERE r.source = sch.source AND r.status = 'done'), - sch.created_at - ) AS since, - (NOT EXISTS (SELECT 1 FROM scrape_runs r - WHERE r.source = sch.source AND r.status = 'done')) AS never_ok + sch.created_at, + r.finished_at, + r.status, + r.counters FROM scrape_schedules sch + LEFT JOIN scrape_runs r + ON r.source = sch.source AND r.finished_at IS NOT NULL WHERE sch.enabled """) @@ -155,6 +165,69 @@ def _schedule_interval_days(raw: Any) -> int: return 1 +def run_brought_data(status: str | None, counters: Mapping[str, Any] | None) -> bool: + """Дал ли завершённый прогон данные — мера свежести источника (#3172). + + Статус НЕ участвует, пока результат измерим: прогон с блоком ('banned', #2657) + успевает записать выдачу так же, как штатный 'done'. Меряем тем же результатным + словарём, которым уже судит сторож нулевого результата (`runs._RESULT_COUNTER_KEYS`, + #2703) — одна мера «сколько объявлений отдала выдача» на оба механизма, а не второй + список ключей, который разъедется с первым. + + НЕ `lots_inserted`: это НОВИЗНА, а не наличие данных. Здоровый sweep, у которого вся + выдача уже в базе, вставляет ноль строк — по такой мере живой источник читался бы + мёртвым, то есть ровно дефект #3118 с другой стороны. + + `_run_result_count() is None` — прогон результат НЕ СООБЩИЛ (28 источников на проде + не имеют результатного ключа вовсе: refresh_search_matview, deactivate_stale_*, + мониторы). Судить нечем → зачитываем прежнюю меру, успешный статус. Иначе «не + измерено» схлопнулось бы с «измерено, ноль» и все они разом стали бы просроченными + навсегда — та же ложная тревога, только оптом. + """ + result = _kit_runs._run_result_count(counters) + if result is None: + return status == "done" + return result > 0 + + +@dataclass(frozen=True) +class FreshnessRow: + """Свёртка прогонов одного расписания: когда источник ПОСЛЕДНИЙ РАЗ дал данные.""" + + source: str + interval_days: Any + since: datetime + never_ok: bool + + +def freshness_rows(rows: list[Any]) -> list[FreshnessRow]: + """Строки `_STALE_SOURCES_SQL` (прогон на строку) → свежесть на источник. + + since = последний прогон, ПРИНЁСШИЙ ДАННЫЕ (любой статус); нет такого — created_at + расписания, как и раньше. never_ok считается ТОЙ ЖЕ мерой: иначе «никогда не + собирал» соврал бы в другую сторону — источник с одними лишь 'banned'-прогонами с + данными числился бы ни разу не собиравшим. + """ + meta: dict[str, tuple[Any, datetime]] = {} + last_data: dict[str, datetime] = {} + for row in rows: + meta.setdefault(row.source, (row.interval_days, row.created_at)) + if row.finished_at is None or not run_brought_data(row.status, row.counters): + continue + prev = last_data.get(row.source) + if prev is None or row.finished_at > prev: + last_data[row.source] = row.finished_at + return [ + FreshnessRow( + source=source, + interval_days=interval_days, + since=last_data.get(source) or created_at, + never_ok=source not in last_data, + ) + for source, (interval_days, created_at) in meta.items() + ] + + def stale_sources(rows: list[Any], now: datetime) -> list[StaleSource]: """Чистая часть сводки: какие расписания просрочены и на сколько (свежие — внизу). @@ -193,7 +266,8 @@ def emit_stale_digest(db: Session, *, now: datetime | None = None) -> list[Stale ): return [] try: - stale = stale_sources(list(db.execute(_STALE_SOURCES_SQL).fetchall()), now) + run_rows = list(db.execute(_STALE_SOURCES_SQL).fetchall()) + stale = stale_sources(freshness_rows(run_rows), now) _last_stale_digest_at = now if not stale: logger.info("scheduler: stale-digest — просроченных источников нет (#2670)") -- 2.45.3