feat: bulk-дампы открытых данных ФНС по юрлицам + lookup по ИНН #3429

Merged
lekss361 merged 2 commits from feat/fns-legal-entities into main 2026-09-08 22:29:11 +00:00
7 changed files with 1081 additions and 0 deletions

View file

@ -0,0 +1,58 @@
"""ФНС opendata lookup по ИНН — тонкое чтение `fns_legal_entity_facts`.
CONTEXT: см. `fns_opendata_loader.py`. Этот модуль единственная точка чтения
загруженных фактов. ПОТРЕБИТЕЛЯ У НЕГО ПОКА НЕТ: ничто в продукте (скоринг
застройщика/УК, estimator и т.п.) сюда не ходит подключение решается отдельной
задачей вне этого PR.
psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`.
"""
from __future__ import annotations
from collections import defaultdict
from datetime import date
from typing import TypedDict
from sqlalchemy import text
from sqlalchemy.orm import Session
_LOOKUP_SQL = text(
"""
SELECT dataset, series, period, value, org_name
FROM fns_legal_entity_facts
WHERE inn = CAST(:inn AS text)
ORDER BY dataset, series, period
"""
)
class FnsFact(TypedDict):
series: str
period: date
value: float
org_name: str | None
def get_facts_by_inn(db: Session, inn: str) -> dict[str, list[FnsFact]]:
"""Факты по ИНН, сгруппированные по dataset (`revexp`/`sshr2019`/`debtam`/`snr`).
Пустой/пробельный ИНН и отсутствие данных `{}` (не исключение вызывающий код
не обязан оборачивать lookup в try/except ради нормального «нет данных»).
"""
normalized = (inn or "").strip()
if not normalized:
return {}
rows = db.execute(_LOOKUP_SQL, {"inn": normalized}).mappings().all()
out: dict[str, list[FnsFact]] = defaultdict(list)
for row in rows:
out[row["dataset"]].append(
{
"series": row["series"],
"period": row["period"],
"value": row["value"],
"org_name": row["org_name"],
}
)
return dict(out)

View file

