"""Regression coverage for official registry ingestion without upstream IO."""
from __future__ import annotations
import zipfile
from collections import Counter
from io import BytesIO
from tempfile import TemporaryDirectory
from types import SimpleNamespace
from unittest.mock import patch
from apps.core.models import BackgroundJob, JobStatus
from apps.parsers.budget_registry import BUDGET_SOURCE, refresh_budget_registry
from apps.parsers.models import ParserLoadLog, ParserSourceArtifact, ParserStagedRecord
from apps.parsers.registry_snapshots import SnapshotValidationError, publish_snapshot
from apps.parsers.services import ParserLoadLogService
from apps.parsers.sme_support import (
SME_SOURCE,
discover_sme_archive,
import_sme_support,
iter_sme_measures,
)
from apps.parsers.tasks_registry_snapshots import _run_snapshot
from django.test import TestCase, override_settings
from django.utils import timezone
from lxml import etree
from organizations.models import Organization, OrganizationSourceRecord
from organizations.serializers import (
OrganizationSourceRecordListSerializer,
OrganizationSourceRecordSerializer,
)
from organizations.source_record_export import _source_group_queryset
from rest_framework.test import APIClient
from tests.apps.user.factories import UserFactory
def _xml(
*,
numbers=("old-support", "new-support"),
inn="1234567890",
documents=1,
units=("1", "3"),
):
root = etree.Element(
"Файл",
ИдФайл="fixture",
ВерсФорм="4.04",
ТипИнф="РЕЕСТРМСП-ПП",
КолДок=str(documents),
)
doc = etree.SubElement(
root, "Документ", ИдДок="recipient-doc", ДатаСост="15.08.2026"
)
etree.SubElement(
doc, "СвЮЛ", НаимОрг="АО Проверка", ИННЮЛ=inn, ОГРН="1027700132195"
)
for number in numbers:
measure = etree.SubElement(
doc,
"СвПредПод",
НомерПод=number,
ДатаСвед="04.10.2024",
ДатаОбнов="01.07.2026",
ВидПП="1",
НаимОрг="Поставщик",
ИННЮЛ="9876543210",
КатСуб="1",
СрокПод="01.01.2027",
ДатаПрин="01.10.2024",
ИнфНаруш="1",
)
etree.SubElement(measure, "ФормПод", КодФорм="0001", НаимФорм="Форма")
etree.SubElement(measure, "ВидПод", КодВид="0002", НаимВид="Вид")
for unit in units:
etree.SubElement(measure, "РазмПод", ЕдПод=unit, РазмПод="2.50")
etree.SubElement(measure, "Нарушения", ВидНаруш="1", ДатаНаруш="01.02.2025")
etree.SubElement(measure, "РегДок", НомерРД="A", НаимРД="Первый документ")
etree.SubElement(measure, "РегДок", НомерРД="B", НаимРД="Второй документ")
return etree.tostring(root, encoding="UTF-8", xml_declaration=True)
def _archive(xml: bytes) -> BytesIO:
handle = BytesIO()
with zipfile.ZipFile(handle, "w", compression=zipfile.ZIP_DEFLATED) as archive:
archive.writestr("support.xml", xml)
handle.seek(0)
return handle
def _budget_record(identifier="1", *, inn="1234567890", status="2"):
return {
"id": identifier,
"info": {
"inn": inn,
"ogrn": "1027700132195",
"fullName": "АО Проверка",
"okpoCode": "12345678",
"statusCode": status,
"dateUpdate": "2024-10-01T15:00:00",
"code": "own-code",
"recordNum": "00000001",
},
"activities": [{"name": "activity"}],
"unknown_new_block": [{"raw": "preserved"}],
}
class RegistrySnapshotTest(TestCase):
def setUp(self):
self.storage = TemporaryDirectory()
self.addCleanup(self.storage.cleanup)
self.settings = override_settings(MEDIA_ROOT=self.storage.name)
self.settings.enable()
self.addCleanup(self.settings.disable)
self.organization = Organization.objects.create(
name="АО Проверка",
inn="1234567890",
ogrn="1027700132195",
okpo="12345678",
directory_imported_at=timezone.now(),
opk_registry_membership=True,
)
def import_sme(self, xml=None, batch=1):
return import_sme_support(
handle=_archive(xml or _xml()),
original_name="snapshot.zip",
snapshot_date="2026-08-15",
load_batch=batch,
)
def test_sme_multiple_measures_nested_data_and_stable_ids(self):
artifact, result = self.import_sme()
self.assertEqual((result.parsed, result.published), (2, 2))
records = list(OrganizationSourceRecord.objects.order_by("external_id"))
self.assertEqual(records[0].record_date, "2024-10-01")
self.assertEqual(str(records[0].amount), "2.50")
self.assertEqual(len(records[0].payload["support_sizes"]), 2)
self.assertEqual(len(records[0].payload["regulatory_documents"]), 2)
self.assertEqual(records[0].payload["source_updated_date"], "2026-07-01")
self.assertEqual(len(records[0].payload["violations"]), 1)
before = {record.external_id: record.uid for record in records}
self.import_sme(batch=2)
self.assertEqual(
dict(OrganizationSourceRecord.objects.values_list("external_id", "uid")),
before,
)
self.assertEqual(Organization.objects.count(), 1)
self.assertEqual(artifact.status, ParserSourceArtifact.Status.PUBLISHED)
def test_sme_non_rub_amount_is_null_and_foreign_subject_not_created(self):
self.import_sme(_xml(units=("3", "4")))
self.assertIsNone(OrganizationSourceRecord.objects.first().amount)
_, result = self.import_sme(_xml(inn="5678901234"), batch=2)
self.assertEqual(result.published, 0)
self.assertEqual(result.skipped, 2)
self.assertEqual(Organization.objects.count(), 1)
self.assertEqual(OrganizationSourceRecord.objects.count(), 0)
def test_sme_invalid_truncated_or_conflicting_snapshot_keeps_old_records(self):
self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
for xml in (_xml(numbers=("same", "same")), _xml()[:-15]):
with self.assertRaises((SnapshotValidationError, etree.XMLSyntaxError)):
self.import_sme(xml, batch=2)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)),
before,
)
def test_official_document_count_mismatch_is_retained_as_diagnostic(self):
root = etree.fromstring(_xml()) # noqa: S320 - locally generated fixture
other = etree.fromstring(_xml(numbers=("third-measure",))) # noqa: S320
root.append(other.find("Документ"))
artifact, result = self.import_sme(etree.tostring(root, encoding="UTF-8"))
self.assertEqual(result.published, 3)
self.assertEqual(artifact.metadata["documents_count"], 2)
self.assertEqual(artifact.metadata["declared_documents_count"], 1)
self.assertEqual(artifact.metadata["document_count_mismatch_files"], 1)
def test_sme_publish_failure_rolls_back_updates_and_deletes(self):
self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
with self.assertRaises(RuntimeError), patch(
"apps.parsers.registry_snapshots.OrganizationSourceIngestionService.save_records",
side_effect=RuntimeError("stop"),
):
self.import_sme(_xml(numbers=("replacement",)), batch=2)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), before
)
def test_repeated_publication_keeps_previous_records_and_counts(self):
artifact, original = self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
self.assertEqual(publish_snapshot(artifact), original)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), before
)
artifact.refresh_from_db()
self.assertEqual(artifact.published_count, 2)
def test_cache_failure_rolls_back_publication(self):
self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
with self.assertRaises(RuntimeError), patch(
"apps.parsers.registry_snapshots.invalidate_source_data_cache",
side_effect=RuntimeError("cache"),
):
self.import_sme(_xml(numbers=("replacement",)), batch=2)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), before
)
def test_cache_is_invalidated_again_only_after_commit(self):
with patch(
"apps.parsers.registry_snapshots.invalidate_source_data_cache"
) as invalidate, self.captureOnCommitCallbacks(execute=True) as callbacks:
artifact, _ = self.import_sme()
self.assertEqual(invalidate.call_count, 1)
artifact.refresh_from_db()
self.assertEqual(artifact.status, ParserSourceArtifact.Status.PUBLISHED)
self.assertEqual(len(callbacks), 1)
self.assertEqual(invalidate.call_count, 2)
def test_after_commit_cache_failure_does_not_reject_published_data(self):
with patch(
"apps.parsers.registry_snapshots.invalidate_source_data_cache",
side_effect=[None, RuntimeError("cache")],
), self.assertLogs(
"apps.parsers.registry_snapshots", level="ERROR"
), self.captureOnCommitCallbacks(execute=True):
artifact, result = self.import_sme()
artifact.refresh_from_db()
self.assertEqual(artifact.status, ParserSourceArtifact.Status.PUBLISHED)
self.assertEqual(OrganizationSourceRecord.objects.count(), result.published)
def test_job_finalization_failure_rolls_back_publication(self):
self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
task = SimpleNamespace(
name="task", request=SimpleNamespace(id="finalize-failure")
)
def refresh(**kwargs):
return import_sme_support(
handle=_archive(_xml(numbers=("replacement",))),
original_name="next.zip",
snapshot_date="2026-08-15",
**kwargs,
)
with self.assertRaises(RuntimeError), patch.object(
BackgroundJob, "complete", side_effect=RuntimeError("finalization")
):
_run_snapshot(
task, source=SME_SOURCE, refresh=refresh, requested_by_id=None
)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), before
)
self.assertEqual(
BackgroundJob.objects.get(task_id="finalize-failure").status,
JobStatus.FAILURE,
)
self.assertEqual(
ParserLoadLog.objects.get().status, ParserLoadLog.Status.FAILED
)
def test_revoked_job_cannot_replace_previous_published_snapshot(self):
self.import_sme()
before = set(OrganizationSourceRecord.objects.values_list("uid", flat=True))
task = SimpleNamespace(
name="task", request=SimpleNamespace(id="revoked-import")
)
def refresh(**kwargs):
BackgroundJob.objects.filter(task_id="revoked-import").update(
status=JobStatus.REVOKED
)
return import_sme_support(
handle=_archive(_xml(numbers=("replacement",))),
original_name="next.zip",
snapshot_date="2026-08-15",
**kwargs,
)
with self.assertRaisesMessage(
SnapshotValidationError, "snapshot_job_no_longer_active"
):
_run_snapshot(
task, source=SME_SOURCE, refresh=refresh, requested_by_id=None
)
self.assertEqual(
set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), before
)
self.assertEqual(
BackgroundJob.objects.get(task_id="revoked-import").status,
JobStatus.REVOKED,
)
self.assertEqual(
ParserLoadLog.objects.get().status, ParserLoadLog.Status.FAILED
)
def test_card_list_dashboard_and_export_count_all_published_organizations(self):
Organization.objects.create(
name="Справочная организация",
inn="5678901234",
ogrn="1027700132195",
okpo="87654321",
directory_imported_at=timezone.now(),
opk_registry_membership=False,
)
root = etree.fromstring(_xml()) # noqa: S320 - locally generated fixture
other = etree.fromstring(_xml(numbers=("non-opk",), inn="5678901234")) # noqa: S320
root.append(other.find("Документ"))
root.set("КолДок", "2")
_, result = self.import_sme(etree.tostring(root, encoding="UTF-8"))
self.assertEqual(result.published, 3)
client = APIClient()
client.force_authenticate(UserFactory.create_user())
card = client.get("/api/v1/sources/sme-support-recipients-registry/")
records = client.get(
"/api/v2/organization-source-records/", {"source": SME_SOURCE}
)
dashboard = client.get("/api/v1/parsers/dashboard/")
self.assertEqual(
(card.status_code, records.status_code, dashboard.status_code),
(200, 200, 200),
)
self.assertEqual(card.data["data"]["records_count"], 3)
self.assertEqual(records.data["meta"]["pagination"]["total_count"], 3)
self.assertEqual(dashboard.data["data"]["source_counts"][SME_SOURCE], 3)
self.assertEqual(
_source_group_queryset("government_support", export_year=2026).count(), 3
)
def test_list_omits_nested_detail_but_detail_and_export_keep_all_years(self):
self.import_sme()
record = OrganizationSourceRecord.objects.first()
self.assertNotIn(
"violations", OrganizationSourceRecordListSerializer(record).data["payload"]
)
self.assertEqual(
len(
OrganizationSourceRecordSerializer(record).data["payload"]["violations"]
),
1,
)
self.assertEqual(
_source_group_queryset("government_support", export_year=2026).count(), 2
)
def test_budget_scans_foreign_without_detail_or_raw_retention(self):
own, foreign = _budget_record(), _budget_record("2", inn="5678901234")
calls = []
def page(_session, params):
calls.append(params)
if params.get("filterid"):
self.assertEqual(params["filterid"], "1")
return {"data": [own]}
return {
"data": [own, foreign],
"recordCount": 2,
"version": "10",
"pageNum": 1,
}
with patch("apps.parsers.budget_registry.budget_page", side_effect=page):
artifact, result = refresh_budget_registry(load_batch=1, session=object())
self.assertEqual((result.parsed, result.published, result.skipped), (2, 1, 1))
self.assertEqual(calls[0]["blocks"], "info")
self.assertEqual(len(calls), 2)
self.assertEqual(
ParserStagedRecord.objects.filter(artifact=artifact).count(), 1
)
with artifact.file.open("rb") as handle:
raw = handle.read().decode()
self.assertNotIn("5678901234", raw)
record = OrganizationSourceRecord.objects.get(source=BUDGET_SOURCE)
self.assertEqual(record.status, "inactive")
self.assertEqual(
record.payload["upstream"]["unknown_new_block"], [{"raw": "preserved"}]
)
self.assertEqual(
_source_group_queryset("budget_process_registry", export_year=2026).count(),
1,
)
self.assertEqual(Organization.objects.count(), 1)
def test_budget_incomplete_second_page_preserves_published_snapshot(self):
own = _budget_record()
first = {"data": [own], "recordCount": 1, "version": "10", "pageNum": 1}
with patch(
"apps.parsers.budget_registry.budget_page",
side_effect=[first, {"data": [own]}],
):
refresh_budget_registry(load_batch=1, session=object())
before = OrganizationSourceRecord.objects.get().uid
first["recordCount"] = 2
with patch(
"apps.parsers.budget_registry.budget_page",
side_effect=[
first,
{"data": [own]},
{"data": [], "recordCount": 2, "version": "10", "pageNum": 2},
],
), self.assertRaises(SnapshotValidationError):
refresh_budget_registry(load_batch=2, session=object())
self.assertEqual(OrganizationSourceRecord.objects.get().uid, before)
def test_budget_unknown_status_stays_unknown(self):
own = _budget_record(status="1")
with patch(
"apps.parsers.budget_registry.budget_page",
side_effect=[
{"data": [own], "recordCount": 1, "version": "10"},
{"data": [own]},
],
):
refresh_budget_registry(load_batch=1, session=object())
self.assertEqual(OrganizationSourceRecord.objects.get().status, "unknown")
class RegistryTaskLifecycleTest(TestCase):
def test_revoked_after_claim_does_not_attach_late_batch_metadata(self):
task = SimpleNamespace(name="task", request=SimpleNamespace(id="revoked-claim"))
original = ParserLoadLogService.create_load_log_with_next_batch_id
def create_log(**kwargs):
value = original(**kwargs)
BackgroundJob.objects.filter(task_id="revoked-claim").update(
status=JobStatus.REVOKED
)
return value
with self.assertRaisesMessage(
SnapshotValidationError, "snapshot_job_no_longer_active"
), patch.object(
ParserLoadLogService,
"create_load_log_with_next_batch_id",
side_effect=create_log,
), patch(
"apps.parsers.tasks_registry_snapshots.refresh_budget_registry"
) as refresh:
_run_snapshot(
task, source=BUDGET_SOURCE, refresh=refresh, requested_by_id=None
)
refresh.assert_not_called()
job = BackgroundJob.objects.get(task_id="revoked-claim")
self.assertEqual(job.status, JobStatus.REVOKED)
self.assertNotIn("batch_id", job.meta)
def test_validation_failure_reports_bounded_reason(self):
task = SimpleNamespace(name="task", request=SimpleNamespace(id="invalid-page"))
with self.assertRaises(SnapshotValidationError), patch(
"apps.parsers.tasks_registry_snapshots.refresh_budget_registry",
side_effect=SnapshotValidationError("incomplete_budget_snapshot"),
) as refresh:
_run_snapshot(
task, source=BUDGET_SOURCE, refresh=refresh, requested_by_id=None
)
job = BackgroundJob.objects.get(task_id="invalid-page")
self.assertIn("incomplete_budget_snapshot", job.error)
def test_duplicate_completed_delivery_returns_result_without_new_batch(self):
job = BackgroundJob.objects.create(
task_id="delivery",
task_name="task",
status=JobStatus.SUCCESS,
result={"published": 5},
)
task = SimpleNamespace(name="task", request=SimpleNamespace(id=job.task_id))
with patch(
"apps.parsers.tasks_registry_snapshots.refresh_sme_support"
) as refresh:
self.assertEqual(
_run_snapshot(
task, source=SME_SOURCE, refresh=refresh, requested_by_id=None
),
{"published": 5},
)
refresh.assert_not_called()
self.assertFalse(ParserLoadLog.objects.exists())
def test_duplicate_started_delivery_does_not_repeat_import(self):
BackgroundJob.objects.create(
task_id="delivery", task_name="task", status=JobStatus.STARTED
)
task = SimpleNamespace(name="task", request=SimpleNamespace(id="delivery"))
with patch(
"apps.parsers.tasks_registry_snapshots.refresh_sme_support"
) as refresh:
result = _run_snapshot(
task, source=SME_SOURCE, refresh=refresh, requested_by_id=None
)
self.assertTrue(result["duplicate_delivery"])
refresh.assert_not_called()
self.assertFalse(ParserLoadLog.objects.exists())
def test_sequence_failure_marks_precreated_job_failed(self):
job = BackgroundJob.objects.create(task_id="delivery", task_name="task")
task = SimpleNamespace(name="task", request=SimpleNamespace(id=job.task_id))
with self.assertRaises(RuntimeError), patch(
"apps.parsers.tasks_registry_snapshots.ParserLoadLogService.create_load_log_with_next_batch_id",
side_effect=RuntimeError("sequence"),
):
_run_snapshot(
task,
source=SME_SOURCE,
refresh=lambda **_: None,
requested_by_id=None,
)
job.refresh_from_db()
self.assertEqual(job.status, JobStatus.FAILURE)
def test_sme_discovery_selects_latest_official_snapshot():
url, snapshot = discover_sme_archive(
b'oldnew'
)
assert snapshot == "2026-08-15"
assert "data-20260815" in url
def test_sme_stream_releases_each_measure_before_reading_next_recipient():
class BoundedReader(BytesIO):
def read(self, size=-1):
assert 0 < size <= 65536
return super().read(size)
counts = Counter()
stream = iter_sme_measures(
BoundedReader(_xml(numbers=tuple(f"measure-{i}" for i in range(1200)))), counts
)
_, recipient, previous = next(stream)
assert recipient["ИННЮЛ"] == "1234567890"
_, recipient, current = next(stream)
assert len(previous) == 0
assert not previous.attrib
assert current.attrib["НомерПод"] == "measure-1"
remaining = 0
for _, recipient, _ in stream:
assert recipient["ИННЮЛ"] == "1234567890"
remaining += 1
assert remaining == 1198
assert counts == {"documents": 1, "measures": 1200, "declared_documents": 1}