fix(ptica): свежесть data-table источника считается по успешным строкам (#2956) (#2957)
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m28s
Deploy / build-worker (push) Successful in 5m26s
Deploy / deploy (push) Successful in 1m25s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 9s
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m28s
Deploy / build-worker (push) Successful in 5m26s
Deploy / deploy (push) Successful in 1m25s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 9s
This commit is contained in:
parent
feff8214f7
commit
f7e8228550
3 changed files with 234 additions and 5 deletions
|
|
@ -1414,6 +1414,15 @@ class FreshnessSource(BaseModel):
|
|||
# внутри окна и не ложно-срабатывает (#1947 fix). Default 1 → флагует только если
|
||||
# суммарный выход цикла = 0 (безопасный минимальный catch).
|
||||
min_output_rows: int = 1
|
||||
# Data-table режим: условие «строка означает УСПЕХ». Без него свежесть считается по
|
||||
# факту записи строки, а не по факту получения данных — и провалившийся загрузчик,
|
||||
# исправно пишущий строку с ошибкой, вечно выглядит свежим. Ровно это и случилось с
|
||||
# nspd: последний успешный дамп 27.07.2026, а монитор молчал 24 суток, потому что
|
||||
# каждый упавший harvest обновлял fetched_at_utc (#2956).
|
||||
# Run-ledger режиму не нужно: там успех уже отделён через FILTER (WHERE status='done').
|
||||
# Значение — статическая SQL-строка ИЗ КОДА (не из пользовательского ввода), она
|
||||
# подставляется в FILTER (WHERE ...) как есть.
|
||||
success_where: str | None = None
|
||||
|
||||
|
||||
# Реестр источников. Run-ledger таблицы (kn/objective/nspd_geo/cadastre) проверены на
|
||||
|
|
@ -1503,6 +1512,12 @@ _FRESHNESS_SOURCES: list[FreshnessSource] = [
|
|||
# defunct nspd_scrape_runs (manual WAF-ban 2026-04-30) больше НЕ источник истины.
|
||||
table="nspd_quarter_dumps",
|
||||
timestamp_col="fetched_at_utc",
|
||||
# Свежесть — по УСПЕШНЫМ дампам. Упавший harvest всё равно пишет строку
|
||||
# (fetched_at_utc проставлен, harvest_error заполнен, счётчики нулевые), и без
|
||||
# этого условия каждый провал обновлял часы свежести. С 03.08.2026 провалились
|
||||
# все 61 дамп подряд, последний успешный — 27.07, а источник числился fresh
|
||||
# (#2956).
|
||||
success_where="harvest_error IS NULL",
|
||||
# В timestamp-режиме не используется — оставляем валидное имя колонки.
|
||||
work_col="total_features",
|
||||
# Медленный кадастровый + lazy-refresh источник: дампы освежаются по мере
|
||||
|
|
@ -1591,7 +1606,9 @@ def compute_freshness(db: Session) -> dict[str, Any]:
|
|||
|
||||
Для data-table источников (src.timestamp_col задан, напр. nspd → nspd_quarter_dumps)
|
||||
нет run-ledger семантики (status/started/finished отсутствуют), поэтому:
|
||||
- last_success_at = last_attempt_at = MAX(timestamp_col)
|
||||
- last_attempt_at = MAX(timestamp_col); last_success_at — то же, но по
|
||||
строкам, прошедшим success_where (у источника с колонкой ошибки это
|
||||
отделяет «строку записали» от «данные получили», #2956)
|
||||
- objects_updated_24h / _7d = COUNT(*) строк, обновлённых в окне
|
||||
- last_status = NULL (косметика только для run-ledger'ов)
|
||||
Остальной downstream (age_days / _classify_freshness / status-маппинг) — общий.
|
||||
|
|
@ -1608,16 +1625,23 @@ def compute_freshness(db: Session) -> dict[str, Any]:
|
|||
# все временные границы передаются параметрами (:d1/:d7).
|
||||
if src.timestamp_col is not None:
|
||||
# Data-table режим: плоская контент-таблица без run-ledger семантики
|
||||
# (нет status/started/finished). Свежесть = MAX(timestamp_col),
|
||||
# upd_24h/_7d = COUNT(*) строк, обновлённых в окне. last_status=NULL
|
||||
# (косметика только для run-ledger'ов).
|
||||
# (нет status/started/finished). Свежесть = MAX(timestamp_col) по строкам,
|
||||
# прошедшим success_where (если задан; иначе по всем), upd_24h/_7d =
|
||||
# COUNT(*) строк, обновлённых в окне. last_status=NULL (косметика только
|
||||
# для run-ledger'ов).
|
||||
ts = src.timestamp_col
|
||||
# Успех vs попытка. last_attempt_at — всегда MAX(ts) (строка записана),
|
||||
# last_success_at — только по строкам, прошедшим success_where. Это тот же
|
||||
# раздел, что в run-ledger ветке ниже (FILTER (WHERE status = 'done')):
|
||||
# без него упавший загрузчик, который исправно пишет строку с ошибкой,
|
||||
# выглядит свежим вечно (#2956).
|
||||
success_filter = f" FILTER (WHERE {src.success_where})" if src.success_where else ""
|
||||
row = (
|
||||
db.execute(
|
||||
text(
|
||||
f"""
|
||||
SELECT
|
||||
MAX({ts}) AS last_success_at,
|
||||
MAX({ts}){success_filter} AS last_success_at,
|
||||
MAX({ts}) AS last_attempt_at,
|
||||
COALESCE(COUNT(*) FILTER (
|
||||
WHERE {ts} > NOW() - CAST(:d1 AS interval)
|
||||
|
|
|
|||
|
|
@ -96,3 +96,13 @@ tests/services/scrapers/test_ekb_ppt_tep_parser.py::TestParsePptTepRealPdf::test
|
|||
tests/services/scrapers/test_ekb_ppt_tep_parser.py::TestParsePptTepRealPdf::test_tep_has_rows
|
||||
tests/services/scrapers/test_ekb_ppt_tep_parser.py::TestParsePptTepRealPdf::test_zone_balance_has_itogo
|
||||
tests/services/scrapers/test_ekb_ppt_tep_parser.py::TestParsePptTepRealPdf::test_zone_balance_has_rows
|
||||
|
||||
# ── #2956: свежесть data-table источника считается по успешным строкам ────────
|
||||
# Нужен живой Postgres: тест создаёт ВРЕМЕННЫЕ копии всех таблиц реестра freshness
|
||||
# и гоняет по ним настоящий SQL compute_freshness. В CI ЭТИ ТЕСТЫ ИДУТ — postgres-
|
||||
# сервис поднят (#2745), как и для соседних tests/sql/*. Записи нужны только для
|
||||
# машины без БД и без туннеля на 15432.
|
||||
tests/sql/test_2956_freshness_ignores_failed_dumps.py::test_failed_dumps_do_not_refresh_the_clock
|
||||
tests/sql/test_2956_freshness_ignores_failed_dumps.py::test_successful_dump_still_counts_as_fresh
|
||||
tests/sql/test_2956_freshness_ignores_failed_dumps.py::test_attempt_is_still_recorded
|
||||
tests/sql/test_2956_freshness_ignores_failed_dumps.py::test_only_failures_means_no_success_at_all
|
||||
|
|
|
|||
195
backend/tests/sql/test_2956_freshness_ignores_failed_dumps.py
Normal file
195
backend/tests/sql/test_2956_freshness_ignores_failed_dumps.py
Normal file
|
|
@ -0,0 +1,195 @@
|
|||
"""Свежесть data-table источника считается по УСПЕШНЫМ строкам, а не по факту записи (#2956).
|
||||
|
||||
Что случилось на проде. `nspd_quarter_dumps` — контент-таблица: harvest пишет строку и
|
||||
когда всё получилось, и когда упал (тогда `harvest_error` заполнен, счётчики нулевые, но
|
||||
`fetched_at_utc` всё равно проставлен). Монитор свежести брал `MAX(fetched_at_utc)` без
|
||||
разбора — и каждый ПРОВАЛ обновлял часы свежести.
|
||||
|
||||
Последний успешный дамп — 27.07.2026. С 03.08 провалились все 61 подряд. Источник при
|
||||
этом числился `ok` (возраст ~3 дня при пороге 14), сторож молчал 24 суток. Слепота по
|
||||
построению: монитор мерил «записали ли мы строку», а не «получили ли мы данные».
|
||||
|
||||
Правильный образец лежал в соседней ветке того же `if`: run-ledger режим отделяет успех
|
||||
от попытки через `FILTER (WHERE status = 'done')`.
|
||||
|
||||
Тест герметичный: все таблицы реестра создаются ВРЕМЕННЫМИ в своей же сессии, поэтому
|
||||
он не зависит от схемы CI-базы и ничего не читает из настоящих таблиц. Скипается, если
|
||||
Postgres недоступен (та же конвенция, что в tests/sql/test_velocity_alerts.py).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
import pytest
|
||||
from sqlalchemy import create_engine, text
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
|
||||
def _dsn() -> str:
|
||||
raw = os.environ.get("TEST_DATABASE_URL") or os.environ.get(
|
||||
"DATABASE_URL",
|
||||
"postgresql+psycopg://gendesign@localhost:15432/gendesign",
|
||||
)
|
||||
return (
|
||||
raw
|
||||
if raw.startswith("postgresql+")
|
||||
else raw.replace("postgresql://", "postgresql+psycopg://")
|
||||
)
|
||||
|
||||
|
||||
def _db_reachable() -> tuple[bool, str]:
|
||||
try:
|
||||
eng = create_engine(_dsn(), connect_args={"connect_timeout": 3})
|
||||
with eng.connect() as c:
|
||||
c.execute(text("SELECT 1"))
|
||||
return True, ""
|
||||
except Exception as exc:
|
||||
return False, str(exc)
|
||||
|
||||
|
||||
_DB_OK, _DB_ERR = _db_reachable()
|
||||
pytestmark = pytest.mark.skipif(not _DB_OK, reason=f"Postgres недоступен: {_DB_ERR}")
|
||||
|
||||
# Таблицы реестра freshness. Создаём ВСЕ, иначе compute_freshness упадёт на первой же
|
||||
# отсутствующей — он проходит по всему реестру, а не только по проверяемому источнику.
|
||||
_TEMP_SCHEMA = """
|
||||
CREATE TEMP TABLE kn_scrape_runs (
|
||||
status text, started_at timestamptz, finished_at timestamptz,
|
||||
objects_count int, flats_count int) ON COMMIT DROP;
|
||||
CREATE TEMP TABLE objective_scrape_runs (
|
||||
status text, started_at timestamptz, finished_at timestamptz, rows_lots int) ON COMMIT DROP;
|
||||
CREATE TEMP TABLE nspd_geo_jobs (
|
||||
status text, started_at timestamptz, finished_at timestamptz,
|
||||
created_at timestamptz, targets_done int) ON COMMIT DROP;
|
||||
CREATE TEMP TABLE cadastre_jobs (
|
||||
status text, started_at timestamptz, finished_at timestamptz,
|
||||
created_at timestamptz, targets_done int) ON COMMIT DROP;
|
||||
CREATE TEMP TABLE gisogd_permits (id bigint, fetched_at timestamptz) ON COMMIT DROP;
|
||||
CREATE TEMP TABLE nspd_quarter_dumps (
|
||||
quarter_cad text, fetched_at_utc timestamptz,
|
||||
total_features int, harvest_error text) ON COMMIT DROP;
|
||||
"""
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def db():
|
||||
engine = create_engine(_dsn())
|
||||
session = sessionmaker(bind=engine)()
|
||||
try:
|
||||
session.execute(text(_TEMP_SCHEMA))
|
||||
# Защита, а не украшение: если ЛЮБАЯ из таблиц реестра не создалась временной
|
||||
# (опечатка в имени, изменившийся реестр), запросы уйдут в НАСТОЯЩУЮ таблицу —
|
||||
# тест станет зависеть от боевых данных и сможет их читать. По конвенции репо
|
||||
# (tests/sql/*) DSN по умолчанию смотрит в туннель к прод-базе, так что цена
|
||||
# такой опечатки реальна. Пустая таблица сразу после создания — признак того,
|
||||
# что затенение сработало: у настоящих таблиц строки есть.
|
||||
for table in (
|
||||
"kn_scrape_runs",
|
||||
"objective_scrape_runs",
|
||||
"nspd_geo_jobs",
|
||||
"cadastre_jobs",
|
||||
"gisogd_permits",
|
||||
"nspd_quarter_dumps",
|
||||
):
|
||||
n = session.execute(text(f"SELECT count(*) FROM {table}")).scalar()
|
||||
assert n == 0, (
|
||||
f"{table}: запрос попал НЕ во временную таблицу ({n} строк) — "
|
||||
"тест читал бы боевые данные, а его выводы были бы про них"
|
||||
)
|
||||
yield session
|
||||
finally:
|
||||
session.rollback()
|
||||
session.close()
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _nspd(db) -> dict:
|
||||
from app.api.v1.admin_scrape import compute_freshness
|
||||
|
||||
payload = compute_freshness(db)
|
||||
rows = [s for s in payload["sources"] if s["source"] == "nspd"]
|
||||
assert (
|
||||
len(rows) == 1
|
||||
), f"источник nspd не найден в реестре: {[s['source'] for s in payload['sources']]}"
|
||||
return rows[0]
|
||||
|
||||
|
||||
def _insert_dump(db, *, days_ago: float, error: str | None, features: int) -> None:
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO nspd_quarter_dumps
|
||||
(quarter_cad, fetched_at_utc, total_features, harvest_error)
|
||||
VALUES (:cad, :ts, :feat, :err)
|
||||
"""
|
||||
),
|
||||
{
|
||||
"cad": f"66:41:{int(days_ago * 100):07d}",
|
||||
"ts": datetime.now(UTC) - timedelta(days=days_ago),
|
||||
"feat": features,
|
||||
"err": error,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def test_failed_dumps_do_not_refresh_the_clock(db) -> None:
|
||||
"""Прод-картина: успех 24 дня назад, дальше только провалы каждый день.
|
||||
|
||||
На origin/main источник получает age_days ≈ 1 и статус ok — ровно то молчание,
|
||||
что длилось 24 суток. С правкой возраст считается от последнего УСПЕХА.
|
||||
"""
|
||||
_insert_dump(db, days_ago=24.0, error=None, features=1831)
|
||||
for day in range(1, 18): # 17 провалов подряд, самый свежий — вчера
|
||||
_insert_dump(
|
||||
db, days_ago=float(day), error="TimeoutError: The read operation timed out", features=0
|
||||
)
|
||||
|
||||
src = _nspd(db)
|
||||
assert src["age_days"] is not None, "возраст не посчитан — успехов не нашлось вовсе"
|
||||
assert src["age_days"] > 20, (
|
||||
f"возраст {src['age_days']}д взят от ПРОВАЛИВШЕГОСЯ дампа: провал обновил часы "
|
||||
"свежести, монитор ослеп (#2956)"
|
||||
)
|
||||
assert src["status"] in {"stale", "failed"}, (
|
||||
f"статус {src['status']!r} при последнем успехе 24 дня назад и пороге "
|
||||
f"fresh_days=14 — источник числится живым, алёрта не будет"
|
||||
)
|
||||
|
||||
|
||||
def test_successful_dump_still_counts_as_fresh(db) -> None:
|
||||
"""Контроль: свежий УСПЕХ по-прежнему даёт ok — правка не ужесточила лишнего."""
|
||||
_insert_dump(db, days_ago=1.0, error=None, features=1200)
|
||||
_insert_dump(db, days_ago=0.5, error="TimeoutError", features=0)
|
||||
|
||||
src = _nspd(db)
|
||||
assert src["status"] == "ok", f"свежий успешный дамп даёт статус {src['status']!r}"
|
||||
assert src["age_days"] < 2
|
||||
|
||||
|
||||
def test_attempt_is_still_recorded(db) -> None:
|
||||
"""Контроль: попытки не потеряны — last_attempt_at видит и провалившиеся строки.
|
||||
|
||||
Разделение успех/попытка должно быть именно разделением, а не отбрасыванием: иначе
|
||||
в UI пропадёт признак «загрузчик ходит, но не приносит».
|
||||
"""
|
||||
_insert_dump(db, days_ago=24.0, error=None, features=1831)
|
||||
_insert_dump(db, days_ago=0.2, error="TimeoutError", features=0)
|
||||
|
||||
src = _nspd(db)
|
||||
assert src["last_attempt_at"] is not None
|
||||
assert src["last_success_at"] is not None
|
||||
assert src["last_attempt_at"] > src["last_success_at"], (
|
||||
"последняя попытка должна быть новее последнего успеха — иначе провалы " "не видны вообще"
|
||||
)
|
||||
|
||||
|
||||
def test_only_failures_means_no_success_at_all(db) -> None:
|
||||
"""Ни одного успеха → failed, а не «свежо, потому что строки пишутся»."""
|
||||
for day in range(1, 5):
|
||||
_insert_dump(db, days_ago=float(day), error="TimeoutError", features=0)
|
||||
|
||||
src = _nspd(db)
|
||||
assert src["last_success_at"] is None, "успехом сочтена строка с harvest_error"
|
||||
assert src["status"] == "failed", f"статус {src['status']!r} при полном отсутствии успехов"
|
||||
Loading…
Add table
Reference in a new issue