@ -0,0 +1,497 @@
"""ФНС opendata loader: доходы/расходы, ССЧ, недоимка, спецрежимы юрлиц → lookup по ИНН.
CONTEXT: открытых данных ЕГРН по правообладателям-физлицам не существует (218-ФЗ
ст. 62) это жёсткий блокер. По ЮРЛИЦАМ данные открыты, лицензия ФНС
(nalog.gov.ru/opendata) разрешает переработку и перераспространение легальный обход
для того среза, где он в принципе доступен. Наборы (маска slug'а: 7707329152-<slug>):
revexp доходы и расходы;
sshr2019 среднесписочная численность;
debtam недоимка и задолженность;
snr спецрежимы.
ЕГРЮЛ-данных здесь НЕТ: ни адреса, ни ОКВЭД, ни учредителей только ИНН + наименование
+ показатели набора.
ПОТРЕБИТЕЛЯ У ЭТОГО СЕРВИСА ПОКА НЕТ. Этот модуль (+ fns_lookup.py) только
загрузка и lookup по ИНН. Подключение к продукту (например, скоринг застройщика/УК
в оценке) сюда сознательно не входит и не реализовано это решается отдельной
задачей вне текущего PR. `estimator.py` не тронут.
ИСТОЧНИК: https://www.nalog.gov.ru/opendata/7707329152-<slug>/ HTML-каталог набора.
Прямая ссылка на .zip вида
https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip
резолвится ДИНАМИЧЕСКИ парсом href со страницы каталога (`resolve_dataset_file_url`):
дата в имени файла меняется с каждой публикацией набора хардкодить URL нельзя,
протухнет на следующей публикации. Внутри .zip множество XML-файлов.
ПАРСИНГ XML generic-харвестер, а не хардкод конкретных тегов набора. Точная схема
атрибутов XSD structure-20180110 варьируется по набору и НЕ была вживую сверена в
этом окружении (нет сетевого доступа к file.nalog.ru отсюда). Вместо этого: любой
XML-элемент, несущий ИНН-подобный атрибут (ИННЮЛ/ИНН данные ФНС opendata, как и
ГАР/ФИАС, лежат в атрибутах элементов, не в тексте, см. gar_flats_loader.py), даёт
identity строки; ЛЮБОЙ его собственный числовой атрибут (кроме ИНН/КПП/ОКПО/ОКТМО/
ОКВЭД/ОГРН и атрибутов-имени) становится отдельным фактом (series=имя атрибута,
value=число). Это устойчиво к точным названиям показателей набора ценой чуть более
широкого набора series, чем «официальный» словарь показателей ПЕРЕД первым боевым
прогоном на реальном .zip стоит свериться с фактическим XML и, если харвестер тянет
лишнее (например служебные коды), сузить `_ID_ATTRS_EXCLUDE`.
ПАМЯТЬ: .zip целиком лежит в памяти как байты (httpx возвращает `.content` иначе
не проверить целостность архива до распаковки), но НИ ОДИН XML-член НЕ
распаковывается в память/на диск целиком: каждый открывается потоково через
`zf.open(name)` и парсится `lxml.etree.iterparse` с очисткой обработанных элементов
(тот же приём, что в gar_flats_loader.py, там на многогигабайтных standalone XML,
здесь на множестве XML внутри одного архива, открываемых по одному).
TLS: nalog.gov.ru/file.nalog.ru отдают сертификат НУЦ Минцифры httpx с default
trust store не верифицирует verify=False. Открытые данные, без auth/PII.
Прецеденты: sber_index.py, domrf_kapremont_loader.py.
psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`.
"""
from __future__ import annotations
import io
import logging
import re
import zipfile
from collections.abc import Iterator
from dataclasses import dataclass
from datetime import date
from urllib.parse import urljoin
import httpx
from lxml import etree
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
CATALOG_URL_TEMPLATE = "https://www.nalog.gov.ru/opendata/7707329152-{slug}/"
DATASET_SLUGS: tuple[str, ...] = ("revexp", "sshr2019", "debtam", "snr")
DOWNLOAD_TIMEOUT_SEC = 300
UPSERT_CHUNK_SIZE = 500
# Атрибуты-идентификаторы/классификаторы — не показатели, исключаем из харвеста.
_INN_ATTRS = ("ИННЮЛ", "ИНН")
_NAME_ATTRS = ("НаимОрг", "НаимОрганизации", "НаимОрганизацииПолн", "НаимЮЛПолн")
_PERIOD_ATTRS = ("ДатаСост",)
# Служебные атрибуты, которые парсятся как числа, но фактом не являются:
# идентификаторы, коды и даты документа.
_ID_ATTRS_EXCLUDE = frozenset(
{
"ИННЮЛ",
"ИНН",
"КПП",
"ОКПО",
"ОКТМО",
"ОКВЭД",
"ОГРН",
"ИдДок",
"ИдФайл",
"ДатаДок",
"ДатаСост",
"ВерсФорм",
"ВерсПрог",
"КолДок",
"ТипИнф",
*_NAME_ATTRS,
}
)
_RECORD_TAGS = frozenset({"Документ", "Док", "СвЮЛ", "Сведения"})
"""Теги, на которых собирается запись. Реальные выгрузки ФНС используют `Документ`;
остальные запас под соседние наборы. Ограничение обязательно: иначе `Файл`
переоткрывал бы уже собранные документы и удваивал записи."""
_ZIP_HREF_RE = re.compile(r'href="([^"]+\.zip)"', re.IGNORECASE)
_PUBLISH_DATE_RE = re.compile(r"data-(\d{8})-")
# ─────────────────────────────────────────────────────────────────────────────
# Резолв ссылки со страницы каталога (URL меняется на каждой публикации)
# ─────────────────────────────────────────────────────────────────────────────
def resolve_dataset_file_url(html: str, *, base_url: str) -> tuple[str, date | None]:
"""Ссылка на .zip актуальной версии набора, снятая со страницы каталога.
Ищет href, оканчивающийся на .zip и несущий токен даты `data-YYYYMMDD-` эта
дата меняется с каждой публикацией (см. docstring модуля), поэтому URL нельзя
хардкодить. Если на странице несколько подходящих ссылок (несколько версий
структуры/публикаций), берёт с МАКСИМАЛЬНОЙ датой детерминированно, без сети.
Возвращает (абсолютный_url, published_on). published_on None, если ссылка
найдена, но токен даты не распарсился (защитный случай, не должен происходить
для ссылок, прошедших регэксп даты).
Raises ValueError, если на странице нет ни одной ссылки вида *data-YYYYMMDD-*.zip.
"""
candidates: list[tuple[str, date]] = []
for m in _ZIP_HREF_RE.finditer(html):
href = m.group(1)
date_m = _PUBLISH_DATE_RE.search(href)
if not date_m:
continue
token = date_m.group(1)
try:
published = date(int(token[:4]), int(token[4:6]), int(token[6:8]))
except ValueError:
continue
candidates.append((href, published))
if not candidates:
raise ValueError(f"каталог {base_url}: не найдено ссылок вида href=*data-YYYYMMDD-*.zip")
href, published = max(candidates, key=lambda c: c[1])
return urljoin(base_url, href), published
# ─────────────────────────────────────────────────────────────────────────────
# In-memory модель факта
# ─────────────────────────────────────────────────────────────────────────────
@dataclass(slots=True)
class FnsRecord:
inn: str
dataset: str
series: str
period: date
value: float
org_name: str | None
def to_params(self) -> dict[str, object]:
return {
"inn": self.inn,
"dataset": self.dataset,
"series": self.series,
"period": self.period,
"value": self.value,
"org_name": self.org_name,
}
# ─────────────────────────────────────────────────────────────────────────────
# Стриминговый парсер XML внутри ZIP (lxml iterparse + очистка, без extractall)
# ─────────────────────────────────────────────────────────────────────────────
def _find_attr(elem: etree._Element, candidates: tuple[str, ...]) -> str | None:
for attr in candidates:
raw = elem.get(attr)
if raw:
stripped = raw.strip()
if stripped:
return stripped
return None
def _parse_numeric(raw: str) -> float | None:
stripped = raw.strip()
if not stripped:
return None
try:
return float(stripped.replace(",", "."))
except ValueError:
return None
def _find_attr_deep(elem: etree._Element, candidates: tuple[str, ...]) -> str | None:
"""Найти атрибут на самом элементе ИЛИ на любом его потомке.
В реальных выгрузках ФНС ИНН и показатели лежат на РАЗНЫХ соседних элементах
внутри `<Документ>` см. докстринг `harvest_element`.
"""
found = _find_attr(elem, candidates)
if found is not None:
return found
for child in elem.iterdescendants():
found = _find_attr(child, candidates)
if found is not None:
return found
return None
def _period_from_element(elem: etree._Element, fallback: date) -> date:
"""Отчётная дата документа из `ДатаСост` (ДД.ММ.ГГГГ), иначе — переданная."""
raw = _find_attr_deep(elem, _PERIOD_ATTRS)
if raw is None:
return fallback
try:
day, month, year = (int(part) for part in raw.split("."))
return date(year, month, day)
except (ValueError, TypeError):
return fallback
def _numeric_attrs(node: etree._Element, seen: set[str]) -> list[tuple[str, float]]:
"""Числовые не-служебные атрибуты узла, без повторов по имени серии."""
out: list[tuple[str, float]] = []
for key, raw in node.attrib.items():
if key in _ID_ATTRS_EXCLUDE or key in seen:
continue
value = _parse_numeric(raw)
if value is None:
continue
seen.add(key)
out.append((key, value))
return out
def harvest_element(elem: etree._Element, *, dataset: str, period: date) -> list[FnsRecord]:
"""Один XML-элемент → 0..N FnsRecord (один на каждый числовой не-id атрибут).
Реальная форма выгрузки (проверено на data-20260825 набора revexp, 2118 XML
в архиве, 40 002 факта на первых 20 001 организаций)::
<Документ ИдДок="..." ДатаДок="25.08.2026" ДатаСост="31.12.2025">
<СведНП НаимОрг="ООО ..." ИННЮЛ="4205406898"/>
<СведДохРасх СумДоход="341864000.00" СумРасход="282224000.00"/>
</Документ>
То есть ИНН живёт на `СведНП`, а показатели на СОСЕДНЕМ элементе. Поэтому
носители ИНН ищутся по всему поддереву, а значения без собственного ИНН
(`СведДохРасх` и подобные) привязываются к организации ТОЛЬКО когда носитель в
поддереве один иначе непонятно, чьи это цифры, и мы их не выдумываем.
Отчётный период берётся из `ДатаСост` документа, если он есть; переданный
`period` запасное значение.
Чистая функция (без БД/сети). Поддерево без ИНН даёт пустой список.
"""
nodes = [elem, *elem.iterdescendants()]
holders = [n for n in nodes if _find_attr(n, _INN_ATTRS) is not None]
if not holders:
return []
doc_period = _period_from_element(elem, period)
shared = [n for n in nodes if n not in holders] if len(holders) == 1 else []
out: list[FnsRecord] = []
for holder in holders:
inn = _find_attr(holder, _INN_ATTRS)
if inn is None: # pragma: no cover - отфильтровано выше
continue
org_name = _find_attr(holder, _NAME_ATTRS) or (
_find_attr_deep(elem, _NAME_ATTRS) if len(holders) == 1 else None
)
seen: set[str] = set()
for node in (holder, *shared):
for series, value in _numeric_attrs(node, seen):
out.append(
FnsRecord(
inn=inn,
dataset=dataset,
series=series,
period=doc_period,
value=value,
org_name=org_name,
)
)
return out
def _iter_xml_records(fh: object, *, dataset: str, period: date) -> Iterator[FnsRecord]:
"""Стримит FnsRecord из одного XML-потока `fh` (открытый член ZIP или файл).
`events=("end",)` + `elem.clear()` + срез предыдущих сиблингов bounded memory
на произвольно большом XML (приём из gar_flats_loader._stream_rows).
"""
context = etree.iterparse(
fh, events=("end",), recover=True, huge_tree=True, resolve_entities=False
)
for _event, elem in context:
if not isinstance(elem.tag, str) or elem.tag not in _RECORD_TAGS:
continue
yield from harvest_element(elem, dataset=dataset, period=period)
# Чистим ТОЛЬКО на границе записи. `end`-события детей приходят раньше
# родительского, поэтому безусловный `elem.clear()` на каждом элементе
# вычищал `<СведНП>`/`<СведДохРасх>` ДО закрытия `<Документ>` — и парсер
# молча извлекал ноль записей из реального дампа.
elem.clear()
parent = elem.getparent()
if parent is not None:
while elem.getprevious() is not None:
del parent[0]
del context
def iter_dataset_records(zip_bytes: bytes, *, dataset: str, period: date) -> Iterator[FnsRecord]:
"""Обходит все *.xml внутри `zip_bytes`, элемент за элементом, без extractall.
Каждый XML-член открывается через `zf.open(name)` как поток (НЕ `zf.read()`
целиком, НЕ `zf.extractall()`) см. docstring модуля § ПАМЯТЬ.
Raises ValueError, если в архиве нет ни одного .xml.
"""
with zipfile.ZipFile(io.BytesIO(zip_bytes)) as zf:
names = [n for n in zf.namelist() if n.lower().endswith(".xml")]
if not names:
raise ValueError(f"zip набора {dataset} не содержит .xml записей: {zf.namelist()!r}")
for name in names:
with zf.open(name) as fh:
yield from _iter_xml_records(fh, dataset=dataset, period=period)
# ─────────────────────────────────────────────────────────────────────────────
# UPSERT фактов + версия набора
# ─────────────────────────────────────────────────────────────────────────────
_UPSERT_FACTS_SQL = text(
"""
INSERT INTO fns_legal_entity_facts (inn, dataset, series, period, value, org_name, loaded_at)
VALUES (
CAST(:inn AS text), CAST(:dataset AS text), CAST(:series AS text),
CAST(:period AS date), CAST(:value AS numeric), CAST(:org_name AS text), now()
)
ON CONFLICT (inn, dataset, series, period) DO UPDATE SET
value = EXCLUDED.value,
org_name = EXCLUDED.org_name,
loaded_at = now()
WHERE fns_legal_entity_facts.value IS DISTINCT FROM EXCLUDED.value
"""
)
_SELECT_VERSION_SQL = text(
"SELECT file_url, published_on FROM fns_dataset_versions WHERE dataset = CAST(:dataset AS text)"
)
_UPSERT_VERSION_SQL = text(
"""
INSERT INTO fns_dataset_versions (dataset, file_url, published_on, loaded_at)
VALUES (CAST(:dataset AS text), CAST(:file_url AS text), CAST(:published_on AS date), now())
ON CONFLICT (dataset) DO UPDATE SET
file_url = EXCLUDED.file_url,
published_on = EXCLUDED.published_on,
loaded_at = now()
"""
)
def _chunks(items: list[FnsRecord], size: int) -> Iterator[list[FnsRecord]]:
for i in range(0, len(items), size):
yield items[i : i + size]
def upsert_records(
db: Session, records: list[FnsRecord], *, chunk_size: int = UPSERT_CHUNK_SIZE
) -> int:
"""Батчевый UPSERT в fns_legal_entity_facts. НЕ коммитит (коммитит caller).
SAVEPOINT на каждый батч сбойный батч откатывается изолированно, остальные
доезжают (тот же приём, что и gar_flats_loader.upsert_gar_houses).
"""
upserted = 0
failed_batches = 0
for batch in _chunks(records, chunk_size):
params = [r.to_params() for r in batch]
try:
with db.begin_nested():
db.execute(_UPSERT_FACTS_SQL, params)
upserted += len(batch)
except Exception:
failed_batches += 1
logger.warning(
"fns_opendata upsert: батч из %d строк сбойнул (пропущен)",
len(batch),
exc_info=True,
)
if failed_batches:
logger.warning("fns_opendata upsert: сбойных батчей=%d", failed_batches)
return upserted
def get_loaded_version(db: Session, dataset: str) -> tuple[str, date | None] | None:
"""Уже загруженная версия набора (file_url, published_on) или None."""
row = db.execute(_SELECT_VERSION_SQL, {"dataset": dataset}).first()
if row is None:
return None
return row[0], row[1]
def record_dataset_version(
db: Session, dataset: str, file_url: str, published_on: date | None
) -> None:
"""UPSERT версии набора. НЕ коммитит (коммитит caller)."""
db.execute(
_UPSERT_VERSION_SQL,
{"dataset": dataset, "file_url": file_url, "published_on": published_on},
)
# ─────────────────────────────────────────────────────────────────────────────
# Orchestration
# ─────────────────────────────────────────────────────────────────────────────
def _download(url: str, *, client: httpx.Client) -> bytes:
resp = client.get(url, timeout=DOWNLOAD_TIMEOUT_SEC)
resp.raise_for_status()
return resp.content
def load_dataset(
db: Session,
slug: str,
*,
client: httpx.Client | None = None,
dry_run: bool = False,
force: bool = False,
chunk_size: int = UPSERT_CHUNK_SIZE,
) -> dict[str, int | str]:
"""Резолвит актуальную ссылку набора `slug` → скачивает .zip → парсит → UPSERT.
Пропускает скачивание, если `fns_dataset_versions` уже содержит ТУ ЖЕ ссылку
(force=True форсирует перекачку). dry_run резолв ссылки происходит (проверяем
каталог доступен), скачивание/парс/запись нет. `client` для тестов
(httpx.MockTransport); если не передан, открывается и закрывается свой.
НЕ коммитит коммитит caller.
"""
if slug not in DATASET_SLUGS:
raise ValueError(f"неизвестный slug набора: {slug!r}, ожидались {DATASET_SLUGS}")
owns_client = client is None
if client is None:
# TLS: см. docstring модуля § TLS — verify=False, открытые данные без auth/PII.
client = httpx.Client(timeout=DOWNLOAD_TIMEOUT_SEC, verify=False)
try:
catalog_url = CATALOG_URL_TEMPLATE.format(slug=slug)
resp = client.get(catalog_url, timeout=DOWNLOAD_TIMEOUT_SEC)
resp.raise_for_status()
file_url, published_on = resolve_dataset_file_url(resp.text, base_url=catalog_url)
existing = get_loaded_version(db, slug)
if not force and existing is not None and existing[0] == file_url:
logger.info("fns_opendata %s: версия %s уже загружена, skip", slug, file_url)
return {"dataset": slug, "skipped": 1, "records": 0, "upserted": 0}
if dry_run:
logger.info(
"fns_opendata %s: dry_run — резолвлен %s (published_on=%s), скачивание пропущено",
slug,
file_url,
published_on,
)
return {"dataset": slug, "skipped": 0, "records": 0, "upserted": 0}
zip_bytes = _download(file_url, client=client)
period = published_on or date.today()
records = list(iter_dataset_records(zip_bytes, dataset=slug, period=period))
upserted = upsert_records(db, records, chunk_size=chunk_size)
record_dataset_version(db, slug, file_url, published_on)
result: dict[str, int | str] = {
"dataset": slug,
"skipped": 0,
"records": len(records),
"upserted": upserted,
}
logger.info("fns_opendata load %s DONE (dry_run=%s): %s", slug, dry_run, result)
return result
finally:
if owns_client:
client.close()
def load_datasets(
db: Session,
slugs: tuple[str, ...] = DATASET_SLUGS,
*,
client: httpx.Client | None = None,
dry_run: bool = False,
force: bool = False,
) -> dict[str, dict[str, int | str]]:
"""load_dataset для каждого slug из `slugs`. НЕ коммитит между наборами (caller)."""
return {
slug: load_dataset(db, slug, client=client, dry_run=dry_run, force=force) for slug in slugs
}

