feat(nspd): harvest_quarter Celery task + beat + admin endpoint (#94 pt.4 part A) #111

Merged
lekss361 merged 2 commits from feat/nspd-harvest-quarter-task into main 2026-05-12 15:45:58 +00:00
3 changed files with 10 additions and 10 deletions
Showing only changes of commit aca4b443aa - Show all commits

View file

@ -302,6 +302,9 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
) )
) )
_db.commit() _db.commit()
except Exception:
_db.rollback()
raise
finally: finally:
_db.close() _db.close()
except Exception as _e: except Exception as _e:

View file

@ -18,6 +18,7 @@ UPSERT pattern:
from __future__ import annotations from __future__ import annotations
import datetime as _dt
import json import json
import logging import logging
import time import time
@ -179,7 +180,7 @@ def _upsert_dump(
"risks_count": _build_risks_count(dump), "risks_count": _build_risks_count(dump),
"total_features": dump.total_features, "total_features": dump.total_features,
"features_json": json.dumps(features_json or [], ensure_ascii=False), "features_json": json.dumps(features_json or [], ensure_ascii=False),
"layers_fetched": "{" + ",".join(dump.layers_fetched) + "}", "layers_fetched": list(dump.layers_fetched),
"fetched_at_utc": dump.fetched_at_utc, "fetched_at_utc": dump.fetched_at_utc,
"harvest_duration_ms": duration_ms, "harvest_duration_ms": duration_ms,
"harvest_error": harvest_error, "harvest_error": harvest_error,
@ -187,8 +188,6 @@ def _upsert_dump(
} }
else: else:
# Error-only row: нулевые счётчики, harvest_error filled # Error-only row: нулевые счётчики, harvest_error filled
import datetime as _dt
params = { params = {
"quarter_cad": quarter_cad, "quarter_cad": quarter_cad,
"geom_json": None, "geom_json": None,
@ -205,7 +204,7 @@ def _upsert_dump(
"risks_count": 0, "risks_count": 0,
"total_features": 0, "total_features": 0,
"features_json": "[]", "features_json": "[]",
"layers_fetched": "{}", "layers_fetched": [],
"fetched_at_utc": _dt.datetime.now(_dt.UTC).isoformat(), "fetched_at_utc": _dt.datetime.now(_dt.UTC).isoformat(),
"harvest_duration_ms": duration_ms, "harvest_duration_ms": duration_ms,
"harvest_error": harvest_error, "harvest_error": harvest_error,
@ -214,6 +213,9 @@ def _upsert_dump(
db.execute(_UPSERT_SQL, params) db.execute(_UPSERT_SQL, params)
db.commit() db.commit()
except Exception:
db.rollback()
raise
finally: finally:
db.close() db.close()
@ -360,8 +362,7 @@ def harvest_stale_quarters(
""" """
SELECT cad_number SELECT cad_number
FROM cad_quarters_geom FROM cad_quarters_geom
WHERE region_code = :rc WHERE cad_number NOT IN (
AND cad_number NOT IN (
SELECT quarter_cad SELECT quarter_cad
FROM nspd_quarter_dumps FROM nspd_quarter_dumps
WHERE region_code = :rc WHERE region_code = :rc

View file

@ -296,8 +296,6 @@ def test_harvest_stale_quarters_fanout(
# Mock DB session # Mock DB session
mock_db = MagicMock() mock_db = MagicMock()
mock_session_cls.return_value = mock_db mock_session_cls.return_value = mock_db
mock_db.__enter__ = MagicMock(return_value=mock_db)
mock_db.__exit__ = MagicMock(return_value=False)
mock_db.execute.return_value.all.return_value = [(c,) for c in stale_cads] mock_db.execute.return_value.all.return_value = [(c,) for c in stale_cads]
# harvest_quarter.apply_async — проверяем что вызван для каждого cad # harvest_quarter.apply_async — проверяем что вызван для каждого cad
@ -324,8 +322,6 @@ def test_harvest_stale_quarters_empty(
"""Нет stale кварталов → 0 enqueued, apply_async не вызван.""" """Нет stale кварталов → 0 enqueued, apply_async не вызван."""
mock_db = MagicMock() mock_db = MagicMock()
mock_session_cls.return_value = mock_db mock_session_cls.return_value = mock_db
mock_db.__enter__ = MagicMock(return_value=mock_db)
mock_db.__exit__ = MagicMock(return_value=False)
mock_db.execute.return_value.all.return_value = [] mock_db.execute.return_value.all.return_value = []
mock_task.apply_async = MagicMock() mock_task.apply_async = MagicMock()