"""Synthetic source-first HTML fixtures; no live query/search traffic.""" from __future__ import annotations import importlib import json import uuid import zipfile from datetime import timedelta from tempfile import TemporaryDirectory from types import SimpleNamespace from unittest.mock import Mock, patch from urllib.parse import parse_qs, urlsplit import requests from apps.core.exceptions import ConflictError from apps.core.models import BackgroundJob, JobStatus from apps.parsers.models import ( ParserSourceArtifact, SroOrganizationLookup, ) from apps.parsers.registry_snapshots import SnapshotValidationError from apps.parsers.source_artifacts import cleanup_parser_source_artifacts from apps.parsers.sro_http import SroHttpClient, SroPage, safe_sro_url from apps.parsers.sro_membership import ( SRO_LOOKUP_URL, SRO_SOURCE, SroResolver, admission_date, parse_sro_lookup, refresh_sro_membership, ) from apps.parsers.tasks_registry_snapshots import _run_snapshot from django.test import SimpleTestCase, TestCase, override_settings from django.utils import timezone from organizations.models import Organization, OrganizationSourceRecord from organizations.sro_payload_serializers import SroMembershipDetailPayloadSerializer APPROVED = { "SRO_UPSTREAM_ACCESS_APPROVED": True, "SRO_UPSTREAM_APPROVAL_REFERENCE": "synthetic-fixture-approval", } def lookup_html( organization, memberships=("001",), *, status="Является членом", href=None, inn=None ): rows = [] for identifier in memberships: url = ( href if href is not None else f"https://moskva.reestr-sro.ru/sro-v-proektirovanii/sro-id-{identifier}/" ) rows.append( f'Краткое: {organization.name}
Полное: Полное {organization.name}{inn or organization.inn}
{organization.ogrn}{status}МоскваСРО {identifier}' ) return ( '

Реестр обновлен 15 августа 2026

' + "".join(rows) + "
НаименованиеИНН, ОГРНСтатусРегион регистрацииСРО
" ).encode() def members_html(inn, date="01.02.2020"): return f'
ИННДата получения допуска
{inn}{date}
'.encode() class FixtureSite: def __init__(self, organizations, overrides=None, admitted="01.02.2020"): self.organizations = {org.ogrn: org for org in organizations} self.overrides = overrides or {} self.admitted = admitted self.urls = [] def get(self, url): self.urls.append(url) parts = urlsplit(url) query = parse_qs(parts.query)["q"][0] if "/members/" in parts.path: return SroPage(url, members_html(query, self.admitted)) value = self.overrides.get(query) if isinstance(value, Exception): raise value return SroPage( url, value if value is not None else lookup_html(self.organizations[query]) ) class SroHttpTest(SimpleTestCase): def test_gate_off_or_missing_reference_never_opens_session(self): for approved, reference in ((False, ""), (True, ""), (False, "written")): with self.subTest( approved=approved, reference=reference ), override_settings( SRO_UPSTREAM_ACCESS_APPROVED=approved, SRO_UPSTREAM_APPROVAL_REFERENCE=reference, ), patch("apps.parsers.sro_http.requests.Session") as session: with self.assertRaises(ConflictError) as caught: SroHttpClient() self.assertEqual( (caught.exception.code, caught.exception.status_code), ("upstream_access_not_approved", 409), ) session.assert_not_called() def test_url_allowlist_rejects_lookalikes_credentials_and_ports(self): for url in ( "http://reestr-sro.ru/", "https://reestr-sro.ru.evil.test/", "https://evilreestr-sro.ru/", "https://user@reestr-sro.ru/", "https://reestr-sro.ru:8443/", "https://127.0.0.1/", ): with self.subTest(url=url), self.assertRaises(SnapshotValidationError): safe_sro_url(url) self.assertEqual( safe_sro_url("/members/", base="https://moskva.reestr-sro.ru/"), "https://moskva.reestr-sro.ru/members/", ) @override_settings(**APPROVED) def test_unsafe_redirect_never_followed(self): response = Mock( status_code=302, is_redirect=True, headers={"Location": "https://example.invalid/"}, ) session = Mock() session.get.return_value = response client = SroHttpClient(session=session) with self.assertRaisesMessage(SnapshotValidationError, "unsafe_sro_url"): client.get(SRO_LOOKUP_URL) self.assertEqual(session.get.call_count, 1) self.assertFalse(session.get.call_args.kwargs["allow_redirects"]) response.close.assert_called_once() @override_settings(**APPROVED) def test_retries_and_redirects_are_paced_and_bounded(self): now = [0.0] pauses = [] def sleep(value): pauses.append(value) now[0] += value response = Mock(status_code=200, is_redirect=False) response.iter_content.return_value = [b""] session = Mock() session.get.side_effect = [requests.Timeout(), response] client = SroHttpClient(session=session, clock=lambda: now[0], sleep=sleep) result = client.get(SRO_LOOKUP_URL) self.assertEqual(result.body, b"") self.assertEqual(pauses, [3.0]) self.assertEqual(client.http_errors_count, 1) @override_settings(**APPROVED, SRO_HTTP_MAX_RESPONSE_BYTES=4) def test_response_size_limit(self): response = Mock(status_code=200, is_redirect=False) response.iter_content.return_value = [b"12345"] session = Mock() session.get.return_value = response with self.assertRaisesMessage( SnapshotValidationError, "source_response_too_large" ): SroHttpClient(session=session).get(SRO_LOOKUP_URL) response.close.assert_called_once() def test_parser_does_not_treat_unknown_html_as_empty(self): with self.assertRaisesMessage( SnapshotValidationError, "sro_members_table_missing" ): parse_sro_lookup(SroPage(SRO_LOOKUP_URL, b"captcha")) self.assertEqual( parse_sro_lookup(SroPage(SRO_LOOKUP_URL, "Ничего не найдено".encode()))[0], [], ) def test_nullable_admission_reason(self): self.assertEqual( admission_date( SroPage(SRO_LOOKUP_URL, members_html("1234567890", "—")), "1234567890" ), (None, "not_found"), ) self.assertEqual( admission_date( SroPage(SRO_LOOKUP_URL, members_html("1234567890", "31.02.2020")), "1234567890", ), (None, "parse_error"), ) def test_sitemap_id_resolution_without_name_guessing(self): card = "https://moskva.reestr-sro.ru/sro-v-proektirovanii/sro-id-001/" pages = { "https://www.reestr-sro.ru/robots.txt": b"Sitemap: https://www.reestr-sro.ru/sitemap.xml", "https://www.reestr-sro.ru/sitemap.xml": f'{card}'.encode(), } resolver = SroResolver(lambda url: SroPage(url, pages[url])) row = {"sro_href": "", "sro_id_hint": "001"} self.assertEqual( resolver.resolve(row, SRO_LOOKUP_URL), ("001", card, "same_site_sitemap") ) self.assertEqual(resolver.resolve(row, SRO_LOOKUP_URL)[2], "same_site_id") @override_settings(**APPROVED) def test_permission_revocation_before_following_redirect_stops_io(self): response = Mock( status_code=302, is_redirect=True, headers={"Location": "/next/"} ) session = Mock() def revoke(*args, **kwargs): from django.conf import settings settings.SRO_UPSTREAM_ACCESS_APPROVED = False return response session.get.side_effect = revoke with self.assertRaises(ConflictError): SroHttpClient(session=session).get(SRO_LOOKUP_URL) session.get.assert_called_once() response.close.assert_called_once() @override_settings(**APPROVED) class SroSnapshotTest(TestCase): def setUp(self): self.storage = TemporaryDirectory() self.addCleanup(self.storage.cleanup) self.settings = override_settings( MEDIA_ROOT=f"{self.storage.name}/media", PARSER_PRIVATE_ARTIFACT_ROOT=f"{self.storage.name}/private", ) self.settings.enable() self.addCleanup(self.settings.disable) self.org = Organization.objects.create( name="АО Фикстура", inn="1234567890", ogrn="1027700132195", okpo="00123456", directory_imported_at=timezone.now(), ) def load(self, site=None, *, mode="full", **kwargs): return refresh_sro_membership( load_batch=ParserSourceArtifact.objects.count() + 1, mode=mode, client=site or FixtureSite([self.org]), **kwargs, ) def test_multiple_memberships_stable_uid_and_status_transition(self): site = FixtureSite( [self.org], {self.org.ogrn: lookup_html(self.org, ("001", "002"))} ) artifact, result = self.load(site) self.assertEqual((result.parsed, result.published), (2, 2)) ids = set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) self.load( FixtureSite( [self.org], { self.org.ogrn: lookup_html( self.org, ("001", "002"), status="Исключен" ) }, ) ) self.assertEqual( ids, set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) ) self.assertEqual( set(OrganizationSourceRecord.objects.values_list("status", flat=True)), {"inactive"}, ) self.assertEqual(Organization.objects.count(), 1) with zipfile.ZipFile(artifact.file.path) as raw: self.assertEqual(raw.testzip(), None) self.assertEqual(len(json.loads(raw.read("manifest.json"))), 3) payload = OrganizationSourceRecord.objects.first().payload schema = SroMembershipDetailPayloadSerializer(data=payload) self.assertTrue(schema.is_valid(), schema.errors) self.assertNotIn("source_organization", schema.data) self.assertEqual( payload["source_organization"]["full_name"], "Полное АО Фикстура" ) def test_incremental_does_not_delete_unqueried_org_and_full_still_replaces_all( self, ): other = Organization.objects.create( name="АО Вторая", inn="2234567890", ogrn="2027700132195", okpo="00123457", directory_imported_at=timezone.now(), ) self.load(FixtureSite([self.org, other])) other_id = OrganizationSourceRecord.objects.get( extension__organization=other ).uid Organization.objects.filter(pk=self.org.pk).update(name="АО Изменённая") site = FixtureSite([self.org], {self.org.ogrn: lookup_html(self.org, ())}) artifact, result = self.load(site, mode="incremental") self.assertEqual(result.published, 1) self.assertEqual(artifact.metadata["updated_records_count"], 0) self.assertEqual( list(OrganizationSourceRecord.objects.values_list("uid", flat=True)), [other_id], ) self.assertEqual(len(site.urls), 1) self.load( FixtureSite( [self.org, other], { self.org.ogrn: lookup_html(self.org, ()), other.ogrn: lookup_html(other, ()), }, ) ) self.assertFalse(OrganizationSourceRecord.objects.exists()) def test_incremental_failure_preserves_records_and_checkpoints(self): self.load() before = list(OrganizationSourceRecord.objects.values_list("uid", "payload")) checkpoint = SroOrganizationLookup.objects.get(organization=self.org).checked_at Organization.objects.filter(pk=self.org.pk).update(name="АО Изменённая") with self.assertRaisesMessage( SnapshotValidationError, "sro_members_table_missing" ): self.load( FixtureSite([self.org], {self.org.ogrn: b"blocked"}), mode="incremental", ) self.assertEqual( before, list(OrganizationSourceRecord.objects.values_list("uid", "payload")) ) self.assertEqual( checkpoint, SroOrganizationLookup.objects.get(organization=self.org).checked_at, ) self.assertEqual(ParserSourceArtifact.objects.first().status, "rejected") rejected = ParserSourceArtifact.objects.first() self.assertEqual(rejected.metadata["parse_errors_count"], 1) with zipfile.ZipFile(rejected.file.path) as raw: self.assertIsNone(raw.testzip()) self.assertEqual(len(json.loads(raw.read("manifest.json"))), 1) def test_failed_finalization_rolls_back_publication_and_checkpoint(self): self.load() before = list(OrganizationSourceRecord.objects.values_list("uid", "payload")) checkpoint = SroOrganizationLookup.objects.get(organization=self.org).checked_at def fail(*args): raise SnapshotValidationError("snapshot_job_no_longer_active") with self.assertRaises(SnapshotValidationError): self.load( FixtureSite([self.org], {self.org.ogrn: lookup_html(self.org, ())}), on_publish=fail, ) self.assertEqual( before, list(OrganizationSourceRecord.objects.values_list("uid", "payload")) ) self.assertEqual( checkpoint, SroOrganizationLookup.objects.get(organization=self.org).checked_at, ) def test_incremental_unchanged_directory_does_not_query_or_drop_memberships(self): previous, _ = self.load() previous.refresh_from_db() original = OrganizationSourceRecord.objects.get().uid site = FixtureSite([self.org]) artifact, result = self.load(site, mode="incremental") self.assertEqual(site.urls, []) self.assertEqual(result.published, 1) self.assertEqual(artifact.metadata["batch_published_records_count"], 0) self.assertEqual(artifact.metadata["snapshot_organizations_count"], 1) self.assertEqual(OrganizationSourceRecord.objects.get().uid, original) self.assertEqual(artifact.source_published_at, previous.source_published_at) self.assertEqual(artifact.version, previous.version) self.assertEqual( artifact.metadata["base_snapshot_artifact_id"], str(previous.uid) ) def test_semantic_errors_quarantined_without_creating_organizations(self): self.load() old_ids = set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) cases = ( ({"status": "Новый статус"}, "unknown_membership_status"), ({"href": "https://other.invalid/card"}, "unsafe_sro_url"), ({"inn": "9999999999"}, "organization_identifier_conflict"), ) for params, reason in cases: with self.subTest(reason=reason): with self.assertRaisesMessage( SnapshotValidationError, "sro_incomplete_membership_scan" ): self.load( FixtureSite( [self.org], {self.org.ogrn: lookup_html(self.org, **params)} ) ) artifact = ParserSourceArtifact.objects.first() self.assertEqual( (artifact.status, artifact.quarantined_count), ("rejected", 1) ) self.assertEqual(artifact.staged_records.get().reason, reason) self.assertEqual( old_ids, set(OrganizationSourceRecord.objects.values_list("uid", flat=True)), ) self.org.okpo = "" self.org.save() with self.assertRaisesMessage( SnapshotValidationError, "sro_incomplete_membership_scan" ): self.load() artifact = ParserSourceArtifact.objects.first() self.assertEqual(artifact.staged_records.get().reason, "required_okpo_missing") self.assertEqual(Organization.objects.count(), 1) def test_incremental_quarantine_preserves_memberships_and_does_not_checkpoint(self): self.load() old_ids = set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) checkpoint = SroOrganizationLookup.objects.get(organization=self.org).checked_at self.org.name = "АО Измененная" self.org.save(update_fields=["name"]) site = FixtureSite( [self.org], {self.org.ogrn: lookup_html(self.org, status="Новый статус")} ) with self.assertRaisesMessage( SnapshotValidationError, "sro_incomplete_membership_scan" ): self.load(site, mode="incremental") self.assertEqual( old_ids, set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) ) self.assertEqual( SroOrganizationLookup.objects.get(organization=self.org).checked_at, checkpoint, ) from apps.parsers.sro_membership import _sro_candidates self.assertEqual( [item.uid for item in _sro_candidates("incremental")], [self.org.uid] ) rejected = ParserSourceArtifact.objects.first() self.assertEqual(rejected.status, "rejected") self.assertEqual( rejected.staged_records.get().reason, "unknown_membership_status" ) def test_duplicate_membership_rejects_whole_snapshot(self): self.load() with self.assertRaisesMessage( SnapshotValidationError, "duplicate_registry_number" ): self.load( FixtureSite( [self.org], {self.org.ogrn: lookup_html(self.org, ("001", "001"))} ) ) self.assertEqual(OrganizationSourceRecord.objects.count(), 1) @override_settings(SRO_UPSTREAM_ACCESS_APPROVED=False) def test_task_gate_closes_precreated_job_without_external_io(self): task_id = str(uuid.uuid4()) job = BackgroundJob.objects.create( task_id=task_id, task_name="parsers.sro_membership_check.refresh" ) task = SimpleNamespace(request=SimpleNamespace(id=task_id), name=job.task_name) with patch( "apps.parsers.sro_membership.SroHttpClient" ) as http, self.assertRaises(ConflictError): _run_snapshot( task, source=SRO_SOURCE, refresh=refresh_sro_membership, requested_by_id=None, ) job.refresh_from_db() self.assertEqual(job.status, JobStatus.FAILURE) self.assertIn("upstream_access_not_approved", job.error) http.assert_not_called() self.assertFalse(ParserSourceArtifact.objects.exists()) def test_sro_retention_protects_ten_successful_versions(self): artifacts = [ ParserSourceArtifact.objects.create(source=SRO_SOURCE, status="published") for _ in range(11) ] for _ in range(12): ParserSourceArtifact.objects.create(source=SRO_SOURCE, status="rejected") ParserSourceArtifact.objects.update( created_at=timezone.now() - timedelta(days=100) ) # Ensure deterministic source ordering independent of UUID/identical timestamps. for index, artifact in enumerate(artifacts): ParserSourceArtifact.objects.filter(pk=artifact.pk).update( created_at=timezone.now() - timedelta(days=120 - index) ) cleanup_parser_source_artifacts() self.assertEqual( ParserSourceArtifact.objects.filter( source=SRO_SOURCE, status="published" ).count(), 10, ) def test_schedule_seed_is_disabled_idempotent_and_reversible(self): from django.apps import apps from django_celery_beat.models import PeriodicTask migration = importlib.import_module( "apps.parsers.migrations.0036_sro_disabled_schedules" ) migration.seed(apps, None) migration.seed(apps, None) scheduled = PeriodicTask.objects.filter( task="parsers.sro_membership_check.refresh" ) self.assertEqual(scheduled.count(), 2) self.assertFalse(scheduled.filter(enabled=True).exists()) self.assertEqual( { json.loads(value)["mode"] for value in scheduled.values_list("kwargs", flat=True) }, {"incremental", "full"}, ) migration.unseed(apps, None) self.assertFalse(scheduled.exists()) def test_scheduled_task_cannot_compete_with_started_source(self): active = BackgroundJob.objects.create( task_id=str(uuid.uuid4()), task_name="parsers.sro_membership_check.refresh", status=JobStatus.STARTED, ) task = SimpleNamespace( request=SimpleNamespace(id=str(uuid.uuid4())), name=active.task_name ) refresh = Mock() with self.assertRaises(ConflictError) as caught: _run_snapshot( task, source=SRO_SOURCE, refresh=refresh, requested_by_id=None ) self.assertEqual(caught.exception.code, "refresh_already_running") refresh.assert_not_called() active.refresh_from_db() self.assertEqual(active.status, JobStatus.STARTED) def test_cancellation_before_publication_preserves_old_snapshot_and_checkpoint( self, ): self.load() old_ids = set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) checkpoint = SroOrganizationLookup.objects.get(organization=self.org).checked_at task = SimpleNamespace( request=SimpleNamespace(id=str(uuid.uuid4())), name="parsers.sro_membership_check.refresh", ) site = FixtureSite([self.org], {self.org.ogrn: lookup_html(self.org, ())}) def refresh(**kwargs): return refresh_sro_membership( client=site, mode="full", on_progress=lambda *_: BackgroundJob.objects.get( task_id=task.request.id ).revoke(), **kwargs, ) with self.assertRaisesMessage( SnapshotValidationError, "snapshot_job_no_longer_active" ): _run_snapshot( task, source=SRO_SOURCE, refresh=refresh, requested_by_id=None ) self.assertEqual( BackgroundJob.objects.get(task_id=task.request.id).status, JobStatus.REVOKED ) self.assertEqual( old_ids, set(OrganizationSourceRecord.objects.values_list("uid", flat=True)) ) self.assertEqual( checkpoint, SroOrganizationLookup.objects.get(organization=self.org).checked_at, ) def test_export_enqueue_runs_only_after_commit_and_failure_keeps_success(self): task = SimpleNamespace( request=SimpleNamespace(id=str(uuid.uuid4())), name="parsers.sro_membership_check.refresh", ) def refresh(**kwargs): return refresh_sro_membership(client=FixtureSite([self.org]), **kwargs) with patch( "organizations.tasks.refresh_source_record_export_artifacts.delay", side_effect=RuntimeError("broker unavailable"), ) as enqueue: with self.captureOnCommitCallbacks(execute=False) as callbacks: result = _run_snapshot( task, source=SRO_SOURCE, refresh=refresh, requested_by_id=None ) enqueue.assert_not_called() for callback in callbacks: callback() enqueue.assert_called_once() self.assertEqual( BackgroundJob.objects.get(task_id=task.request.id).status, JobStatus.SUCCESS ) self.assertEqual(result["snapshot_records_count"], 1) self.assertEqual(result["snapshot_organizations_count"], 1) self.assertEqual(result["records_count"], 1) self.assertEqual(result["organizations_count"], 1) self.assertEqual(result["raw_memberships_count"], 1) self.assertEqual(result["http_errors_count"], 0) self.assertEqual(result["parse_errors_count"], 0) meta = BackgroundJob.objects.get(task_id=task.request.id).meta self.assertEqual(meta["published_records_count"], 1) self.assertEqual(meta["queried_organizations_count"], 1)