View file

@ -674,6 +674,53 @@ async def _job_domrf_kapremont_load(
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
# ── fns_opendata_load — bulk-дампы ФНС по юрлицам → fns_legal_entity_facts ────
# Открытые данные ФНС по ЮРЛИЦАМ (лицензия nalog.gov.ru/opendata разрешает
# переработку/перераспространение) — легальный обход того, что открытых данных
# ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62). Потребителя у
# fns_legal_entity_facts на момент добавления НЕТ (см. docstring
# app/services/fns_opendata_loader.py и fns_lookup.py) — расписание сидируется
# enabled=false, включение отдельным осознанным шагом.
async def _job_fns_opendata_load(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
"""Скачать наборы ФНС opendata (revexp/sshr2019/debtam/snr) → fns_legal_entity_facts.
Тело переиспользует load_dataset (тот же дизайн-инвариант модуля, что у
_job_domrf_kapremont_load: не дублируем логику CLI). commit() после каждого
набора сбой на debtam не должен откатывать уже загруженный revexp.
"""
from app.services.fns_opendata_loader import DATASET_SLUGS, load_dataset
datasets = params.get("datasets") or list(DATASET_SLUGS)
def _run() -> dict[str, int]:
total_records = 0
total_upserted = 0
for slug in datasets:
result = load_dataset(db, slug)
db.commit()
total_records += int(result["records"])
total_upserted += int(result["upserted"])
return {
"records": total_records,
"upserted": total_upserted,
# см. докстринг _job_domrf_kapremont_load: выделенные колонки прогона +
# гейт «три подряд нулевых прогона».
"total_seen": total_records,
"new_count": total_upserted,
}
loop = asyncio.get_event_loop()
try:
counters = await loop.run_in_executor(None, _run)
ctx.runs.mark_done(db, run_id, counters)
except Exception as exc:
logger.exception("scheduler: fns_opendata_load crashed run_id=%d", run_id)
db.rollback()
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
# ── frt_mkd_load — АИС ППК ФРТ, реестр МКД region 66 (issue #frt-mkd) ────────
async def _job_frt_mkd_load(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
@ -929,6 +976,7 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]:
"domrf_kapremont_load": Handler(_job_domrf_kapremont_load, "domrf_kapremont_load"),
"frt_mkd_load": Handler(_job_frt_mkd_load, "frt_mkd_load"),
"cbr_macro_pull": Handler(_job_cbr_macro_pull, "cbr_macro_pull"),
"fns_opendata_load": Handler(_job_fns_opendata_load, "fns_opendata_load"),
"purge_expired_trade_in_data": Handler(
_job_purge_expired_trade_in_data, "purge_expired_trade_in_data"
),

View file

@ -0,0 +1,71 @@
"""CLI: ФНС opendata (revexp/sshr2019/debtam/snr) → fns_legal_entity_facts.
Только загрузка + lookup (`app/services/fns_lookup.py`) потребителя у данных пока
нет, см. docstring `app/services/fns_opendata_loader.py`.
Запуск из контейнера tradein-backend/tradein-scraper (нужен интернет к nalog.gov.ru):
python -m app.tasks.fns_opendata_load # все 4 набора
python -m app.tasks.fns_opendata_load --datasets revexp # один набор
python -m app.tasks.fns_opendata_load --dry-run # без записи, только резолв ссылки
python -m app.tasks.fns_opendata_load --force # перекачать, даже если версия та же
В режиме --dry-run скачивание/парс/запись не происходят только резолв актуальной
ссылки со страницы каталога (проверка доступности + логирование того, что было бы
скачано).
"""
from __future__ import annotations
import argparse
import logging
from app.core.db import SessionLocal
from app.services.fns_opendata_loader import DATASET_SLUGS, load_dataset
logger = logging.getLogger(__name__)
def build_parser() -> argparse.ArgumentParser:
"""Парсер CLI (вынесен для тестируемости флагов без запуска main)."""
parser = argparse.ArgumentParser(
description="ФНС opendata loader: revexp/sshr2019/debtam/snr → fns_legal_entity_facts"
)
parser.add_argument(
"--datasets",
nargs="+",
choices=list(DATASET_SLUGS),
default=list(DATASET_SLUGS),
help="список наборов (по умолчанию — все 4)",
)
parser.add_argument(
"--dry-run", action="store_true", help="без записи в БД — только резолв ссылки/подсчёт"
)
parser.add_argument(
"--force", action="store_true", help="перекачать, даже если версия совпадает с загруженной"
)
return parser
def main() -> None:
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
)
args = build_parser().parse_args()
db = SessionLocal()
try:
results: dict[str, object] = {}
for slug in args.datasets:
result = load_dataset(db, slug, dry_run=args.dry_run, force=args.force)
if not args.dry_run:
db.commit()
results[slug] = result
logger.info("fns_opendata_load DONE: dry_run=%s %s", args.dry_run, results)
finally:
db.close()
if __name__ == "__main__":
main()

View file

@ -0,0 +1,71 @@
-- 296_fns_legal_entity_facts.sql
-- ФНС open data по юрлицам: доходы/расходы, среднесписочная численность, недоимка,
-- спецрежимы (issue: bulk-дампы ФНС → lookup по ИНН, PR "feat/fns-legal-entities").
--
-- ПОЧЕМУ ЭТОТ ИСТОЧНИК
-- Открытых данных ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62).
-- По ЮРЛИЦАМ данные открыты, лицензия ФНС (nalog.gov.ru/opendata) разрешает
-- переработку и перераспределение. Наборы (все по маске 7707329152-<slug>):
-- revexp (доходы и расходы), sshr2019 (среднесписочная численность), debtam
-- (недоимка и задолженность), snr (спецрежимы). ЕГРЮЛ (адрес/ОКВЭД/учредители) в
-- этих наборах НЕТ — только ИНН + наименование + показатели набора.
--
-- ЭТОТ PR — ТОЛЬКО ЗАГРУЗКА И LOOKUP ПО ИНН. Потребителя (скоринг застройщика/УК в
-- оценке) здесь нет и не планируется в рамках этого PR — подключение отдельной
-- задачей, если решение об этом будет принято отдельно.
--
-- ФОРМА ХРАНЕНИЯ: одна узкая long-format таблица на все 4 набора, а не таблица на
-- набор. Наборы различаются только НАБОРОМ показателей (атрибутов XML), сама форма
-- «ИНН + показатель + период + значение» у всех одинаковая — long-format не плодит
-- 4 почти идентичные таблицы и не требует миграции при добавлении 5-го набора в
-- будущем (не в этом PR).
--
-- PRIMARY KEY (inn, dataset, series, period), А НЕ (inn, series, period) — так
-- предлагал скелет задачи, но series здесь = имя XML-атрибута источника, и без
-- dataset в ключе одноимённые атрибуты в двух разных наборах молча схлопнутся в
-- одну строку (мы НЕ проверяли по XSD, что имена атрибутов между revexp/sshr2019/
-- debtam/snr не пересекаются, поэтому закладываемся на худший случай).
--
-- Идемпотентно: CREATE TABLE/INDEX IF NOT EXISTS. lock_timeout — на случай, если
-- таблицы уже существуют от прогона на другой ветке (ALTER не планируется, но
-- дисциплина сохраняется).
BEGIN;
SET LOCAL lock_timeout = '5s';
CREATE TABLE IF NOT EXISTS fns_legal_entity_facts (
inn text NOT NULL,
dataset text NOT NULL, -- 'revexp' | 'sshr2019' | 'debtam' | 'snr'
series text NOT NULL, -- имя показателя (атрибут XML источника)
period date NOT NULL, -- период показателя; см. fns_opendata_loader.py
value numeric NOT NULL,
org_name text, -- наименование юрлица на момент загрузки (не ЕГРЮЛ, provenance)
loaded_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (inn, dataset, series, period)
);
COMMENT ON TABLE fns_legal_entity_facts IS
'ФНС opendata (nalog.gov.ru/opendata/7707329152-*) по юрлицам: показатели наборов '
'revexp/sshr2019/debtam/snr. Long-format: одна строка = один показатель одного ИНН '
'за один период. ЕГРЮЛ-данных (адрес/ОКВЭД/учредители) здесь нет. Без потребителя '
'на момент создания — см. app/services/fns_lookup.py.';
COMMENT ON COLUMN fns_legal_entity_facts.inn IS 'ИНН юрлица (10 цифр), как в XML источника';
COMMENT ON COLUMN fns_legal_entity_facts.dataset IS 'slug набора ФНС opendata: revexp/sshr2019/debtam/snr';
COMMENT ON COLUMN fns_legal_entity_facts.series IS 'имя показателя = имя XML-атрибута источника (не переименовываем)';
COMMENT ON COLUMN fns_legal_entity_facts.period IS 'период показателя: дата публикации набора (см. fns_dataset_versions.published_on)';
COMMENT ON COLUMN fns_legal_entity_facts.value IS 'значение показателя как есть из XML (numeric — источник не различает валюту/единицы в атрибуте)';
COMMENT ON COLUMN fns_legal_entity_facts.org_name IS 'наименование юрлица на момент загрузки, provenance only — НЕ авторитетный источник (используй ЕГРЮЛ)';
CREATE TABLE IF NOT EXISTS fns_dataset_versions (
dataset text PRIMARY KEY, -- 'revexp' | 'sshr2019' | 'debtam' | 'snr'
file_url text NOT NULL, -- прямая ссылка на .zip, снятая со страницы каталога
published_on date, -- дата публикации, распарсенная из имени файла (data-YYYYMMDD-...)
loaded_at timestamptz NOT NULL DEFAULT now()
);
COMMENT ON TABLE fns_dataset_versions IS
'Последняя загруженная версия каждого набора ФНС opendata — гейт от перекачки '
'уже загруженной версии (URL меняется с каждой публикацией, дата в имени файла).';
COMMENT ON COLUMN fns_dataset_versions.file_url IS 'резолвится динамически со страницы каталога набора — НЕ хардкодить, протухает';
COMMENT ON COLUMN fns_dataset_versions.published_on IS 'дата публикации из имени файла (data-YYYYMMDD-structure-...)';
COMMIT;

