Merge pull request 'fix(scraper-kit): свежесть источника — по последнему прогону, принёсшему данные, а не по статусу done' (#3365) from fix/3172-freshness-by-data into main
Some checks failed
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m8s
Deploy Trade-In / build-backend (push) Successful in 1m59s
Deploy Trade-In / deploy (push) Successful in 2m5s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Failing after 12s
Some checks failed
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m8s
Deploy Trade-In / build-backend (push) Successful in 1m59s
Deploy Trade-In / deploy (push) Successful in 2m5s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Failing after 12s
This commit is contained in:
commit
4a7e648068
3 changed files with 290 additions and 15 deletions
|
|
@ -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)
|
@pytest.fixture(autouse=True)
|
||||||
def _reset_digest_clock() -> Any:
|
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:
|
def test_digest_emits_one_event_listing_all_stale_sources() -> None:
|
||||||
sentry = MagicMock()
|
sentry = MagicMock()
|
||||||
with patch.object(sched, "sentry_sdk", sentry):
|
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
|
assert [s.source for s in stale] == PROD_STALE
|
||||||
sentry.capture_message.assert_called_once()
|
sentry.capture_message.assert_called_once()
|
||||||
msg = sentry.capture_message.call_args[0][0]
|
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 не мешает сводке — она меряет календарь, а не серию."""
|
"""Главное свойство: стрик 0 не мешает сводке — она меряет календарь, а не серию."""
|
||||||
sentry = MagicMock()
|
sentry = MagicMock()
|
||||||
with patch.object(sched, "sentry_sdk", sentry):
|
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"}
|
zero_streak = {"avito_full_load_exhaustive", "cian_history_backfill"}
|
||||||
assert zero_streak <= {s.source for s in stale}
|
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()
|
sentry = MagicMock()
|
||||||
fresh = [_row("avito_city_sweep", None, 1.1), _row("sber_index_pull", 7, 4.1)]
|
fresh = [_row("avito_city_sweep", None, 1.1), _row("sber_index_pull", 7, 4.1)]
|
||||||
with patch.object(sched, "sentry_sdk", sentry):
|
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()
|
sentry.capture_message.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
def test_digest_is_daily_not_per_tick() -> None:
|
def test_digest_is_daily_not_per_tick() -> None:
|
||||||
"""Планировщик тикает раз в минуту; сводка обязана выходить раз в сутки."""
|
"""Планировщик тикает раз в минуту; сводка обязана выходить раз в сутки."""
|
||||||
sentry = MagicMock()
|
sentry = MagicMock()
|
||||||
db = _db(PROD_ROWS)
|
db = _db(PROD_RUN_ROWS)
|
||||||
with patch.object(sched, "sentry_sdk", sentry):
|
with patch.object(sched, "sentry_sdk", sentry):
|
||||||
sched.emit_stale_digest(db, now=NOW)
|
sched.emit_stale_digest(db, now=NOW)
|
||||||
sched.emit_stale_digest(db, now=NOW + timedelta(minutes=1))
|
sched.emit_stale_digest(db, now=NOW + timedelta(minutes=1))
|
||||||
|
|
|
||||||
181
tradein-mvp/backend/tests/test_freshness_by_data_3172.py
Normal file
181
tradein-mvp/backend/tests/test_freshness_by_data_3172.py
Normal file
|
|
@ -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
|
||||||
|
|
@ -108,20 +108,30 @@ SKIP_UNKNOWN_SOURCE = "unknown_source"
|
||||||
STALE_DIGEST_INTERVAL_FACTOR = 3
|
STALE_DIGEST_INTERVAL_FACTOR = 3
|
||||||
STALE_DIGEST_PERIOD_H = 24
|
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("""
|
_STALE_SOURCES_SQL = text("""
|
||||||
SELECT sch.source,
|
SELECT sch.source,
|
||||||
sch.default_params->>'interval_days' AS interval_days,
|
sch.default_params->>'interval_days' AS interval_days,
|
||||||
COALESCE(
|
sch.created_at,
|
||||||
(SELECT max(r.finished_at) FROM scrape_runs r
|
r.finished_at,
|
||||||
WHERE r.source = sch.source AND r.status = 'done'),
|
r.status,
|
||||||
sch.created_at
|
r.counters
|
||||||
) AS since,
|
|
||||||
(NOT EXISTS (SELECT 1 FROM scrape_runs r
|
|
||||||
WHERE r.source = sch.source AND r.status = 'done')) AS never_ok
|
|
||||||
FROM scrape_schedules sch
|
FROM scrape_schedules sch
|
||||||
|
LEFT JOIN scrape_runs r
|
||||||
|
ON r.source = sch.source AND r.finished_at IS NOT NULL
|
||||||
WHERE sch.enabled
|
WHERE sch.enabled
|
||||||
""")
|
""")
|
||||||
|
|
||||||
|
|
@ -155,6 +165,69 @@ def _schedule_interval_days(raw: Any) -> int:
|
||||||
return 1
|
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]:
|
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 []
|
return []
|
||||||
try:
|
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
|
_last_stale_digest_at = now
|
||||||
if not stale:
|
if not stale:
|
||||||
logger.info("scheduler: stale-digest — просроченных источников нет (#2670)")
|
logger.info("scheduler: stale-digest — просроченных источников нет (#2670)")
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue