feat(site-finder): live CBR key-rate scraper → macro_indicator (#945 PR B)
Site Finder v2 / GG-форсайт. Deterministic, no LLM. Fills the key_rate series (empty after PR A) into macro_indicator from the live CBR source. - services/scrapers/cbr_macro.py: fetch_key_rate via CBR SOAP DailyInfoWebServ KeyRateXML(fromDate, ToDate) (POST soap+xml; ToDate capital-T + dateTime suffix required — flat GET returns empty). Pure parse_key_rate_xml (TZ-safe date, empty/malformed → []). Live-verified: 85 rows, current rate 14.5%. - workers/tasks/cbr_macro_sync.py: upsert (key_rate, rf, daily, %) into macro_indicator ON CONFLICT DO UPDATE; SAVEPOINT per row; fetch-fail → raise. - beat: cbr-macro-sync-weekly (Mon 05:30 МСК, Europe/Moscow tz). - 15 tests (pure parse + mocked upsert). No migration (table from #963 PR A).
This commit is contained in:
parent
dbae4b0bda
commit
ea42d29a4a
6 changed files with 680 additions and 0 deletions
282
backend/app/services/scrapers/cbr_macro.py
Normal file
282
backend/app/services/scrapers/cbr_macro.py
Normal file
|
|
@ -0,0 +1,282 @@
|
||||||
|
"""CBR (Банк России) macro indicators scraper — key rate history.
|
||||||
|
|
||||||
|
Fills the ``key_rate`` series of the ``macro_indicator`` table (region='rf')
|
||||||
|
from the live CBR SOAP service ``DailyInfoWebServ``. Deterministic, no LLM.
|
||||||
|
|
||||||
|
Endpoint / method (verified live 2026-06-02):
|
||||||
|
POST https://www.cbr.ru/DailyInfoWebServ/DailyInfo.asmx
|
||||||
|
SOAP 1.2 action ``KeyRateXML(fromDate, ToDate)``.
|
||||||
|
|
||||||
|
⚠ Параметры называются ``fromDate`` и ``ToDate`` (вторая — с ЗАГЛАВНОЙ T,
|
||||||
|
так в WSDL). Оба типа ``xsd:dateTime`` — нужен суффикс времени
|
||||||
|
``T00:00:00``; голый ``YYYY-MM-DD`` сервис молча игнорит и возвращает
|
||||||
|
пустой ``KeyRateXMLResponse``. (Поэтому плоский GET-form
|
||||||
|
``/KeyRate?fromDate=..&toDate=..`` тоже не работает — у DailyInfoWebServ нет
|
||||||
|
HTTP-GET binding для этого метода, только SOAP.)
|
||||||
|
|
||||||
|
Response shape (``KeyRateXML`` — без diffgram-обёртки, в отличие от ``KeyRate``):
|
||||||
|
<KeyRateXMLResult>
|
||||||
|
<KeyRate>
|
||||||
|
<KR><DT>2023-08-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>
|
||||||
|
...
|
||||||
|
</KeyRate>
|
||||||
|
</KeyRateXMLResult>
|
||||||
|
Строки идут в обратном хронологическом порядке; ``parse_key_rate_xml``
|
||||||
|
сортировку не навязывает — отдаёт в порядке появления.
|
||||||
|
|
||||||
|
Парсинг вынесен в чистую ``parse_key_rate_xml`` (offline-тестируемую на
|
||||||
|
фикстуре); HTTP-обвязка ``fetch_key_rate`` тонкая (httpx, timeout, UA, retry).
|
||||||
|
|
||||||
|
Standalone smoke (для ручной проверки против живого cbr.ru):
|
||||||
|
python -m app.services.scrapers.cbr_macro
|
||||||
|
|
||||||
|
Mirror conventions: httpx.Client + явный User-Agent (как objective.py /
|
||||||
|
ekburg_permits.py), logger вместо print, типизированные сигнатуры.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
from datetime import date, datetime, timedelta
|
||||||
|
from decimal import Decimal, InvalidOperation
|
||||||
|
from xml.etree import ElementTree as ET
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# ── constants ────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
CBR_DAILYINFO_URL = "https://www.cbr.ru/DailyInfoWebServ/DailyInfo.asmx"
|
||||||
|
# SOAP method that returns a clean <KeyRate><KR>... structure (no ADO diffgram).
|
||||||
|
CBR_KEYRATE_METHOD = "KeyRateXML"
|
||||||
|
|
||||||
|
USER_AGENT = "GenDesign/1.0 (+https://gendsgn.ru) macro indicators scraper"
|
||||||
|
DEFAULT_TIMEOUT_S = 30.0
|
||||||
|
DEFAULT_RETRIES = 3
|
||||||
|
|
||||||
|
# Полный бэкфилл истории ключевой ставки. ON CONFLICT в таске делает повторные
|
||||||
|
# прогоны идемпотентными, так что диапазон можно держать широким.
|
||||||
|
DEFAULT_FROM_DATE = date(2019, 1, 1)
|
||||||
|
|
||||||
|
|
||||||
|
class CBRScraperError(RuntimeError):
|
||||||
|
"""Сетевая / протокольная ошибка обращения к DailyInfoWebServ."""
|
||||||
|
|
||||||
|
|
||||||
|
# ── SOAP envelope ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _soap_envelope(from_date: date, to_date: date) -> str:
|
||||||
|
"""Собрать SOAP 1.2 envelope для ``KeyRateXML(fromDate, ToDate)``.
|
||||||
|
|
||||||
|
Даты сериализуются как ``YYYY-MM-DDT00:00:00`` — сервис принимает только
|
||||||
|
``xsd:dateTime`` (см. module docstring). ``ToDate`` — с заглавной T.
|
||||||
|
"""
|
||||||
|
return (
|
||||||
|
'<?xml version="1.0" encoding="utf-8"?>'
|
||||||
|
'<soap12:Envelope xmlns:soap12="http://www.w3.org/2003/05/soap-envelope">'
|
||||||
|
"<soap12:Body>"
|
||||||
|
f'<{CBR_KEYRATE_METHOD} xmlns="http://web.cbr.ru/">'
|
||||||
|
f"<fromDate>{from_date.isoformat()}T00:00:00</fromDate>"
|
||||||
|
f"<ToDate>{to_date.isoformat()}T00:00:00</ToDate>"
|
||||||
|
f"</{CBR_KEYRATE_METHOD}>"
|
||||||
|
"</soap12:Body>"
|
||||||
|
"</soap12:Envelope>"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── pure parse (unit-testable offline) ──────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _parse_dt(raw: str) -> date | None:
|
||||||
|
"""Распарсить ``<DT>`` (``2023-08-15T00:00:00+03:00``) в ``date``.
|
||||||
|
|
||||||
|
Берём только календарную дату; tz-суффикс игнорируем (ключевая ставка —
|
||||||
|
дневной показатель, время всегда полночь МСК). Поддерживаем и формы без
|
||||||
|
суффикса / с 'Z'.
|
||||||
|
"""
|
||||||
|
raw = (raw or "").strip()
|
||||||
|
if not raw:
|
||||||
|
return None
|
||||||
|
# Дата всегда в первых 10 символах ISO-8601 (YYYY-MM-DD...). Это устойчиво к
|
||||||
|
# любому tz-суффиксу (+03:00 / Z / отсутствию) без зависимости от Python <3.11
|
||||||
|
# fromisoformat ограничений.
|
||||||
|
head = raw[:10]
|
||||||
|
try:
|
||||||
|
return datetime.strptime(head, "%Y-%m-%d").date()
|
||||||
|
except ValueError:
|
||||||
|
logger.warning("CBR KeyRate: не удалось распарсить DT=%r", raw)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _parse_rate(raw: str) -> float | None:
|
||||||
|
"""Распарсить ``<Rate>`` в float. Принимаем и запятую как десятичный
|
||||||
|
разделитель (на всякий случай — фактически CBR отдаёт точку)."""
|
||||||
|
raw = (raw or "").strip().replace(",", ".")
|
||||||
|
if not raw:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
# Через Decimal — чтобы не словить артефакты float-парсинга строки.
|
||||||
|
return float(Decimal(raw))
|
||||||
|
except (InvalidOperation, ValueError):
|
||||||
|
logger.warning("CBR KeyRate: не удалось распарсить Rate=%r", raw)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _localname(tag: str) -> str:
|
||||||
|
"""ElementTree-тег без namespace-префикса: ``{ns}KR`` → ``KR``."""
|
||||||
|
return tag.rsplit("}", 1)[-1] if "}" in tag else tag
|
||||||
|
|
||||||
|
|
||||||
|
def parse_key_rate_xml(xml: str) -> list[tuple[date, float]]:
|
||||||
|
"""Pure-парсер SOAP-ответа CBR KeyRate → отсортированный список (date, rate).
|
||||||
|
|
||||||
|
Находит все элементы ``<KR>`` где угодно в дереве (устойчиво и к чистому
|
||||||
|
``KeyRateXML``, и к diffgram-варианту ``KeyRate``, и к разным namespace'ам).
|
||||||
|
Из каждого берёт дочерние ``<DT>`` и ``<Rate>``.
|
||||||
|
|
||||||
|
Возвращает список ``(date, rate)``, отсортированный по возрастанию даты.
|
||||||
|
Пустой / битый / без KR XML → ``[]`` (не бросает) — единственное исключение
|
||||||
|
логируется как warning. Дубли по дате схлопываются (последний выигрывает).
|
||||||
|
"""
|
||||||
|
if not xml or not xml.strip():
|
||||||
|
return []
|
||||||
|
try:
|
||||||
|
root = ET.fromstring(xml)
|
||||||
|
except ET.ParseError as e:
|
||||||
|
logger.warning("CBR KeyRate: XML parse error: %s", e)
|
||||||
|
return []
|
||||||
|
|
||||||
|
by_date: dict[date, float] = {}
|
||||||
|
for kr in root.iter():
|
||||||
|
if _localname(kr.tag) != "KR":
|
||||||
|
continue
|
||||||
|
dt_val: date | None = None
|
||||||
|
rate_val: float | None = None
|
||||||
|
for child in kr:
|
||||||
|
name = _localname(child.tag)
|
||||||
|
if name == "DT":
|
||||||
|
dt_val = _parse_dt(child.text or "")
|
||||||
|
elif name == "Rate":
|
||||||
|
rate_val = _parse_rate(child.text or "")
|
||||||
|
if dt_val is not None and rate_val is not None:
|
||||||
|
by_date[dt_val] = rate_val
|
||||||
|
|
||||||
|
return sorted(by_date.items())
|
||||||
|
|
||||||
|
|
||||||
|
# ── HTTP fetch (thin) ───────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def fetch_key_rate(
|
||||||
|
from_date: date | None = None,
|
||||||
|
to_date: date | None = None,
|
||||||
|
*,
|
||||||
|
timeout_s: float = DEFAULT_TIMEOUT_S,
|
||||||
|
retries: int = DEFAULT_RETRIES,
|
||||||
|
) -> list[tuple[date, Decimal]]:
|
||||||
|
"""Загрузить историю ключевой ставки ЦБ за период [from_date, to_date].
|
||||||
|
|
||||||
|
Args:
|
||||||
|
from_date: начало периода. По умолчанию ``DEFAULT_FROM_DATE`` (2019-01-01)
|
||||||
|
— полный бэкфилл истории.
|
||||||
|
to_date: конец периода. По умолчанию сегодня (UTC-date достаточно;
|
||||||
|
ставка дневная).
|
||||||
|
timeout_s: httpx timeout на запрос.
|
||||||
|
retries: число повторов на сетевых / 5xx ошибках (экспоненциальный backoff).
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Отсортированный по дате список ``(date, Decimal)``. Decimal — чтобы
|
||||||
|
отдать в numeric-колонку без потери точности.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
CBRScraperError: при сетевой ошибке / не-200 / SOAP fault после всех
|
||||||
|
повторов. Caller (Celery task) логирует и пробрасывает — surfaces
|
||||||
|
в GlitchTip.
|
||||||
|
"""
|
||||||
|
from_date = from_date or DEFAULT_FROM_DATE
|
||||||
|
to_date = to_date or datetime.now().date()
|
||||||
|
|
||||||
|
body = _soap_envelope(from_date, to_date)
|
||||||
|
headers = {
|
||||||
|
"Content-Type": "application/soap+xml; charset=utf-8",
|
||||||
|
"User-Agent": USER_AGENT,
|
||||||
|
}
|
||||||
|
|
||||||
|
last_exc: Exception | None = None
|
||||||
|
with httpx.Client(timeout=timeout_s, headers=headers) as client:
|
||||||
|
for attempt in range(retries + 1):
|
||||||
|
try:
|
||||||
|
resp = client.post(CBR_DAILYINFO_URL, content=body.encode("utf-8"))
|
||||||
|
except httpx.HTTPError as e:
|
||||||
|
last_exc = e
|
||||||
|
logger.warning(
|
||||||
|
"CBR KeyRate fetch network error (attempt %d/%d): %s",
|
||||||
|
attempt + 1,
|
||||||
|
retries + 1,
|
||||||
|
e,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
if resp.status_code == 200:
|
||||||
|
rows = parse_key_rate_xml(resp.text)
|
||||||
|
logger.info(
|
||||||
|
"CBR KeyRate fetched %d rows for [%s..%s]",
|
||||||
|
len(rows),
|
||||||
|
from_date.isoformat(),
|
||||||
|
to_date.isoformat(),
|
||||||
|
)
|
||||||
|
# Decimal(str(float)) — нормализуем через строку, чтобы
|
||||||
|
# numeric получил аккуратное значение (12.0 → '12.0').
|
||||||
|
return [(d, Decimal(str(v))) for d, v in rows]
|
||||||
|
# 5xx — повторяем; 4xx — нет смысла.
|
||||||
|
last_exc = CBRScraperError(
|
||||||
|
f"CBR DailyInfoWebServ HTTP {resp.status_code}: {resp.text[:300]}"
|
||||||
|
)
|
||||||
|
if resp.status_code < 500:
|
||||||
|
break
|
||||||
|
logger.warning(
|
||||||
|
"CBR KeyRate HTTP %d (attempt %d/%d)",
|
||||||
|
resp.status_code,
|
||||||
|
attempt + 1,
|
||||||
|
retries + 1,
|
||||||
|
)
|
||||||
|
# backoff перед следующей попыткой (не на последней итерации)
|
||||||
|
if attempt < retries:
|
||||||
|
import time
|
||||||
|
|
||||||
|
time.sleep(min(2**attempt, 10))
|
||||||
|
|
||||||
|
raise CBRScraperError(f"CBR KeyRate fetch failed after {retries + 1} attempts: {last_exc}")
|
||||||
|
|
||||||
|
|
||||||
|
# ── national mortgage weighted rate (region='rf') ───────────────────────────────
|
||||||
|
# TODO(#945): национальная средневзвешенная ставка по ипотеке (region='rf').
|
||||||
|
# Намеренно НЕ реализовано в этом PR. У ЦБ ряд есть, но стабильного
|
||||||
|
# машиночитаемого endpoint без авторизации / Excel-парсинга не нашлось
|
||||||
|
# (формы 0409316 публикуются как .xlsx с плавающей структурой листов —
|
||||||
|
# рискованно для детерминированного scraper). key_rate — must-have, закрыт.
|
||||||
|
# Ипотека deferrable: можно добавить отдельным fetch_*-методом по образцу
|
||||||
|
# ekburg_permits.parse_xlsx, когда определимся со стабильным источником.
|
||||||
|
|
||||||
|
|
||||||
|
# ── standalone smoke ────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _main() -> None:
|
||||||
|
"""``python -m app.services.scrapers.cbr_macro`` — печатает несколько свежих
|
||||||
|
строк ключевой ставки для ручной верификации против живого cbr.ru."""
|
||||||
|
logging.basicConfig(level=logging.INFO, format="%(levelname)s %(name)s: %(message)s")
|
||||||
|
# Узкий диапазон (последние ~120 дней) — быстрый smoke без полного бэкфилла.
|
||||||
|
today = datetime.now().date()
|
||||||
|
rows = fetch_key_rate(from_date=today - timedelta(days=120), to_date=today)
|
||||||
|
logger.info("KeyRate rows fetched: %d", len(rows))
|
||||||
|
# Печатаем последние 5 (самые свежие) — sorted asc, поэтому хвост.
|
||||||
|
for d, v in rows[-5:]:
|
||||||
|
logger.info(" %s -> %s%%", d.isoformat(), v)
|
||||||
|
if not rows:
|
||||||
|
logger.warning("No rows returned — проверь endpoint/период.")
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
_main()
|
||||||
|
|
@ -291,4 +291,18 @@ def build_beat_schedule() -> dict:
|
||||||
"options": {"queue": "celery"},
|
"options": {"queue": "celery"},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
# CBR макро-показатели (ключевая ставка) → macro_indicator (#945 PR B,
|
||||||
|
# GG-форсайт). Ставка меняется редко (8 заседаний ЦБ в год) — еженедельного
|
||||||
|
# прогона с запасом достаточно. ON CONFLICT делает прогон идемпотентным
|
||||||
|
# (полный бэкфилл с 2019-01-01 каждый раз; дёшево — ~1 SOAP-запрос + ~1500
|
||||||
|
# upsert'ов). Понедельник 05:30 МСК (Celery conf.timezone=Europe/Moscow, т.е.
|
||||||
|
# crontab трактуется в МСК) — после publish итогов пятничного заседания. Задача
|
||||||
|
# техническая, не управляется через job_settings — поэтому добавлена сюда (как
|
||||||
|
# poi-sync / refresh-analytics). Оператор может перенести расписание правкой crontab.
|
||||||
|
schedule["cbr-macro-sync-weekly"] = {
|
||||||
|
"task": "tasks.cbr_macro_sync.cbr_macro_sync",
|
||||||
|
"schedule": _parse_cron("30 5 * * mon"),
|
||||||
|
"options": {"queue": "celery"},
|
||||||
|
}
|
||||||
|
|
||||||
return schedule
|
return schedule
|
||||||
|
|
|
||||||
|
|
@ -60,6 +60,7 @@ celery_app = Celery(
|
||||||
"app.workers.tasks.pzz_sync",
|
"app.workers.tasks.pzz_sync",
|
||||||
"app.workers.tasks.scrape_cadastre",
|
"app.workers.tasks.scrape_cadastre",
|
||||||
"app.workers.tasks.ekburg_permits_sync",
|
"app.workers.tasks.ekburg_permits_sync",
|
||||||
|
"app.workers.tasks.cbr_macro_sync",
|
||||||
],
|
],
|
||||||
)
|
)
|
||||||
celery_app.conf.timezone = "Europe/Moscow"
|
celery_app.conf.timezone = "Europe/Moscow"
|
||||||
|
|
|
||||||
115
backend/app/workers/tasks/cbr_macro_sync.py
Normal file
115
backend/app/workers/tasks/cbr_macro_sync.py
Normal file
|
|
@ -0,0 +1,115 @@
|
||||||
|
"""Celery task: синхронизация макро-показателей ЦБ (ключевая ставка) в
|
||||||
|
``macro_indicator`` (#945 PR B, GG-форсайт / Site Finder v2).
|
||||||
|
|
||||||
|
Тянет историю ключевой ставки через ``fetch_key_rate`` (живой SOAP CBR
|
||||||
|
DailyInfoWebServ) и апсертит каждую (дату, ставку) в ``macro_indicator`` как
|
||||||
|
ряд ``indicator_type='key_rate'``, ``region='rf'``, ``source='cbr'``.
|
||||||
|
|
||||||
|
Детерминированно, без LLM. ON CONFLICT делает повторные прогоны идемпотентными
|
||||||
|
(re-run обновляет value + updated_at). Расписание — еженедельно (ключевая
|
||||||
|
ставка меняется редко), регистрируется в beat_schedule.py.
|
||||||
|
|
||||||
|
Mirror conventions: ``SessionLocal()`` + try/finally close, ``logger`` (не
|
||||||
|
print), SAVEPOINT per-row в цикле upsert (backend.md), ``CAST(:x AS type)`` —
|
||||||
|
никогда ``:x::type`` (psycopg v3). На fetch-фейле — логируем и пробрасываем
|
||||||
|
(surfaces в Celery/GlitchTip), не глотаем.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
from datetime import date
|
||||||
|
from decimal import Decimal
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.core.db import SessionLocal
|
||||||
|
from app.services.scrapers.cbr_macro import fetch_key_rate
|
||||||
|
from app.workers.celery_app import celery_app
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# psycopg v3: CAST(:x AS type) — НИКОГДА :x::type (SQLAlchemy+psycopg3 роняет
|
||||||
|
# синтаксис на ::). Контракт колонок совпадает с macro_indicator (migration 123):
|
||||||
|
# (indicator_type, region, obs_date, value, source, frequency, unit, comment,
|
||||||
|
# updated_at) PK (indicator_type, region, obs_date).
|
||||||
|
UPSERT_KEY_RATE_SQL = text(
|
||||||
|
"""
|
||||||
|
INSERT INTO macro_indicator (
|
||||||
|
indicator_type, region, obs_date, value,
|
||||||
|
source, frequency, unit, comment
|
||||||
|
) VALUES (
|
||||||
|
'key_rate', 'rf', CAST(:d AS date), CAST(:v AS numeric),
|
||||||
|
'cbr', 'daily', '%', 'CBR key rate'
|
||||||
|
)
|
||||||
|
ON CONFLICT (indicator_type, region, obs_date) DO UPDATE SET
|
||||||
|
value = EXCLUDED.value,
|
||||||
|
updated_at = now()
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _upsert_key_rate(db: Session, rows: list[tuple[date, Decimal]]) -> int:
|
||||||
|
"""Апсертит (date, rate) в macro_indicator. SAVEPOINT per-row, чтобы один
|
||||||
|
битый ряд не откатывал всю транзакцию. Возвращает число успешных upsert'ов."""
|
||||||
|
upserted = 0
|
||||||
|
for d, v in rows:
|
||||||
|
try:
|
||||||
|
with db.begin_nested():
|
||||||
|
db.execute(UPSERT_KEY_RATE_SQL, {"d": d, "v": v})
|
||||||
|
upserted += 1
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("upsert key_rate %s=%s failed: %s", d, v, e)
|
||||||
|
db.commit()
|
||||||
|
return upserted
|
||||||
|
|
||||||
|
|
||||||
|
@celery_app.task(
|
||||||
|
bind=True,
|
||||||
|
name="tasks.cbr_macro_sync.cbr_macro_sync",
|
||||||
|
max_retries=2,
|
||||||
|
)
|
||||||
|
def cbr_macro_sync(
|
||||||
|
self: Any,
|
||||||
|
from_date: str | None = None,
|
||||||
|
to_date: str | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""Загрузить историю ключевой ставки ЦБ и апсертить в macro_indicator.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
from_date: ISO-дата начала периода ('YYYY-MM-DD'). None → дефолт
|
||||||
|
scraper'а (2019-01-01, полный бэкфилл).
|
||||||
|
to_date: ISO-дата конца периода. None → сегодня.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Счётчики ``{"fetched": N, "upserted": M, "from_date": ..., "to_date": ...}``.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
Пробрасывает любую ошибку fetch'а (CBRScraperError и пр.) — чтобы фейл
|
||||||
|
был виден в Celery / GlitchTip, а не молча проглочен.
|
||||||
|
"""
|
||||||
|
fd = date.fromisoformat(from_date) if from_date else None
|
||||||
|
td = date.fromisoformat(to_date) if to_date else None
|
||||||
|
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
rows = fetch_key_rate(from_date=fd, to_date=td)
|
||||||
|
fetched = len(rows)
|
||||||
|
logger.info("cbr_macro_sync: fetched %d key_rate rows", fetched)
|
||||||
|
|
||||||
|
upserted = _upsert_key_rate(db, rows)
|
||||||
|
logger.info("cbr_macro_sync: upserted %d/%d key_rate rows", upserted, fetched)
|
||||||
|
|
||||||
|
return {
|
||||||
|
"fetched": fetched,
|
||||||
|
"upserted": upserted,
|
||||||
|
"from_date": fd.isoformat() if fd else None,
|
||||||
|
"to_date": td.isoformat() if td else None,
|
||||||
|
}
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception("cbr_macro_sync failed: %s", e)
|
||||||
|
raise
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
131
backend/tests/scrapers/test_cbr_macro.py
Normal file
131
backend/tests/scrapers/test_cbr_macro.py
Normal file
|
|
@ -0,0 +1,131 @@
|
||||||
|
"""Тесты pure-парсера CBR KeyRate (offline, без живой сети).
|
||||||
|
|
||||||
|
Покрывают: парсинг (date, value), обработку tz-суффикса, сортировку,
|
||||||
|
пустой / битый / без-KR XML → []. Фикстуры — из реального формата
|
||||||
|
KeyRateXML-ответа DailyInfoWebServ (verified live 2026-06-02).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from datetime import date
|
||||||
|
|
||||||
|
from app.services.scrapers.cbr_macro import parse_key_rate_xml
|
||||||
|
|
||||||
|
# Реальный shape KeyRateXML (без diffgram). 3 строки, descending по дате,
|
||||||
|
# tz-суффикс +03:00, точка как десятичный разделитель.
|
||||||
|
FIXTURE_KEYRATE_XML = (
|
||||||
|
'<?xml version="1.0" encoding="utf-8"?>'
|
||||||
|
'<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">'
|
||||||
|
"<soap:Body>"
|
||||||
|
'<KeyRateXMLResponse xmlns="http://web.cbr.ru/">'
|
||||||
|
"<KeyRateXMLResult>"
|
||||||
|
'<KeyRate xmlns="">'
|
||||||
|
"<KR><DT>2023-09-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>"
|
||||||
|
"<KR><DT>2023-08-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>"
|
||||||
|
"<KR><DT>2023-08-14T00:00:00+03:00</DT><Rate>8.50</Rate></KR>"
|
||||||
|
"</KeyRate>"
|
||||||
|
"</KeyRateXMLResult>"
|
||||||
|
"</KeyRateXMLResponse>"
|
||||||
|
"</soap:Body>"
|
||||||
|
"</soap:Envelope>"
|
||||||
|
)
|
||||||
|
|
||||||
|
# Diffgram-вариант (метод KeyRate) — парсер должен его тоже понимать (находит KR
|
||||||
|
# где угодно), несмотря на schema-блок и namespace на KR.
|
||||||
|
FIXTURE_KEYRATE_DIFFGRAM = (
|
||||||
|
'<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">'
|
||||||
|
"<soap:Body><KeyRateResponse xmlns=\"http://web.cbr.ru/\"><KeyRateResult>"
|
||||||
|
'<diffgr:diffgram xmlns:msdata="urn:schemas-microsoft-com:xml-msdata"'
|
||||||
|
' xmlns:diffgr="urn:schemas-microsoft-com:xml-diffgram-v1">'
|
||||||
|
'<KeyRate xmlns="">'
|
||||||
|
'<KR diffgr:id="KR1" msdata:rowOrder="0">'
|
||||||
|
"<DT>2023-09-29T00:00:00+03:00</DT><Rate>13.00</Rate></KR>"
|
||||||
|
'<KR diffgr:id="KR2" msdata:rowOrder="1">'
|
||||||
|
"<DT>2023-09-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>"
|
||||||
|
"</KeyRate></diffgr:diffgram></KeyRateResult></KeyRateResponse>"
|
||||||
|
"</soap:Body></soap:Envelope>"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_basic_rows() -> None:
|
||||||
|
"""3 KR-строки → 3 кортежа (date, float), отсортированы по возрастанию даты."""
|
||||||
|
rows = parse_key_rate_xml(FIXTURE_KEYRATE_XML)
|
||||||
|
assert rows == [
|
||||||
|
(date(2023, 8, 14), 8.5),
|
||||||
|
(date(2023, 8, 15), 12.0),
|
||||||
|
(date(2023, 9, 15), 12.0),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_returns_sorted_ascending() -> None:
|
||||||
|
"""Вход descending — выход обязан быть ascending по дате."""
|
||||||
|
rows = parse_key_rate_xml(FIXTURE_KEYRATE_XML)
|
||||||
|
dates = [d for d, _ in rows]
|
||||||
|
assert dates == sorted(dates)
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_strips_timezone_suffix() -> None:
|
||||||
|
"""DT с tz-суффиксом '+03:00' парсится в чистую дату (без сдвига)."""
|
||||||
|
rows = parse_key_rate_xml(FIXTURE_KEYRATE_XML)
|
||||||
|
assert (date(2023, 8, 15), 12.0) in rows
|
||||||
|
# значение float, не строка
|
||||||
|
assert all(isinstance(v, float) for _, v in rows)
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_value_types() -> None:
|
||||||
|
"""Дробная ставка 8.50 → 8.5 (float)."""
|
||||||
|
rows = dict(parse_key_rate_xml(FIXTURE_KEYRATE_XML))
|
||||||
|
assert rows[date(2023, 8, 14)] == 8.5
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_diffgram_variant() -> None:
|
||||||
|
"""Diffgram-обёртка (метод KeyRate) тоже парсится — KR находится в любом
|
||||||
|
namespace / вложенности."""
|
||||||
|
rows = parse_key_rate_xml(FIXTURE_KEYRATE_DIFFGRAM)
|
||||||
|
assert rows == [
|
||||||
|
(date(2023, 9, 15), 12.0),
|
||||||
|
(date(2023, 9, 29), 13.0),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_empty_string() -> None:
|
||||||
|
assert parse_key_rate_xml("") == []
|
||||||
|
assert parse_key_rate_xml(" ") == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_empty_response_no_kr() -> None:
|
||||||
|
"""Валидный, но пустой KeyRateXMLResponse (период без данных) → []."""
|
||||||
|
xml = (
|
||||||
|
'<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">'
|
||||||
|
'<soap:Body><KeyRateXMLResponse xmlns="http://web.cbr.ru/" />'
|
||||||
|
"</soap:Body></soap:Envelope>"
|
||||||
|
)
|
||||||
|
assert parse_key_rate_xml(xml) == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_malformed_xml() -> None:
|
||||||
|
"""Битый XML → [] (не бросает)."""
|
||||||
|
assert parse_key_rate_xml("<KeyRate><KR><DT>2023") == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_skips_kr_with_missing_fields() -> None:
|
||||||
|
"""KR без DT или без Rate пропускается; валидные остаются."""
|
||||||
|
xml = (
|
||||||
|
"<KeyRate>"
|
||||||
|
"<KR><DT>2023-08-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>"
|
||||||
|
"<KR><DT>2023-08-16T00:00:00+03:00</DT></KR>" # нет Rate
|
||||||
|
"<KR><Rate>11.00</Rate></KR>" # нет DT
|
||||||
|
"</KeyRate>"
|
||||||
|
)
|
||||||
|
assert parse_key_rate_xml(xml) == [(date(2023, 8, 15), 12.0)]
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_dedups_by_date_last_wins() -> None:
|
||||||
|
"""Дубли по дате схлопываются — последний выигрывает."""
|
||||||
|
xml = (
|
||||||
|
"<KeyRate>"
|
||||||
|
"<KR><DT>2023-08-15T00:00:00+03:00</DT><Rate>12.00</Rate></KR>"
|
||||||
|
"<KR><DT>2023-08-15T00:00:00+03:00</DT><Rate>13.00</Rate></KR>"
|
||||||
|
"</KeyRate>"
|
||||||
|
)
|
||||||
|
assert parse_key_rate_xml(xml) == [(date(2023, 8, 15), 13.0)]
|
||||||
137
backend/tests/workers/test_cbr_macro_sync.py
Normal file
137
backend/tests/workers/test_cbr_macro_sync.py
Normal file
|
|
@ -0,0 +1,137 @@
|
||||||
|
"""Тесты Celery-таски cbr_macro_sync (mock session + mock fetch, без живой сети).
|
||||||
|
|
||||||
|
Покрывают: контракт upsert-SQL (indicator_type='key_rate', region='rf',
|
||||||
|
source='cbr', ON CONFLICT, CAST not ::), счётчики fetched/upserted, проброс
|
||||||
|
ошибки fetch'а.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from datetime import date
|
||||||
|
from decimal import Decimal
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import MagicMock, patch
|
||||||
|
|
||||||
|
|
||||||
|
def _make_mock_db() -> tuple[MagicMock, list[dict[str, Any]]]:
|
||||||
|
"""Mock DB session с поддержкой db.begin_nested() context manager.
|
||||||
|
|
||||||
|
Возвращает (db, captured_params) — captured содержит params каждого
|
||||||
|
db.execute(..., params).
|
||||||
|
"""
|
||||||
|
captured: list[dict[str, Any]] = []
|
||||||
|
|
||||||
|
def execute_router(stmt: Any, params: Any = None) -> MagicMock:
|
||||||
|
if params is not None:
|
||||||
|
captured.append(params)
|
||||||
|
return MagicMock()
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute = execute_router
|
||||||
|
# begin_nested() → context manager (SAVEPOINT)
|
||||||
|
db.begin_nested.return_value.__enter__ = MagicMock(return_value=None)
|
||||||
|
db.begin_nested.return_value.__exit__ = MagicMock(return_value=False)
|
||||||
|
return db, captured
|
||||||
|
|
||||||
|
|
||||||
|
def test_upsert_sql_contract() -> None:
|
||||||
|
"""Smoke: UPSERT_KEY_RATE_SQL содержит обязательные литералы контракта
|
||||||
|
macro_indicator и НЕ использует :x::type (psycopg v3)."""
|
||||||
|
from app.workers.tasks.cbr_macro_sync import UPSERT_KEY_RATE_SQL
|
||||||
|
|
||||||
|
sql = str(UPSERT_KEY_RATE_SQL)
|
||||||
|
upper = sql.upper()
|
||||||
|
|
||||||
|
assert "INTO MACRO_INDICATOR" in upper
|
||||||
|
# литералы контракта — в SQL они lowercase
|
||||||
|
assert "'key_rate'" in sql
|
||||||
|
assert "'rf'" in sql
|
||||||
|
assert "'cbr'" in sql
|
||||||
|
assert "'daily'" in sql
|
||||||
|
# ON CONFLICT по полному PK + обновление value/updated_at
|
||||||
|
assert "ON CONFLICT (INDICATOR_TYPE, REGION, OBS_DATE) DO UPDATE" in upper
|
||||||
|
assert "VALUE = EXCLUDED.VALUE" in upper
|
||||||
|
assert "UPDATED_AT = NOW()" in upper
|
||||||
|
# psycopg v3: CAST(:x AS type), НИКОГДА :x::type
|
||||||
|
assert "CAST(:D AS DATE)" in upper
|
||||||
|
assert "CAST(:V AS NUMERIC)" in upper
|
||||||
|
assert "::" not in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_upserts_each_row_with_correct_params() -> None:
|
||||||
|
"""Каждая (date, rate) → один execute с bind-параметрами :d и :v."""
|
||||||
|
from app.workers.tasks import cbr_macro_sync as task_mod
|
||||||
|
|
||||||
|
db, captured = _make_mock_db()
|
||||||
|
fake_rows = [
|
||||||
|
(date(2023, 8, 14), Decimal("8.5")),
|
||||||
|
(date(2023, 8, 15), Decimal("12.0")),
|
||||||
|
]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(task_mod, "SessionLocal", return_value=db),
|
||||||
|
patch.object(task_mod, "fetch_key_rate", return_value=fake_rows) as mock_fetch,
|
||||||
|
):
|
||||||
|
result = task_mod.cbr_macro_sync.run()
|
||||||
|
|
||||||
|
mock_fetch.assert_called_once()
|
||||||
|
assert result["fetched"] == 2
|
||||||
|
assert result["upserted"] == 2
|
||||||
|
# параметры каждого upsert
|
||||||
|
assert captured == [
|
||||||
|
{"d": date(2023, 8, 14), "v": Decimal("8.5")},
|
||||||
|
{"d": date(2023, 8, 15), "v": Decimal("12.0")},
|
||||||
|
]
|
||||||
|
db.commit.assert_called()
|
||||||
|
db.close.assert_called_once()
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_passes_date_range_to_fetch() -> None:
|
||||||
|
"""from_date / to_date (ISO-строки) парсятся в date и прокидываются в fetch."""
|
||||||
|
from app.workers.tasks import cbr_macro_sync as task_mod
|
||||||
|
|
||||||
|
db, _ = _make_mock_db()
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(task_mod, "SessionLocal", return_value=db),
|
||||||
|
patch.object(task_mod, "fetch_key_rate", return_value=[]) as mock_fetch,
|
||||||
|
):
|
||||||
|
result = task_mod.cbr_macro_sync.run(from_date="2020-01-01", to_date="2020-12-31")
|
||||||
|
|
||||||
|
mock_fetch.assert_called_once_with(from_date=date(2020, 1, 1), to_date=date(2020, 12, 31))
|
||||||
|
assert result["from_date"] == "2020-01-01"
|
||||||
|
assert result["to_date"] == "2020-12-31"
|
||||||
|
assert result["fetched"] == 0
|
||||||
|
assert result["upserted"] == 0
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_reraises_fetch_failure() -> None:
|
||||||
|
"""Ошибка fetch'а пробрасывается (surfaces в Celery/GlitchTip), не глотается.
|
||||||
|
DB-сессия при этом закрывается."""
|
||||||
|
from app.workers.tasks import cbr_macro_sync as task_mod
|
||||||
|
|
||||||
|
db, _ = _make_mock_db()
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(task_mod, "SessionLocal", return_value=db),
|
||||||
|
patch.object(task_mod, "fetch_key_rate", side_effect=RuntimeError("CBR down")),
|
||||||
|
):
|
||||||
|
try:
|
||||||
|
task_mod.cbr_macro_sync.run()
|
||||||
|
raise AssertionError("expected RuntimeError to propagate")
|
||||||
|
except RuntimeError as e:
|
||||||
|
assert "CBR down" in str(e)
|
||||||
|
|
||||||
|
db.close.assert_called_once()
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_registered_in_beat_weekly() -> None:
|
||||||
|
"""Задача зарегистрирована в beat как еженедельная (понедельник)."""
|
||||||
|
from app.workers.beat_schedule import build_beat_schedule
|
||||||
|
|
||||||
|
schedule = build_beat_schedule()
|
||||||
|
assert "cbr-macro-sync-weekly" in schedule
|
||||||
|
entry = schedule["cbr-macro-sync-weekly"]
|
||||||
|
assert entry["task"] == "tasks.cbr_macro_sync.cbr_macro_sync"
|
||||||
|
# crontab day_of_week = понедельник (1)
|
||||||
|
assert 1 in entry["schedule"].day_of_week
|
||||||
Loading…
Add table
Reference in a new issue