View file

@ -0,0 +1,36 @@
-- 297_scrape_schedules_seed_fns_opendata_load.sql
-- Расписание для fns_opendata_load (см. 296_fns_legal_entity_facts.sql — контекст источника).
--
-- enabled=false: как domclick_detail_backfill (мигр. 175) — включение отдельным
-- осознанным шагом после деплоя, дымовой пробы и ручного --dry-run прогона. Наборы
-- ФНС весят десятки-сотни МБ каждый, первый прогон надо смотреть глазами.
--
-- interval_days=30: наборы ФНС публикуются раз в квартал/год (revexp — по итогам
-- отчётного периода), суточный/недельный опрос каталога бессмысленно частый.
-- Реализовано через default_params (interval_minutes читает post_claim
-- reschedule_after_minutes, здесь используем ту же схему в днях: 30 * 24 * 60).
--
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)), 294 (таблицы факта).
-- Идемпотентно: ON CONFLICT (source) DO NOTHING.
BEGIN;
INSERT INTO scrape_schedules (
source,
enabled,
window_start_hour,
window_end_hour,
next_run_at,
default_params
)
VALUES
(
'fns_opendata_load',
false,
2,
5,
((CURRENT_DATE + INTERVAL '1 day')) AT TIME ZONE 'UTC',
'{"interval_minutes": 43200, "datasets": ["revexp", "sshr2019", "debtam", "snr"]}'::jsonb
)
ON CONFLICT (source) DO NOTHING;
COMMIT;

