Compare commits
1 commit
main
...
fix/2956-f
| Author | SHA1 | Date | |
|---|---|---|---|
| 660232ffdc |
3 changed files with 234 additions and 5 deletions
|
|
@ -1414,6 +1414,15 @@ class FreshnessSource(BaseModel):
|
||||||
# внутри окна и не ложно-срабатывает (#1947 fix). Default 1 → флагует только если
|
# внутри окна и не ложно-срабатывает (#1947 fix). Default 1 → флагует только если
|
||||||
# суммарный выход цикла = 0 (безопасный минимальный catch).
|
# суммарный выход цикла = 0 (безопасный минимальный catch).
|
||||||
min_output_rows: int = 1
|
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) проверены на
|
# Реестр источников. 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) больше НЕ источник истины.
|
# defunct nspd_scrape_runs (manual WAF-ban 2026-04-30) больше НЕ источник истины.
|
||||||
table="nspd_quarter_dumps",
|
table="nspd_quarter_dumps",
|
||||||
timestamp_col="fetched_at_utc",
|
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-режиме не используется — оставляем валидное имя колонки.
|
# В timestamp-режиме не используется — оставляем валидное имя колонки.
|
||||||
work_col="total_features",
|
work_col="total_features",
|
||||||
# Медленный кадастровый + lazy-refresh источник: дампы освежаются по мере
|
# Медленный кадастровый + lazy-refresh источник: дампы освежаются по мере
|
||||||
|
|
@ -1591,7 +1606,9 @@ def compute_freshness(db: Session) -> dict[str, Any]:
|
||||||
|
|
||||||
Для data-table источников (src.timestamp_col задан, напр. nspd → nspd_quarter_dumps)
|
Для data-table источников (src.timestamp_col задан, напр. nspd → nspd_quarter_dumps)
|
||||||
нет run-ledger семантики (status/started/finished отсутствуют), поэтому:
|
нет 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(*) строк, обновлённых в окне
|
- objects_updated_24h / _7d = COUNT(*) строк, обновлённых в окне
|
||||||
- last_status = NULL (косметика только для run-ledger'ов)
|
- last_status = NULL (косметика только для run-ledger'ов)
|
||||||
Остальной downstream (age_days / _classify_freshness / status-маппинг) — общий.
|
Остальной downstream (age_days / _classify_freshness / status-маппинг) — общий.
|
||||||
|
|
@ -1608,16 +1625,23 @@ def compute_freshness(db: Session) -> dict[str, Any]:
|
||||||
# все временные границы передаются параметрами (:d1/:d7).
|
# все временные границы передаются параметрами (:d1/:d7).
|
||||||
if src.timestamp_col is not None:
|
if src.timestamp_col is not None:
|
||||||
# Data-table режим: плоская контент-таблица без run-ledger семантики
|
# Data-table режим: плоская контент-таблица без run-ledger семантики
|
||||||
# (нет status/started/finished). Свежесть = MAX(timestamp_col),
|
# (нет status/started/finished). Свежесть = MAX(timestamp_col) по строкам,
|
||||||
# upd_24h/_7d = COUNT(*) строк, обновлённых в окне. last_status=NULL
|
# прошедшим success_where (если задан; иначе по всем), upd_24h/_7d =
|
||||||
# (косметика только для run-ledger'ов).
|
# COUNT(*) строк, обновлённых в окне. last_status=NULL (косметика только
|
||||||
|
# для run-ledger'ов).
|
||||||
ts = src.timestamp_col
|
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 = (
|
row = (
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
f"""
|
f"""
|
||||||
SELECT
|
SELECT
|
||||||
MAX({ts}) AS last_success_at,
|
MAX({ts}){success_filter} AS last_success_at,
|
||||||
MAX({ts}) AS last_attempt_at,
|
MAX({ts}) AS last_attempt_at,
|
||||||
COALESCE(COUNT(*) FILTER (
|
COALESCE(COUNT(*) FILTER (
|
||||||
WHERE {ts} > NOW() - CAST(:d1 AS interval)
|
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_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_itogo
|
||||||
tests/services/scrapers/test_ekb_ppt_tep_parser.py::TestParsePptTepRealPdf::test_zone_balance_has_rows
|
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