View file

@ -0,0 +1,300 @@
"""Тесты ФНС opendata loader'а + lookup (app/services/fns_opendata_loader.py,
app/services/fns_lookup.py, мигр. 296/295, "bulk-дампы ФНС по юрлицам").
Coverage:
- resolve_dataset_file_url резолв .zip-ссылки из литерального куска HTML каталога
(несколько ссылок берём максимальную дату публикации); ValueError, если ссылок нет.
- harvest_element generic-харвестер: ИНН-подобный атрибут даёт identity, числовые
атрибуты (кроме id/name-подобных) становятся FnsRecord; элемент без ИНН [].
- iter_dataset_records round-trip через реально собранный .zip с мини-XML фикстурой.
- load_dataset HTTP замокан через httpx.MockTransport (каталог + zip); dry_run не
пишет; повторная загрузка той же версии skip без сети на .zip.
- Статические asserts по SQL: CAST-дисциплина, ON CONFLICT, IS DISTINCT FROM.
- upsert_records: `db.begin_nested()` в исходнике, модуль НЕ коммитит (коммитит caller).
- fns_lookup.get_facts_by_inn на MagicMock db: группировка по dataset, пустой ИНН.
"""
from __future__ import annotations
import inspect
import io
import os
import re
import zipfile
from datetime import date
from unittest.mock import MagicMock
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
import httpx
import pytest
from lxml import etree
from app.services import fns_lookup as fl
from app.services import fns_opendata_loader as fol
# ─────────────────────────────────────────────────────────────────────────────
# resolve_dataset_file_url
# ─────────────────────────────────────────────────────────────────────────────
_CATALOG_HTML = """
<html><body>
<table>
<tr><td><a href="https://file.nalog.ru/opendata/7707329152-revexp/data-20260101-structure-20180110.zip">2026-01-01</a></td></tr>
<tr><td><a href="https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip">2026-08-25</a></td></tr>
<tr><td><a href="/opendata/7707329152-revexp/structure-20180110.xsd">XSD</a></td></tr>
</table>
</body></html>
"""
def test_resolve_dataset_file_url_picks_latest_publication() -> None:
url, published = fol.resolve_dataset_file_url(
_CATALOG_HTML, base_url="https://www.nalog.gov.ru/opendata/7707329152-revexp/"
)
assert (
url
== "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip"
)
assert published == date(2026, 8, 25)
def test_resolve_dataset_file_url_relative_href_resolved_against_base() -> None:
html = '<a href="data-20260301-structure-1.zip">v</a>'
url, published = fol.resolve_dataset_file_url(
html, base_url="https://www.nalog.gov.ru/opendata/7707329152-snr/"
)
assert url == "https://www.nalog.gov.ru/opendata/7707329152-snr/data-20260301-structure-1.zip"
assert published == date(2026, 3, 1)
def test_resolve_dataset_file_url_raises_when_no_zip_links() -> None:
with pytest.raises(ValueError, match="data-YYYYMMDD"):
fol.resolve_dataset_file_url("<html><body>empty</body></html>", base_url="https://x/")
# ─────────────────────────────────────────────────────────────────────────────
# harvest_element — generic-харвестер
# ─────────────────────────────────────────────────────────────────────────────
def test_harvest_element_extracts_numeric_attrs_excluding_ids() -> None:
elem = etree.fromstring(
'<СвНП ИННЮЛ="7707329152" КПП="770701001" НаимОрг="ООО РОМАШКА" '
'СумДоход="1234567.89" СумРасход="654321,00"/>'
)
records = fol.harvest_element(elem, dataset="revexp", period=date(2026, 8, 25))
by_series = {r.series: r for r in records}
assert set(by_series) == {"СумДоход", "СумРасход"}
assert by_series["СумДоход"].inn == "7707329152"
assert by_series["СумДоход"].value == pytest.approx(1234567.89)
assert by_series["СумРасход"].value == pytest.approx(654321.00) # запятая — decimal separator
assert by_series["СумДоход"].org_name == "ООО РОМАШКА"
assert by_series["СумДоход"].dataset == "revexp"
assert by_series["СумДоход"].period == date(2026, 8, 25)
def test_harvest_element_without_inn_returns_empty() -> None:
elem = etree.fromstring('<Прочее КПП="770701001" Значение="1"/>')
assert fol.harvest_element(elem, dataset="revexp", period=date(2026, 1, 1)) == []
def test_harvest_element_skips_non_numeric_attrs() -> None:
elem = etree.fromstring('<СвНП ИННЮЛ="123" Комментарий="текст, не число"/>')
assert fol.harvest_element(elem, dataset="snr", period=date(2026, 1, 1)) == []
# ─────────────────────────────────────────────────────────────────────────────
# iter_dataset_records — round-trip через собранный .zip
# ─────────────────────────────────────────────────────────────────────────────
def _zip_bytes(xml_by_name: dict[str, str]) -> bytes:
buf = io.BytesIO()
with zipfile.ZipFile(buf, "w") as zf:
for name, content in xml_by_name.items():
zf.writestr(name, content)
return buf.getvalue()
_SAMPLE_XML = """<?xml version="1.0" encoding="UTF-8"?>
<Файл>
<Документ>
<СвНП ИННЮЛ="7707329152" НаимОрг="ООО РОМАШКА" СумДоход="1000"/>
<СвНП ИННЮЛ="6660000001" НаимОрг="ООО ВАСИЛЁК" СумДоход="2000" СумРасход="500"/>
</Документ>
</Файл>
"""
def test_iter_dataset_records_round_trip_from_zip() -> None:
data = _zip_bytes({"revexp_66.xml": _SAMPLE_XML})
records = list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25)))
assert len(records) == 3 # 1000 + (2000, 500)
inns = {r.inn for r in records}
assert inns == {"7707329152", "6660000001"}
assert all(r.dataset == "revexp" and r.period == date(2026, 8, 25) for r in records)
def test_iter_dataset_records_raises_when_zip_has_no_xml() -> None:
data = _zip_bytes({"readme.txt": "no xml here"})
with pytest.raises(ValueError, match="\\.xml"):
list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 1, 1)))
# ─────────────────────────────────────────────────────────────────────────────
# load_dataset — HTTP замокан (никакой живой сети)
# ─────────────────────────────────────────────────────────────────────────────
def _make_client(zip_bytes: bytes, *, file_url: str, catalog_html: str) -> httpx.Client:
def handler(request: httpx.Request) -> httpx.Response:
if str(request.url) == file_url:
return httpx.Response(200, content=zip_bytes)
if "opendata/7707329152-" in str(request.url):
return httpx.Response(200, text=catalog_html)
return httpx.Response(404)
return httpx.Client(transport=httpx.MockTransport(handler))
def test_load_dataset_dry_run_does_not_write(monkeypatch: pytest.MonkeyPatch) -> None:
file_url = "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-1.zip"
catalog_html = f'<a href="{file_url}">zip</a>'
client = _make_client(b"unused", file_url=file_url, catalog_html=catalog_html)
db = MagicMock()
db.execute.return_value.first.return_value = None # get_loaded_version → None
result = fol.load_dataset(db, "revexp", client=client, dry_run=True)
assert result == {"dataset": "revexp", "skipped": 0, "records": 0, "upserted": 0}
# dry_run: единственный execute — SELECT версии (get_loaded_version), UPSERT не звался.
assert db.execute.call_count == 1
def test_load_dataset_skips_when_version_unchanged() -> None:
file_url = "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-1.zip"
catalog_html = f'<a href="{file_url}">zip</a>'
client = _make_client(b"should-not-be-fetched", file_url=file_url, catalog_html=catalog_html)
db = MagicMock()
db.execute.return_value.first.return_value = (file_url, date(2026, 8, 25))
result = fol.load_dataset(db, "revexp", client=client)
assert result == {"dataset": "revexp", "skipped": 1, "records": 0, "upserted": 0}
def test_load_dataset_unknown_slug_raises() -> None:
db = MagicMock()
with pytest.raises(ValueError, match="неизвестный slug"):
fol.load_dataset(db, "not-a-real-dataset")
# ─────────────────────────────────────────────────────────────────────────────
# Статические asserts по SQL (без БД)
# ─────────────────────────────────────────────────────────────────────────────
def test_upsert_facts_sql_uses_cast_never_double_colon() -> None:
sql = re.sub(r"\s+", " ", str(fol._UPSERT_FACTS_SQL.text))
assert not re.search(r":\w+::", sql)
assert "CAST(:inn AS text)" in sql
assert "CAST(:period AS date)" in sql
assert "CAST(:value AS numeric)" in sql
assert "ON CONFLICT (inn, dataset, series, period) DO UPDATE SET" in sql
assert "IS DISTINCT FROM" in sql
def test_upsert_version_sql_uses_cast_and_conflict_target() -> None:
sql = re.sub(r"\s+", " ", str(fol._UPSERT_VERSION_SQL.text))
assert not re.search(r":\w+::", sql)
assert "ON CONFLICT (dataset) DO UPDATE SET" in sql
def test_upsert_records_uses_savepoint_and_module_does_not_commit() -> None:
src = inspect.getsource(fol.upsert_records)
assert "with db.begin_nested():" in src
module_src = inspect.getsource(fol)
assert "db.commit()" not in module_src
# ─────────────────────────────────────────────────────────────────────────────
# fns_lookup.get_facts_by_inn — MagicMock db
# ─────────────────────────────────────────────────────────────────────────────
def test_get_facts_by_inn_groups_by_dataset() -> None:
db = MagicMock()
org = "ООО РОМАШКА"
p2026 = date(2026, 8, 25)
def row(dataset: str, series: str, period: date, value: int) -> dict[str, object]:
return {
"dataset": dataset,
"series": series,
"period": period,
"value": value,
"org_name": org,
}
db.execute.return_value.mappings.return_value.all.return_value = [
row("revexp", "СумДоход", p2026, 1000),
row("revexp", "СумРасход", p2026, 500),
row("sshr2019", "ССЧ", date(2019, 1, 1), 42),
]
result = fl.get_facts_by_inn(db, "7707329152")
assert set(result) == {"revexp", "sshr2019"}
assert len(result["revexp"]) == 2
assert result["sshr2019"][0]["series"] == "ССЧ"
def test_get_facts_by_inn_empty_inn_returns_empty_without_query() -> None:
db = MagicMock()
assert fl.get_facts_by_inn(db, "") == {}
assert fl.get_facts_by_inn(db, " ") == {}
db.execute.assert_not_called()
def test_get_facts_by_inn_no_rows_returns_empty_dict() -> None:
db = MagicMock()
db.execute.return_value.mappings.return_value.all.return_value = []
assert fl.get_facts_by_inn(db, "0000000000") == {}
_REAL_SHAPE_XML = """<?xml version="1.0" encoding="UTF-8"?>
<Файл ИдФайл="VO_OTKRDAN_5_9965_9965_20260825_cf1c08a2" ВерсФорм="4.01" \
ТипИнф="ОТКРДАННЫЕ5" КолДок="2">
<ИдОтпр><ФИООтв Фамилия="_" Имя="_"/></ИдОтпр>
<Документ ИдДок="dea28452-05f3-42bf-a08f-7a648c593473" ДатаДок="25.08.2026" \
ДатаСост="31.12.2025">
<СведНП НаимОрг="ОБЩЕСТВО С ОГРАНИЧЕННОЙ ОТВЕТСТВЕННОСТЬЮ &quot;СПК&quot;" \
ИННЮЛ="4205406898"/>
<СведДохРасх СумДоход="341864000.00" СумРасход="282224000.00"/>
</Документ>
<Документ ИдДок="28be840b-5db9-422b-b80c-3b678b808c96" ДатаДок="25.08.2026" \
ДатаСост="31.12.2025">
<СведНП НаимОрг="ООО ВАСИЛЁК" ИННЮЛ="6660000001"/>
<СведДохРасх СумДоход="1000.00" СумРасход="500.00"/>
</Документ>
</Файл>
"""
def test_iter_dataset_records_parses_real_fns_document_shape() -> None:
"""Боевая форма ФНС: ИНН на `СведНП`, показатели на СОСЕДНЕМ `СведДохРасх`.
Регресс на реальный дефект: харвестер требовал ИНН и числа на ОДНОМ элементе,
поэтому на настоящем дампе (revexp, data-20260825, 2118 XML) извлекал НОЛЬ
записей, а тесты на выдуманной фикстуре этого не показывали. Плюс безусловный
`elem.clear()` вычищал детей до закрытия `<Документ>`.
"""
data = _zip_bytes({"revexp_1.xml": _REAL_SHAPE_XML})
records = list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25)))
assert len(records) == 4 # 2 организации × (СумДоход, СумРасход)
assert {r.inn for r in records} == {"4205406898", "6660000001"}
assert {r.series for r in records} == {"СумДоход", "СумРасход"}
# период берётся из ДатаСост документа, а не из переданного fallback'а
assert {r.period for r in records} == {date(2025, 12, 31)}
first = next(r for r in records if r.inn == "4205406898" and r.series == "СумДоход")
assert first.value == 341864000.0
assert first.org_name is not None and "СПК" in first.org_name
def test_service_ids_are_not_harvested_as_facts() -> None:
"""ВерсФорм/КолДок/ДатаДок — служебные, фактами быть не должны."""
data = _zip_bytes({"revexp_1.xml": _REAL_SHAPE_XML})
series = {
r.series for r in fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25))
}
assert not series & {"ВерсФорм", "КолДок", "ДатаДок", "ДатаСост", "ИННЮЛ"}