From 7c50cf82addbd7c45452e0a577f2184963a41e53 Mon Sep 17 00:00:00 2001 From: Aleksandr Meshchriakov Date: Thu, 20 Aug 2026 15:50:04 +0200 Subject: [PATCH] feat: add ROPK sanctions file source --- src/apps/parsers/api_result_urls.py | 19 +- .../migrations/0033_auto_20260820_1325.py | 33 ++ src/apps/parsers/models.py | 1 + src/apps/parsers/ropk_sanctions.py | 433 +++++++++++++++ src/apps/parsers/serializers.py | 1 + src/apps/parsers/source_cards.py | 26 +- src/apps/parsers/source_registry.py | 19 + src/apps/parsers/tasks.py | 85 +++ src/apps/parsers/views.py | 163 +++++- src/organizations/filters.py | 2 + .../migrations/0011_auto_20260820_1325.py | 31 ++ src/organizations/models.py | 12 + src/organizations/serializers.py | 8 + src/organizations/source_groups.py | 8 + src/organizations/source_ingestion.py | 22 + src/organizations/source_record_export.py | 64 ++- src/organizations/test_companies.py | 8 +- src/organizations/views.py | 19 + .../test_source_record_export.py | 97 +++- tests/apps/organizations/test_tasks.py | 2 +- .../test_test_companies_commands.py | 6 +- tests/apps/parsers/test_ropk_sanctions.py | 500 ++++++++++++++++++ .../apps/parsers/test_source_cards_service.py | 2 + 23 files changed, 1520 insertions(+), 41 deletions(-) create mode 100644 src/apps/parsers/migrations/0033_auto_20260820_1325.py create mode 100644 src/apps/parsers/ropk_sanctions.py create mode 100644 src/organizations/migrations/0011_auto_20260820_1325.py create mode 100644 tests/apps/parsers/test_ropk_sanctions.py diff --git a/src/apps/parsers/api_result_urls.py b/src/apps/parsers/api_result_urls.py index 0050455..4257860 100644 --- a/src/apps/parsers/api_result_urls.py +++ b/src/apps/parsers/api_result_urls.py @@ -9,6 +9,7 @@ from apps.parsers.views import ( MEDIA_NEWS_UPLOAD_FILE_PARAM, RESULT_DETAIL_PARAMS, RESULT_LIST_PARAMS, + ROPK_SANCTIONS_UPLOAD_FILE_PARAM, UPLOAD_FILE_PARAM, ParserUploadView, SourceResultDetailView, @@ -125,13 +126,23 @@ def _upload_view(descriptor: ParserSourceDescriptor): "Файл обрабатывается через Celery." ), manual_parameters=[ - MEDIA_NEWS_UPLOAD_FILE_PARAM - if descriptor.key == "media_news" - else UPLOAD_FILE_PARAM + ( + MEDIA_NEWS_UPLOAD_FILE_PARAM + if descriptor.key == "media_news" + else ROPK_SANCTIONS_UPLOAD_FILE_PARAM + if descriptor.key == "ropk_sanctions" + else UPLOAD_FILE_PARAM + ) ], consumes=["multipart/form-data"], tags=[tag], - responses={202: ParserRunResponseSerializer, 400: "Ошибка валидации"}, + responses={ + 202: ParserRunResponseSerializer, + 400: "Ошибка валидации", + 409: "Загрузка уже выполняется", + 413: "Файл превышает лимит", + 415: "Неподдерживаемый тип файла", + }, ) def post(self, request): return super().post(request, source_key=descriptor.key) diff --git a/src/apps/parsers/migrations/0033_auto_20260820_1325.py b/src/apps/parsers/migrations/0033_auto_20260820_1325.py new file mode 100644 index 0000000..6fadefa --- /dev/null +++ b/src/apps/parsers/migrations/0033_auto_20260820_1325.py @@ -0,0 +1,33 @@ +# Generated by Django 3.2.25 on 2026-08-20 13:25 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('parsers', '0032_seed_gosedo_and_artifact_schedules'), + ] + + operations = [ + migrations.AlterField( + model_name='genericparserrecord', + name='source', + field=models.CharField(choices=[('industrial', 'Сертификаты промышленного производства'), ('industrial_products', 'Реестр промышленной продукции'), ('manufactures', 'Реестр производителей'), ('inspections', 'Единый реестр проверок'), ('procurements', 'Единая информационная система закупок'), ('fns_reports', 'Бухгалтерская отчетность ФНС'), ('procurements_44fz', 'Закупки 44-ФЗ'), ('procurements_223fz', 'Закупки 223-ФЗ'), ('contracts', 'Контракты ЕИС'), ('unfair_suppliers', 'Недобросовестные поставщики'), ('fas_goz', 'Уклонение от ГОЗ'), ('arbitration', 'Арбитражные дела'), ('fedresurs_bankruptcy', 'Банкротства Федресурс'), ('fstec', 'Реестры ФСТЭК'), ('trudvsem', 'Вакансии Работа России'), ('gosedo_address_directory', 'Глобальный адресный справочник ГосЭДО'), ('media_news', 'Новости СМИ'), ('ropk_sanctions', 'РОПК — Санкции'), ('hh', 'Вакансии HeadHunter'), ('superjob', 'Вакансии SuperJob')], db_index=True, help_text='Источник данных', max_length=50, verbose_name='источник'), + ), + migrations.AlterField( + model_name='parserbatchsequence', + name='source', + field=models.CharField(choices=[('industrial', 'Сертификаты промышленного производства'), ('industrial_products', 'Реестр промышленной продукции'), ('manufactures', 'Реестр производителей'), ('inspections', 'Единый реестр проверок'), ('procurements', 'Единая информационная система закупок'), ('fns_reports', 'Бухгалтерская отчетность ФНС'), ('procurements_44fz', 'Закупки 44-ФЗ'), ('procurements_223fz', 'Закупки 223-ФЗ'), ('contracts', 'Контракты ЕИС'), ('unfair_suppliers', 'Недобросовестные поставщики'), ('fas_goz', 'Уклонение от ГОЗ'), ('arbitration', 'Арбитражные дела'), ('fedresurs_bankruptcy', 'Банкротства Федресурс'), ('fstec', 'Реестры ФСТЭК'), ('trudvsem', 'Вакансии Работа России'), ('gosedo_address_directory', 'Глобальный адресный справочник ГосЭДО'), ('media_news', 'Новости СМИ'), ('ropk_sanctions', 'РОПК — Санкции')], help_text='Источник данных', max_length=50, unique=True, verbose_name='источник'), + ), + migrations.AlterField( + model_name='parserloadlog', + name='source', + field=models.CharField(choices=[('industrial', 'Сертификаты промышленного производства'), ('industrial_products', 'Реестр промышленной продукции'), ('manufactures', 'Реестр производителей'), ('inspections', 'Единый реестр проверок'), ('procurements', 'Единая информационная система закупок'), ('fns_reports', 'Бухгалтерская отчетность ФНС'), ('procurements_44fz', 'Закупки 44-ФЗ'), ('procurements_223fz', 'Закупки 223-ФЗ'), ('contracts', 'Контракты ЕИС'), ('unfair_suppliers', 'Недобросовестные поставщики'), ('fas_goz', 'Уклонение от ГОЗ'), ('arbitration', 'Арбитражные дела'), ('fedresurs_bankruptcy', 'Банкротства Федресурс'), ('fstec', 'Реестры ФСТЭК'), ('trudvsem', 'Вакансии Работа России'), ('gosedo_address_directory', 'Глобальный адресный справочник ГосЭДО'), ('media_news', 'Новости СМИ'), ('ropk_sanctions', 'РОПК — Санкции')], db_index=True, help_text='Источник данных', max_length=50, verbose_name='источник'), + ), + migrations.AlterField( + model_name='parsersourceartifact', + name='source', + field=models.CharField(choices=[('industrial', 'Сертификаты промышленного производства'), ('industrial_products', 'Реестр промышленной продукции'), ('manufactures', 'Реестр производителей'), ('inspections', 'Единый реестр проверок'), ('procurements', 'Единая информационная система закупок'), ('fns_reports', 'Бухгалтерская отчетность ФНС'), ('procurements_44fz', 'Закупки 44-ФЗ'), ('procurements_223fz', 'Закупки 223-ФЗ'), ('contracts', 'Контракты ЕИС'), ('unfair_suppliers', 'Недобросовестные поставщики'), ('fas_goz', 'Уклонение от ГОЗ'), ('arbitration', 'Арбитражные дела'), ('fedresurs_bankruptcy', 'Банкротства Федресурс'), ('fstec', 'Реестры ФСТЭК'), ('trudvsem', 'Вакансии Работа России'), ('gosedo_address_directory', 'Глобальный адресный справочник ГосЭДО'), ('media_news', 'Новости СМИ'), ('ropk_sanctions', 'РОПК — Санкции')], db_index=True, max_length=50), + ), + ] diff --git a/src/apps/parsers/models.py b/src/apps/parsers/models.py index 1aa126b..0b3b44e 100644 --- a/src/apps/parsers/models.py +++ b/src/apps/parsers/models.py @@ -42,6 +42,7 @@ class ParserLoadLog(TimestampMixin, models.Model): _("Глобальный адресный справочник ГосЭДО"), ) MEDIA_NEWS = "media_news", _("Новости СМИ") + ROPK_SANCTIONS = "ropk_sanctions", _("РОПК — Санкции") class Status(models.TextChoices): SUCCESS = "success", _("Успешно") diff --git a/src/apps/parsers/ropk_sanctions.py b/src/apps/parsers/ropk_sanctions.py new file mode 100644 index 0000000..453cdbd --- /dev/null +++ b/src/apps/parsers/ropk_sanctions.py @@ -0,0 +1,433 @@ +"""Validated XLSX ingestion for the ROPK organization sanctions snapshot.""" + +from __future__ import annotations + +import hashlib +import uuid +import zipfile +from collections import Counter +from dataclasses import dataclass +from pathlib import Path +from typing import BinaryIO + +from apps.parsers.models import ParserSourceArtifact, ParserStagedRecord +from django.core.files import File +from django.db import transaction +from django.db.models import Q +from openpyxl import load_workbook +from organizations.models import ( + Organization, + OrganizationSourceRecord, + SanctionsExtension, +) +from organizations.resolver import OrganizationDirectoryResolver +from organizations.source_cache import invalidate_source_data_cache +from organizations.source_ingestion import ( + OrganizationSourceIngestionService, + SourceRecordInput, +) + +ROPK_SANCTIONS_SOURCE = "ropk_sanctions" +ROPK_SANCTIONS_RECORD_TYPE = "organization_sanctions" +ROPK_SANCTIONS_MIME = ( + "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet" +) +ROPK_SANCTIONS_MAX_BYTES = 25 * 1024 * 1024 +ROPK_SANCTIONS_MAX_UNCOMPRESSED_BYTES = 100 * 1024 * 1024 +ROPK_SANCTIONS_MAX_ARCHIVE_ENTRIES = 512 +ROPK_SANCTIONS_UUID_NAMESPACE = uuid.UUID("b1d47229-f3ac-41a3-8a2b-2c8827544321") +ROPK_SANCTIONS_HEADERS = ( + "rn", + "ogrn", + "inn", + "okpo", + "Санкции - Великобритания HM Treasury", + "Санкции - Евросоюз", + "Санкции - США", + "Санкции - Швейцария", + "Санкции секторальные - США", + "Санкции - Великобритания UKSL", + "Санкции - Украина", +) +ROPK_SANCTIONS_FLAG_FIELDS = ( + "uk_hm_treasury", + "european_union", + "united_states", + "switzerland", + "united_states_sectoral", + "uk_uksl", + "ukraine", +) + + +class RopkSanctionsValidationError(ValueError): + """The uploaded workbook is unsafe or violates the source contract.""" + + +@dataclass(frozen=True) +class RopkSanctionsImportResult: + parsed: int + published: int + quarantined: int + reasons: dict[str, int] + skipped: bool = False + + +@dataclass +class _ParsedSanctionsRow: + row_number: int + rn: str + ogrn: str + inn: str + okpo: str + flags: dict[str, bool] + raw_data: dict[str, str] + reason: str = "" + + +def stable_ropk_sanctions_uid(rn: str) -> uuid.UUID: + """Return a stable source-record identifier for one sanctions row.""" + return uuid.uuid5(ROPK_SANCTIONS_UUID_NAMESPACE, f"{ROPK_SANCTIONS_SOURCE}:{rn}") + + +def _cell_string(cell) -> str: + value = cell.value + if value is None or isinstance(value, bool): + return "" + if isinstance(value, float): + if not value.is_integer(): + return "" + value = int(value) + text = str(value).strip() + number_format = str(cell.number_format or "") + if text.isdigit() and number_format and set(number_format) <= {"0"}: + text = text.zfill(len(number_format)) + return text + + +def _identifier(cell, *, lengths: set[int]) -> str: + value = _cell_string(cell) + if not value.isdigit() or len(value) not in lengths: + return "" + return value + + +def _flag(value: object) -> bool | None: + if isinstance(value, bool) or value is None: + return None + if isinstance(value, int | float): + if value == 0: + return False + if value == 1: + return True + return None + text = str(value).strip() + if text == "0": + return False + if text == "1": + return True + return None + + +def _validate_xlsx_archive(handle: BinaryIO) -> None: + handle.seek(0) + try: + with zipfile.ZipFile(handle) as archive: + entries = archive.infolist() + if not entries or len(entries) > ROPK_SANCTIONS_MAX_ARCHIVE_ENTRIES: + raise RopkSanctionsValidationError("unsafe_xlsx_archive") + if any(entry.flag_bits & 0x1 for entry in entries): + raise RopkSanctionsValidationError("encrypted_xlsx_not_supported") + if ( + sum(entry.file_size for entry in entries) + > ROPK_SANCTIONS_MAX_UNCOMPRESSED_BYTES + ): + raise RopkSanctionsValidationError("xlsx_uncompressed_size_exceeded") + if archive.testzip() is not None: + raise RopkSanctionsValidationError("corrupt_xlsx_archive") + except zipfile.BadZipFile as exc: + raise RopkSanctionsValidationError("invalid_xlsx") from exc + finally: + handle.seek(0) + + +def _resolve_organization( + *, inn: str, ogrn: str, okpo: str +) -> tuple[Organization | None, str]: + directory = OrganizationDirectoryResolver._directory_queryset() + if bool(inn) != bool(ogrn): + return None, "missing_required_value" + + if inn and ogrn: + candidates = list(directory.filter(inn=inn, ogrn=ogrn)[:2]) + if len(candidates) > 1: + return None, "organization_ambiguous" + if not candidates: + has_related_identity = directory.filter( + Q(inn=inn) | Q(ogrn=ogrn) | Q(okpo=okpo) + ).exists() + return None, ( + "identifier_mismatch" + if has_related_identity + else "organization_not_found" + ) + organization = candidates[0] + if organization.okpo != okpo: + return None, "identifier_mismatch" + else: + candidates = list(directory.filter(okpo=okpo)[:2]) + if len(candidates) > 1: + return None, "organization_ambiguous" + if not candidates: + return None, "organization_not_found" + organization = candidates[0] + + if not all( + ( + organization.name.strip(), + organization.inn, + organization.ogrn, + organization.okpo, + ) + ): + return None, "missing_required_value" + if organization.okpo != okpo: + return None, "identifier_mismatch" + return organization, "" + + +def _parse_row(row_number: int, cells: tuple) -> _ParsedSanctionsRow: + cells = cells[: len(ROPK_SANCTIONS_HEADERS)] + raw_data = { + header: str(cell.value or "") + for header, cell in zip(ROPK_SANCTIONS_HEADERS, cells, strict=True) + } + reason = ( + "formula_not_allowed" if any(cell.data_type == "f" for cell in cells) else "" + ) + rn = _cell_string(cells[0]) + ogrn = _identifier(cells[1], lengths={13}) + inn = _identifier(cells[2], lengths={10, 12}) + okpo = _identifier(cells[3], lengths={8, 14}) + if not reason and not rn: + reason = "missing_required_value" + if not reason and not okpo: + reason = "invalid_okpo" + if not reason and cells[1].value not in (None, "") and not ogrn: + reason = "invalid_ogrn" + if not reason and cells[2].value not in (None, "") and not inn: + reason = "invalid_inn" + + flags: dict[str, bool] = {} + for field_name, cell in zip( + ROPK_SANCTIONS_FLAG_FIELDS, + cells[4:], + strict=True, + ): + parsed_flag = _flag(cell.value) + if parsed_flag is None: + reason = reason or "invalid_flag" + else: + flags[field_name] = parsed_flag + + return _ParsedSanctionsRow( + row_number=row_number, + rn=rn, + ogrn=ogrn, + inn=inn, + okpo=okpo, + flags=flags, + raw_data=raw_data, + reason=reason, + ) + + +def _parse_rows(sheet) -> list[_ParsedSanctionsRow]: + header_cells = next(sheet.iter_rows(min_row=1, max_row=1), ()) + headers = tuple(str(cell.value or "").strip() for cell in header_cells) + if headers != ROPK_SANCTIONS_HEADERS: + raise RopkSanctionsValidationError("invalid_headers") + + rows = [ + _parse_row(row_number, cells) + for row_number, cells in enumerate(sheet.iter_rows(min_row=2), start=2) + if any(cell.value not in (None, "") for cell in cells) + ] + if not rows: + raise RopkSanctionsValidationError("empty_source") + + duplicate_rns = { + rn + for rn, count in Counter(row.rn for row in rows if row.rn).items() + if count > 1 + } + for row in rows: + if not row.reason and row.rn in duplicate_rns: + row.reason = "duplicate_rn" + return rows + + +def import_ropk_sanctions( # noqa: C901 + *, + handle: BinaryIO, + original_name: str, + load_batch: int, + uploaded_by_id: int | None, +) -> tuple[ParserSourceArtifact, RopkSanctionsImportResult]: + """Validate, stage and atomically publish one ROPK sanctions snapshot.""" + handle.seek(0, 2) + size_bytes = handle.tell() + if size_bytes > ROPK_SANCTIONS_MAX_BYTES: + raise RopkSanctionsValidationError("source_too_large") + handle.seek(0) + digest = hashlib.sha256() + for chunk in iter(lambda: handle.read(128 * 1024), b""): + digest.update(chunk) + handle.seek(0) + sha256 = digest.hexdigest() + safe_original_name = Path(str(original_name).replace("\\", "/")).name + artifact = ParserSourceArtifact.objects.create( + source=ROPK_SANCTIONS_SOURCE, + version=sha256, + sha256=sha256, + content_type=ROPK_SANCTIONS_MIME, + size_bytes=size_bytes, + original_name=safe_original_name, + load_batch=load_batch, + uploaded_by_id=uploaded_by_id, + ) + workbook = None + try: + artifact.file.save(artifact.original_name, File(handle), save=True) + duplicate = ( + ParserSourceArtifact.objects.filter( + source=ROPK_SANCTIONS_SOURCE, + sha256=sha256, + status=ParserSourceArtifact.Status.PUBLISHED, + ) + .exclude(uid=artifact.uid) + .order_by("-created_at") + .first() + ) + if duplicate is not None: + artifact.status = ParserSourceArtifact.Status.SKIPPED + artifact.metadata = {"duplicate_of": str(duplicate.uid)} + artifact.save(update_fields=["status", "metadata", "updated_at"]) + return artifact, RopkSanctionsImportResult(0, 0, 0, {}, skipped=True) + + _validate_xlsx_archive(handle) + workbook = load_workbook( + handle, + read_only=True, + data_only=False, + keep_links=False, + ) + if len(workbook.worksheets) != 1: + raise RopkSanctionsValidationError("invalid_worksheet_count") + rows = _parse_rows(workbook.worksheets[0]) + artifact.status = ParserSourceArtifact.Status.PARSED + artifact.save(update_fields=["status", "updated_at"]) + + reasons: Counter[str] = Counter() + staged: list[ParserStagedRecord] = [] + inputs: list[SourceRecordInput] = [] + for row in rows: + organization = None + reason = row.reason + if not reason: + organization, reason = _resolve_organization( + inn=row.inn, + ogrn=row.ogrn, + okpo=row.okpo, + ) + if reason: + reasons[reason] += 1 + else: + assert organization is not None + payload = {"rn": row.rn, **row.flags} + inputs.append( + SourceRecordInput( + uid=stable_ropk_sanctions_uid(row.rn), + organization_uid=organization.uid, + external_id=row.rn, + record_type=ROPK_SANCTIONS_RECORD_TYPE, + title=organization.name, + organization_name=organization.name, + inn=organization.inn, + ogrn=organization.ogrn, + status="active", + payload=payload, + ) + ) + staged.append( + ParserStagedRecord( + artifact=artifact, + row_number=row.row_number, + external_id=row.rn, + record_type=ROPK_SANCTIONS_RECORD_TYPE, + raw_data=row.raw_data, + normalized_data={ + "rn": row.rn, + "ogrn": row.ogrn or None, + "inn": row.inn or None, + "okpo": row.okpo or None, + **row.flags, + }, + organization=organization, + disposition=( + ParserStagedRecord.Disposition.QUARANTINED + if reason + else ParserStagedRecord.Disposition.STAGED + ), + reason=reason, + ) + ) + ParserStagedRecord.objects.bulk_create(staged, batch_size=500) + + result = RopkSanctionsImportResult( + parsed=len(rows), + published=len(inputs), + quarantined=sum(reasons.values()), + reasons=dict(reasons), + ) + keep_uids = [record.uid for record in inputs if record.uid] + with transaction.atomic(): + ingestion = OrganizationSourceIngestionService.save_records( + source=ROPK_SANCTIONS_SOURCE, + load_batch=load_batch, + records=inputs, + ) + if ingestion.unresolved: + raise RopkSanctionsValidationError("organization_resolution_changed") + OrganizationSourceRecord.objects.filter( + source=ROPK_SANCTIONS_SOURCE + ).exclude(uid__in=keep_uids).delete() + SanctionsExtension.objects.filter(records__isnull=True).delete() + ParserStagedRecord.objects.filter( + artifact=artifact, + disposition=ParserStagedRecord.Disposition.STAGED, + ).update(disposition=ParserStagedRecord.Disposition.PUBLISHED) + artifact.status = ParserSourceArtifact.Status.PUBLISHED + artifact.parsed_count = result.parsed + artifact.published_count = result.published + artifact.quarantined_count = result.quarantined + artifact.rejection_reasons = result.reasons + artifact.save( + update_fields=[ + "status", + "parsed_count", + "published_count", + "quarantined_count", + "rejection_reasons", + "updated_at", + ] + ) + invalidate_source_data_cache() + return artifact, result + except Exception: + artifact.status = ParserSourceArtifact.Status.REJECTED + artifact.save(update_fields=["status", "updated_at"]) + raise + finally: + if workbook is not None: + workbook.close() diff --git a/src/apps/parsers/serializers.py b/src/apps/parsers/serializers.py index c33f30e..72139c2 100644 --- a/src/apps/parsers/serializers.py +++ b/src/apps/parsers/serializers.py @@ -735,6 +735,7 @@ class ParserRunResponseSerializer(serializers.Serializer): task_ids = serializers.ListField(child=serializers.CharField()) source = serializers.CharField() task_name = serializers.CharField() + status = serializers.CharField(required=False) # ============================================================================= diff --git a/src/apps/parsers/source_cards.py b/src/apps/parsers/source_cards.py index 903e15a..3e83878 100644 --- a/src/apps/parsers/source_cards.py +++ b/src/apps/parsers/source_cards.py @@ -16,6 +16,7 @@ from apps.parsers.models import ( VACANCY_RECORD_SOURCES, ParserLoadLog, ) +from apps.parsers.source_registry import get_source_by_model_source from django.conf import settings from django.core.cache import cache from django.db.models import Count, Max, Q @@ -34,7 +35,7 @@ ACTIVE_JOB_STATUSES = [JobStatus.PENDING, JobStatus.STARTED, JobStatus.RETRY] STALE_ACTIVE_MAX_AGE_MINUTES = 4 * 60 STALE_PENDING_MAX_AGE_MINUTES = 24 * 60 SOURCE_CARD_STATS_CACHE_TIMEOUT_SECONDS = 7 * 24 * 60 * 60 -SOURCE_CARD_STATS_CACHE_SCHEMA_VERSION = 2 +SOURCE_CARD_STATS_CACHE_SCHEMA_VERSION = 3 @dataclass(frozen=True) @@ -294,6 +295,24 @@ SOURCE_CARD_DEFINITIONS: tuple[SourceCardDefinition, ...] = ( ), refresh_interval=timedelta(days=1), ), + SourceCardDefinition( + slug="ropk-sanctions", + title="РОПК — Санкции", + description="Санкционные признаки организаций по семи перечням.", + order=70, + task_names=("apps.parsers.tasks.parse_ropk_sanctions",), + source_items=( + SourceItemDefinition( + code="ropk_sanctions", + title="РОПК — Санкции", + description="Санкционные признаки организаций из XLSX-выгрузки РОПК.", + parser_source=ParserLoadLog.Source.ROPK_SANCTIONS, + refresh_key="ropk_sanctions", + ), + ), + supports_refresh=False, + upload_url="/api/v1/parsers/upload/ropk_sanctions/", + ), SourceCardDefinition( slug="arbitration-cases", title="Арбитражные дела", @@ -1246,8 +1265,9 @@ class SourceCardService: else last_updated_at ), "upload_url": ( - "/api/v1/parsers/upload/media_news/" - if item.parser_source == ParserLoadLog.Source.MEDIA_NEWS + descriptor.upload_url + if item.parser_source + and (descriptor := get_source_by_model_source(item.parser_source)) else "" ), "latest_load": cls._serialize_load_log(latest_load), diff --git a/src/apps/parsers/source_registry.py b/src/apps/parsers/source_registry.py index d7a533f..40427ad 100644 --- a/src/apps/parsers/source_registry.py +++ b/src/apps/parsers/source_registry.py @@ -329,6 +329,25 @@ PARSER_SOURCES: dict[str, ParserSourceDescriptor] = { supports_refresh=False, admin_only=True, ), + "ropk_sanctions": ParserSourceDescriptor( + key="ropk_sanctions", + source=ParserLoadLog.Source.ROPK_SANCTIONS, + title="РОПК — Санкции", + agency="РОПК", + data_scope="Санкционные признаки организаций по семи перечням", + task_name="apps.parsers.tasks.parse_ropk_sanctions", + mode="manual_upload", + access_method="admin_upload", + parser_strategy="validated_xlsx_snapshot", + source_notes=( + "XLSX загружается вручную; строки связываются только с организациями " + "канонического реестра, неразрешенные строки сохраняются в карантине." + ), + supports_file_upload=True, + upload_route="parsers/upload/ropk_sanctions", + supports_refresh=False, + admin_only=True, + ), } diff --git a/src/apps/parsers/tasks.py b/src/apps/parsers/tasks.py index 729c1bb..a9be7ce 100644 --- a/src/apps/parsers/tasks.py +++ b/src/apps/parsers/tasks.py @@ -56,6 +56,7 @@ from apps.parsers.clients.zakupki import ZakupkiClient from apps.parsers.gosedo import GosedoNotModified, refresh_gosedo from apps.parsers.media_news import import_media_news from apps.parsers.models import CheckoCollectionAttempt, ParserLoadLog +from apps.parsers.ropk_sanctions import import_ropk_sanctions from apps.parsers.services import ( FNSReportOrganizationResolutionSkipped, FNSReportService, @@ -3725,6 +3726,90 @@ def parse_media_news( default_storage.delete(file_path) +@shared_task(bind=True, soft_time_limit=15 * 60, time_limit=20 * 60) +def parse_ropk_sanctions( + self, + *, + file_path: str, + original_name: str | None = None, + requested_by_id: int | None = None, +) -> dict: + """Validate and atomically publish a ROPK sanctions XLSX snapshot.""" + from django.core.files.storage import default_storage + + source = ParserLoadLog.Source.ROPK_SANCTIONS + task_name = "apps.parsers.tasks.parse_ropk_sanctions" + load_log, batch_id = ParserLoadLogService.create_load_log_with_next_batch_id( + source=source, + status=ParserLoadLog.Status.IN_PROGRESS, + ) + task_id = self.request.id or str(uuid.uuid4()) + job = _get_or_create_background_job( + task_id=task_id, + task_name=task_name, + source=source, + batch_id=batch_id, + requested_by_id=requested_by_id, + meta={"source_key": source}, + ) + job.mark_started() + job.update_progress(10, "Проверка файла РОПК — Санкции...") + try: + with default_storage.open(file_path, "rb") as handle: + artifact, result = import_ropk_sanctions( + handle=handle, + original_name=original_name or Path(file_path).name, + load_batch=batch_id, + uploaded_by_id=requested_by_id, + ) + if result.skipped: + ParserLoadLogService.update( + load_log, + status=ParserLoadLog.Status.SKIPPED, + error_message="Файл с таким checksum уже опубликован", + ) + payload = { + "status": "skipped", + "batch_id": batch_id, + "artifact_id": str(artifact.uid), + "raw_records_count": 0, + "published_records_count": 0, + "quarantine_records_count": 0, + } + job.update_progress(100, "Файл уже был опубликован") + job.complete(result=payload) + return payload + + ParserLoadLogService.update( + load_log, + status=ParserLoadLog.Status.SUCCESS, + records_count=result.published, + ) + payload = { + "status": "success", + "batch_id": batch_id, + "load_id": load_log.id, + "artifact_id": str(artifact.uid), + "raw_records_count": result.parsed, + "published_records_count": result.published, + "quarantine_records_count": result.quarantined, + "rejection_reasons": result.reasons, + } + job.update_progress(100, "Санкционные признаки РОПК загружены") + job.complete(result=payload) + return payload + except Exception as exc: + logger.error("ROPK sanctions import failed: %s", exc, exc_info=True) + ParserLoadLogService.mark_failed(load_log, str(exc)) + job.fail(error=str(exc)) + raise + finally: + if file_path.startswith("parser_uploads/") and default_storage.exists( + file_path + ): + default_storage.delete(file_path) + + @shared_task def cleanup_source_artifacts() -> dict: """Apply the 90-day/minimum-ten-versions parser artifact retention policy.""" diff --git a/src/apps/parsers/views.py b/src/apps/parsers/views.py index 993dce4..f1ac5fb 100644 --- a/src/apps/parsers/views.py +++ b/src/apps/parsers/views.py @@ -34,6 +34,10 @@ from apps.parsers.models import ( ProcurementRecord, ) from apps.parsers.organization_enrichment import enrich_parser_result_rows +from apps.parsers.ropk_sanctions import ( + ROPK_SANCTIONS_MAX_BYTES, + ROPK_SANCTIONS_MIME, +) from apps.parsers.serializers import ( FinancialReportDetailSerializer, FinancialReportSerializer, @@ -158,6 +162,7 @@ TASKS_BY_NAME = { tasks.parse_gosedo_address_directory ), "apps.parsers.tasks.parse_media_news": tasks.parse_media_news, + "apps.parsers.tasks.parse_ropk_sanctions": tasks.parse_ropk_sanctions, } PARSER_SOURCE_ALIASES = { @@ -231,6 +236,7 @@ EXISTING_TASK_PARAMS = { "fns_financial": {"requested_by_id"}, "gosedo_address_directory": {"requested_by_id"}, "media_news": {"file_path", "original_name", "requested_by_id"}, + "ropk_sanctions": {"file_path", "original_name", "requested_by_id"}, } @@ -382,6 +388,13 @@ MEDIA_NEWS_UPLOAD_FILE_PARAM = openapi.Parameter( type=openapi.TYPE_FILE, required=True, ) +ROPK_SANCTIONS_UPLOAD_FILE_PARAM = openapi.Parameter( + "file", + openapi.IN_FORM, + description="XLSX-файл РОПК — Санкции размером не более 25 МиБ", + type=openapi.TYPE_FILE, + required=True, +) RESULT_LIST_PARAMS = [ PAGE_PARAM, PAGE_SIZE_PARAM, @@ -1390,6 +1403,17 @@ class SourceCardRefreshView(APIView): }, ) def post(self, request, slug: str): + definition = SourceCardService.get_definition(slug) + if slug == "ropk-sanctions" and not definition.supports_refresh: + return api_error_response( + [ + { + "code": "refresh_not_supported", + "message": "Источник обновляется загрузкой файла", + } + ], + status_code=status.HTTP_405_METHOD_NOT_ALLOWED, + ) serializer = SourceCardRefreshRequestSerializer(data=request.data) serializer.is_valid(raise_exception=True) @@ -2712,6 +2736,16 @@ class ParserRunView(APIView): if descriptor.admin_only and not request.user.is_staff: return Response(status=status.HTTP_403_FORBIDDEN) if not descriptor.supports_refresh: + if canonical_source_key == ParserLoadLog.Source.ROPK_SANCTIONS: + return api_error_response( + [ + { + "code": "parser_run_not_supported", + "message": "Источник обновляется загрузкой файла", + } + ], + status_code=status.HTTP_405_METHOD_NOT_ALLOWED, + ) return api_error_response( [ { @@ -2772,11 +2806,93 @@ class ParserRunView(APIView): "task_ids": [async_result.id], "source": descriptor.source, "task_name": descriptor.task_name, + "status": JobStatus.PENDING, }, status_code=status.HTTP_202_ACCEPTED, ) +def _media_news_upload_error(uploaded_file) -> Response | None: + if not uploaded_file.name.lower().endswith(".xlsx"): + return api_error_response( + [ + { + "code": "invalid_file_type", + "field": "file", + "message": "Для новостей требуется XLSX", + } + ], + status_code=status.HTTP_400_BAD_REQUEST, + ) + if uploaded_file.size > MEDIA_NEWS_MAX_BYTES: + return api_error_response( + [ + { + "code": "file_too_large", + "field": "file", + "message": "Размер XLSX для новостей не должен превышать 25 МиБ", + } + ], + status_code=status.HTTP_400_BAD_REQUEST, + ) + return None + + +def _ropk_sanctions_upload_error(uploaded_file, *, task_name: str) -> Response | None: + if ( + not uploaded_file.name.lower().endswith(".xlsx") + or uploaded_file.content_type != ROPK_SANCTIONS_MIME + ): + return api_error_response( + [ + { + "code": "unsupported_media_type", + "field": "file", + "message": "Для источника РОПК — Санкции требуется XLSX", + } + ], + status_code=status.HTTP_415_UNSUPPORTED_MEDIA_TYPE, + ) + if uploaded_file.size > ROPK_SANCTIONS_MAX_BYTES: + return api_error_response( + [ + { + "code": "file_too_large", + "field": "file", + "message": "Размер XLSX не должен превышать 25 МиБ", + } + ], + status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, + ) + now = timezone.now() + active_job = ( + BackgroundJobService.get_queryset() + .filter(task_name=task_name) + .filter( + Q( + status=JobStatus.PENDING, + created_at__gte=now - timedelta(hours=24), + ) + | Q( + status__in=[JobStatus.STARTED, JobStatus.RETRY], + updated_at__gte=now - timedelta(hours=4), + ) + ) + .exists() + ) + if active_job: + return api_error_response( + [ + { + "code": "upload_already_running", + "message": "Загрузка источника уже выполняется", + } + ], + status_code=status.HTTP_409_CONFLICT, + ) + return None + + class ParserUploadView(APIView): """Ручная загрузка файла только для источников с supports_file_upload.""" @@ -2799,32 +2915,34 @@ class ParserUploadView(APIView): ], status_code=status.HTTP_400_BAD_REQUEST, ) + if ( + source_key == ParserLoadLog.Source.ROPK_SANCTIONS + and "file" not in request.FILES + ): + return api_error_response( + [ + { + "code": "file_required", + "field": "file", + "message": "Файл обязателен", + } + ], + status_code=status.HTTP_400_BAD_REQUEST, + ) serializer = ParserUploadRequestSerializer(data=request.data) serializer.is_valid(raise_exception=True) uploaded_file = serializer.validated_data["file"] if source_key == "media_news": - if not uploaded_file.name.lower().endswith(".xlsx"): - return api_error_response( - [ - { - "code": "invalid_file_type", - "field": "file", - "message": "Для новостей требуется XLSX", - } - ], - status_code=status.HTTP_400_BAD_REQUEST, - ) - if uploaded_file.size > MEDIA_NEWS_MAX_BYTES: - return api_error_response( - [ - { - "code": "file_too_large", - "field": "file", - "message": "Размер XLSX для новостей не должен превышать 25 МиБ", - } - ], - status_code=status.HTTP_400_BAD_REQUEST, - ) + error_response = _media_news_upload_error(uploaded_file) + if error_response is not None: + return error_response + if source_key == ParserLoadLog.Source.ROPK_SANCTIONS: + error_response = _ropk_sanctions_upload_error( + uploaded_file, + task_name=descriptor.task_name, + ) + if error_response is not None: + return error_response file_path = _save_uploaded_parser_file(uploaded_file) run_serializer = ParserRunRequestSerializer(data={"file_path": file_path}) run_serializer.is_valid(raise_exception=True) @@ -2853,6 +2971,7 @@ class ParserUploadView(APIView): "task_ids": [async_result.id], "source": descriptor.source, "task_name": descriptor.task_name, + "status": JobStatus.PENDING, }, status_code=status.HTTP_202_ACCEPTED, ) diff --git a/src/organizations/filters.py b/src/organizations/filters.py index 92ce9d3..b46cbea 100644 --- a/src/organizations/filters.py +++ b/src/organizations/filters.py @@ -29,6 +29,8 @@ SOURCE_FILTER_ALIASES = { "fstec": SourceGroup.SECURITY_REGISTRIES, "vacancies": SourceGroup.VACANCIES, "trudvsem": SourceGroup.VACANCIES, + "sanctions": SourceGroup.SANCTIONS, + "ropk_sanctions": SourceGroup.SANCTIONS, } diff --git a/src/organizations/migrations/0011_auto_20260820_1325.py b/src/organizations/migrations/0011_auto_20260820_1325.py new file mode 100644 index 0000000..645b323 --- /dev/null +++ b/src/organizations/migrations/0011_auto_20260820_1325.py @@ -0,0 +1,31 @@ +# Generated by Django 3.2.25 on 2026-08-20 13:25 + +from django.db import migrations, models +import django.db.models.deletion + + +class Migration(migrations.Migration): + + dependencies = [ + ('organizations', '0010_gosedo_media_source_groups'), + ] + + operations = [ + migrations.CreateModel( + name='SanctionsExtension', + fields=[ + ('organizationsourceextension_ptr', models.OneToOneField(auto_created=True, on_delete=django.db.models.deletion.CASCADE, parent_link=True, primary_key=True, serialize=False, to='organizations.organizationsourceextension')), + ], + options={ + 'verbose_name': 'санкционные признаки', + 'verbose_name_plural': 'санкционные признаки', + 'db_table': 'organizations_sanctions_extension', + }, + bases=('organizations.organizationsourceextension',), + ), + migrations.AlterField( + model_name='organizationsourceextension', + name='source_group', + field=models.CharField(choices=[('financial_indicators', 'Финансово-экономические показатели'), ('government_procurements', 'Государственные закупки'), ('industrial_production', 'Производители и продукция России'), ('planned_inspections', 'Плановые проверки'), ('bankruptcy', 'Сведения о процедурах банкротства'), ('defense_suppliers', 'Недобросовестные поставщики ГОЗ'), ('arbitration', 'Арбитражные дела'), ('security_registries', 'Реестры по информационной безопасности'), ('vacancies', 'Вакансии'), ('electronic_document_exchange', 'Электронный документооборот'), ('media_mentions', 'Упоминания в СМИ'), ('sanctions', 'Санкции')], db_index=True, max_length=64, verbose_name='группа источников'), + ), + ] diff --git a/src/organizations/models.py b/src/organizations/models.py index 9cb8ea9..ce99102 100644 --- a/src/organizations/models.py +++ b/src/organizations/models.py @@ -37,6 +37,7 @@ class SourceGroup(models.TextChoices): _("Электронный документооборот"), ) MEDIA_MENTIONS = "media_mentions", _("Упоминания в СМИ") + SANCTIONS = "sanctions", _("Санкции") class SourceExtensionStatus(models.TextChoices): @@ -691,6 +692,17 @@ class MediaMentionExtension(OrganizationSourceExtension): verbose_name_plural = _("упоминания в СМИ") +class SanctionsExtension(OrganizationSourceExtension): + """Organization sanctions snapshot linked to the canonical directory.""" + + source_group_value = SourceGroup.SANCTIONS + + class Meta: + db_table = "organizations_sanctions_extension" + verbose_name = _("санкционные признаки") + verbose_name_plural = _("санкционные признаки") + + class OrganizationSourceRecord(models.Model): """Subordinate source record stored under a source extension.""" diff --git a/src/organizations/serializers.py b/src/organizations/serializers.py index 76dd1cb..9404dea 100644 --- a/src/organizations/serializers.py +++ b/src/organizations/serializers.py @@ -82,6 +82,14 @@ class OrganizationSourceRecordPayloadSerializer(serializers.Serializer): """Typed optional fields exposed by source-record list payloads.""" artifact_id = serializers.UUIDField(read_only=True, allow_null=True) + rn = serializers.CharField(read_only=True, allow_blank=True, allow_null=True) + uk_hm_treasury = serializers.BooleanField(read_only=True, allow_null=True) + european_union = serializers.BooleanField(read_only=True, allow_null=True) + united_states = serializers.BooleanField(read_only=True, allow_null=True) + switzerland = serializers.BooleanField(read_only=True, allow_null=True) + united_states_sectoral = serializers.BooleanField(read_only=True, allow_null=True) + uk_uksl = serializers.BooleanField(read_only=True, allow_null=True) + ukraine = serializers.BooleanField(read_only=True, allow_null=True) inn = serializers.CharField(read_only=True, allow_blank=True, allow_null=True) kpp = serializers.CharField(read_only=True, allow_blank=True, allow_null=True) ogrn = serializers.CharField(read_only=True, allow_blank=True, allow_null=True) diff --git a/src/organizations/source_groups.py b/src/organizations/source_groups.py index 27b6371..052e259 100644 --- a/src/organizations/source_groups.py +++ b/src/organizations/source_groups.py @@ -17,6 +17,7 @@ from organizations.models import ( MediaMentionExtension, OrganizationSourceExtension, PlannedInspectionExtension, + SanctionsExtension, SecurityRegistryExtension, SourceGroup, VacancyExtension, @@ -154,6 +155,13 @@ SOURCE_GROUP_DESCRIPTORS: dict[str, SourceGroupDescriptor] = { title="Упоминания в СМИ", extension_model=MediaMentionExtension, ), + ParserLoadLog.Source.ROPK_SANCTIONS: SourceGroupDescriptor( + source=ParserLoadLog.Source.ROPK_SANCTIONS, + source_group=SourceGroup.SANCTIONS, + record_type="organization_sanctions", + title="РОПК — Санкции", + extension_model=SanctionsExtension, + ), "hh": SourceGroupDescriptor( source="hh", source_group=SourceGroup.VACANCIES, diff --git a/src/organizations/source_ingestion.py b/src/organizations/source_ingestion.py index d5c1375..ddda17d 100644 --- a/src/organizations/source_ingestion.py +++ b/src/organizations/source_ingestion.py @@ -51,6 +51,7 @@ class SourceRecordInput: organization_name: str record_type: str = "" uid: UUID | None = None + organization_uid: UUID | None = None inn: str = "" kpp: str = "" ogrn: str = "" @@ -84,6 +85,7 @@ class OrganizationSourceIngestionResult: class _NormalizedRecordInput: index: int record: SourceRecordInput + organization_uid: UUID | None rn: str okpo: str inn: str @@ -240,6 +242,7 @@ class OrganizationSourceIngestionService: _NormalizedRecordInput( index=index, record=record_input, + organization_uid=record_input.organization_uid, rn=OrganizationDirectoryResolver._digits( (record_input.payload or {}).get("rn") or (record_input.payload or {}).get("organization_rn"), @@ -268,7 +271,26 @@ class OrganizationSourceIngestionService: skipped_unmatched = 0 skipped_ambiguous = 0 + requested_uids = { + record.organization_uid + for record in normalized_records + if record.organization_uid is not None + } + organizations_by_uid = { + organization.uid: organization + for organization in OrganizationDirectoryResolver._directory_queryset().filter( + uid__in=requested_uids + ) + } + for record in normalized_records: + if record.organization_uid is not None: + organization = organizations_by_uid.get(record.organization_uid) + if organization is not None: + organizations_by_index[record.index] = organization + else: + skipped_unmatched += 1 + continue result = OrganizationDirectoryResolver.resolve( OrganizationDirectoryResolver.identity( rn=record.rn, diff --git a/src/organizations/source_record_export.py b/src/organizations/source_record_export.py index 8c23c4a..98e8843 100644 --- a/src/organizations/source_record_export.py +++ b/src/organizations/source_record_export.py @@ -38,7 +38,10 @@ EXPORT_FORMAT_XLSX = "xlsx" EXPORT_FORMAT_JSON = "json" EXPORT_FORMATS = (EXPORT_FORMAT_CSV, EXPORT_FORMAT_XLSX, EXPORT_FORMAT_JSON) FINANCIAL_SOURCE_GROUP = SourceGroup.FINANCIAL_INDICATORS.value -ALL_HISTORY_SOURCE_GROUPS = {SourceGroup.MEDIA_MENTIONS.value} +ALL_HISTORY_SOURCE_GROUPS = { + SourceGroup.MEDIA_MENTIONS.value, + SourceGroup.SANCTIONS.value, +} EXPORT_MANIFEST_VERSION = 2 CURRENT_EXPORT_MANIFEST_FILE_NAME = "current.json" GENERATION_MANIFEST_FILE_NAME = "manifest.json" @@ -72,6 +75,7 @@ SOURCE_GROUP_EXPORT_FILE_STEMS: dict[str, str] = { SourceGroup.VACANCIES.value: "labor-vacancies", SourceGroup.ELECTRONIC_DOCUMENT_EXCHANGE.value: "gosedo-address-directory", SourceGroup.MEDIA_MENTIONS.value: "media-mentions", + SourceGroup.SANCTIONS.value: "ropk-sanctions", } ORGANIZATION_EXPORT_FIELDS = ["Наименование", "ИНН", "ОГРН", "КПП", "ОКПО"] @@ -90,6 +94,29 @@ SOURCE_RECORD_EXPORT_FIELDS = [ "created_at", "updated_at", ] +SANCTIONS_EXPORT_FIELDS = [ + "Наименование", + "rn", + "ogrn", + "inn", + "okpo", + "Санкции - Великобритания HM Treasury", + "Санкции - Евросоюз", + "Санкции - США", + "Санкции - Швейцария", + "Санкции секторальные - США", + "Санкции - Великобритания UKSL", + "Санкции - Украина", +] +SANCTIONS_EXPORT_FLAG_FIELDS = { + "Санкции - Великобритания HM Treasury": "uk_hm_treasury", + "Санкции - Евросоюз": "european_union", + "Санкции - США": "united_states", + "Санкции - Швейцария": "switzerland", + "Санкции секторальные - США": "united_states_sectoral", + "Санкции - Великобритания UKSL": "uk_uksl", + "Санкции - Украина": "ukraine", +} FINANCIAL_LINE_FIELDS = [ "id", "form_code", @@ -641,6 +668,9 @@ def _spool_source_group_rows( output.write("\n") output.write("]") + if source_group == SourceGroup.SANCTIONS: + return SANCTIONS_EXPORT_FIELDS, records_count + return [ *ORGANIZATION_EXPORT_FIELDS, *SOURCE_RECORD_EXPORT_FIELDS, @@ -696,7 +726,7 @@ def _render_csv_file( writer.writeheader() for row in _iter_spooled_rows(row_spool_path): writer.writerow( - {key: _serialize_flat_value(row.get(key)) for key in headers} + {key: _serialize_spreadsheet_value(row.get(key)) for key in headers} ) @@ -740,7 +770,9 @@ def _render_xlsx_files( part_number += 1 rows_in_file = 0 workbook, worksheet = _new_export_workbook(headers) - worksheet.append([_serialize_flat_value(row.get(key)) for key in headers]) + worksheet.append( + [_serialize_spreadsheet_value(row.get(key)) for key in headers] + ) rows_in_file += 1 workbook.save(part_paths[part_number - 1]) @@ -784,6 +816,21 @@ def _build_record_row( include_financial_lines: bool, ) -> dict[str, Any]: organization = record.extension.organization + if record.extension.source_group == SourceGroup.SANCTIONS: + payload = record.payload if isinstance(record.payload, dict) else {} + return { + "Наименование": organization.full_name + or organization.short_name + or organization.name, + "rn": str(payload.get("rn") or record.external_id), + "ogrn": organization.ogrn, + "inn": organization.inn, + "okpo": organization.okpo, + **{ + header: bool(payload.get(field_name)) + for header, field_name in SANCTIONS_EXPORT_FLAG_FIELDS.items() + }, + } row: dict[str, Any] = { "Наименование": organization.full_name or organization.short_name @@ -866,6 +913,17 @@ def _serialize_flat_value(value: Any) -> str | int | float | bool: return str(value) +def _serialize_spreadsheet_value(value: Any) -> str | int | float: + if isinstance(value, bool): + return int(value) + rendered_value = _serialize_flat_value(value) + if isinstance(rendered_value, str) and rendered_value.startswith( + ("=", "+", "-", "@") + ): + return f"'{rendered_value}" + return rendered_value + + def _serialize_json_value(value: Any) -> Any: if isinstance(value, dict | list | tuple): return json.dumps(value, ensure_ascii=False, sort_keys=True) diff --git a/src/organizations/test_companies.py b/src/organizations/test_companies.py index 86989d3..484de32 100644 --- a/src/organizations/test_companies.py +++ b/src/organizations/test_companies.py @@ -40,6 +40,9 @@ CANONICAL_ONLY_TEST_SOURCES = { ParserLoadLog.Source.GOSEDO_ADDRESS_DIRECTORY, ParserLoadLog.Source.MEDIA_NEWS, } +FILE_UPLOAD_ONLY_TEST_SOURCES = { + ParserLoadLog.Source.ROPK_SANCTIONS, +} TEST_BALANCE_LINE_NAMES = { "1110": "Нематериальные активы", @@ -171,7 +174,10 @@ class TestCompanyDatasetService: expected_external_ids = [] for source, descriptor in SOURCE_GROUP_DESCRIPTORS.items(): - if source not in ParserLoadLog.Source.values: + if ( + source not in ParserLoadLog.Source.values + or source in FILE_UPLOAD_ONLY_TEST_SOURCES + ): continue external_id = f"{TEST_RECORD_PREFIX}:{index:02d}:{source}" expected_external_ids.append(external_id) diff --git a/src/organizations/views.py b/src/organizations/views.py index 48be882..c07ab34 100644 --- a/src/organizations/views.py +++ b/src/organizations/views.py @@ -96,6 +96,7 @@ SOURCE_RECORD_ORDERING_FIELDS = ( "record_date", "record_type", "external_id", + "organization__name", "extension__organization__name", "extension__organization__full_name", "created_at", @@ -110,6 +111,13 @@ SOURCE_RECORD_ORDERING_FIELDS = ( "payload__full_name", "payload__medo_address", "payload__registration_number", + "payload__uk_hm_treasury", + "payload__european_union", + "payload__united_states", + "payload__switzerland", + "payload__united_states_sectoral", + "payload__uk_uksl", + "payload__ukraine", "payload__sentiment", "payload__news_source", ) @@ -651,6 +659,12 @@ class OrganizationSourceRecordViewSet(ReadOnlyModelViewSet): @staticmethod def _order_source_records(queryset, ordering: str): + if ordering.lstrip("-") == "organization__name": + ordering = ordering.replace( + "organization__name", + "extension__organization__name", + ) + return queryset.order_by(ordering, "external_id", "uid") if ordering.lstrip("-") == "record_date": expression = F("canonical_record_date") if ordering.startswith("-"): @@ -694,6 +708,11 @@ class OrganizationSourceRecordViewSet(ReadOnlyModelViewSet): ) ordering = self.request.query_params.get("ordering") + if not ordering and ( + self.request.query_params.get("source_group") == SourceGroup.SANCTIONS + or self.request.query_params.get("source") == "ropk_sanctions" + ): + ordering = "organization__name" if ordering and ordering not in SOURCE_RECORD_ORDERING_VALUES: errors.append( { diff --git a/tests/apps/organizations/test_source_record_export.py b/tests/apps/organizations/test_source_record_export.py index 53d6a25..525018e 100644 --- a/tests/apps/organizations/test_source_record_export.py +++ b/tests/apps/organizations/test_source_record_export.py @@ -22,6 +22,7 @@ from organizations.models import ( OrganizationSourceFinancialLine, OrganizationSourceRecord, PlannedInspectionExtension, + SanctionsExtension, SourceGroup, ) from organizations.source_record_export import ( @@ -92,9 +93,93 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): call_command("build_source_record_exports", stdout=command_output) generation = load_current_source_record_export_generation() - self.assertEqual(generation.artifacts_count, 31) + self.assertEqual(generation.artifacts_count, 34) self.assertEqual(generation.export_year, timezone.localdate().year) - self.assertIn('"artifacts_count": 31', command_output.getvalue()) + self.assertIn('"artifacts_count": 34', command_output.getvalue()) + + def test_sanctions_xlsx_export_uses_source_contract_columns_and_boolean_flags(self): + organization = Organization.objects.create( + name="АО Экспорт Санкции", + inn="0012345678", + ogrn="1027700132195", + okpo="00123456", + ) + extension = SanctionsExtension.objects.create( + organization=organization, + title="РОПК — Санкции", + ) + OrganizationSourceRecord.objects.create( + extension=extension, + source="ropk_sanctions", + record_type="organization_sanctions", + external_id="000001", + title=organization.name, + status="active", + payload={ + "rn": "000001", + "uk_hm_treasury": True, + "european_union": False, + "united_states": True, + "switzerland": False, + "united_states_sectoral": True, + "uk_uksl": False, + "ukraine": True, + }, + ) + build_source_record_export_artifacts() + self.client.force_authenticate(UserFactory.create_superuser()) + + response = self.client.post( + self.url, + {"sources": [SourceGroup.SANCTIONS.value], "format": "xlsx"}, + format="json", + ) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + with zipfile.ZipFile(BytesIO(self._response_body(response))) as archive: + self.assertEqual(archive.namelist(), ["ropk-sanctions.xlsx"]) + workbook = load_workbook( + BytesIO(archive.read("ropk-sanctions.xlsx")), + read_only=True, + data_only=True, + ) + rows = list(workbook["data"].iter_rows(values_only=True)) + workbook.close() + + self.assertEqual( + rows[0], + ( + "Наименование", + "rn", + "ogrn", + "inn", + "okpo", + "Санкции - Великобритания HM Treasury", + "Санкции - Евросоюз", + "Санкции - США", + "Санкции - Швейцария", + "Санкции секторальные - США", + "Санкции - Великобритания UKSL", + "Санкции - Украина", + ), + ) + self.assertEqual( + rows[1], + ( + "АО Экспорт Санкции", + "000001", + "1027700132195", + "0012345678", + "00123456", + 1, + 0, + 1, + 0, + 1, + 0, + 1, + ), + ) def test_generation_contains_only_records_from_its_calendar_year(self): export_year = 2026 @@ -313,7 +398,7 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): self.assertEqual( response["X-Source-Export-Generated-At"], generation.generated_at ) - self.assertEqual(generation.artifacts_count, 31) + self.assertEqual(generation.artifacts_count, 34) self.assertIn( 'filename="planned-inspections__financial-indicators_', response["Content-Disposition"], @@ -508,8 +593,8 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): key=lambda item: item.file_name, ) - self.assertEqual(generation.artifacts_count, 31) - self.assertEqual(generation.files_count, 32) + self.assertEqual(generation.artifacts_count, 34) + self.assertEqual(generation.files_count, 35) self.assertEqual( [item.file_name for item in artifacts], [ @@ -707,7 +792,7 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): first_generation = build_source_record_export_artifacts() current_generation = load_current_source_record_export_generation() - self.assertEqual(first_generation.artifacts_count, 31) + self.assertEqual(first_generation.artifacts_count, 34) self.assertEqual( current_generation.generation_id, first_generation.generation_id ) diff --git a/tests/apps/organizations/test_tasks.py b/tests/apps/organizations/test_tasks.py index 23091c7..133a352 100644 --- a/tests/apps/organizations/test_tasks.py +++ b/tests/apps/organizations/test_tasks.py @@ -145,7 +145,7 @@ class SourceRecordExportArtifactsTaskTest(TestCase): result = refresh_source_record_export_artifacts() self.assertEqual(result["status"], "success") - self.assertEqual(result["artifacts_count"], 31) + self.assertEqual(result["artifacts_count"], 34) self.assertEqual(result["export_year"], timezone.localdate().year) self.assertIsNone(cache.get(settings.SOURCE_RECORD_EXPORT_LOCK_KEY)) diff --git a/tests/apps/organizations/test_test_companies_commands.py b/tests/apps/organizations/test_test_companies_commands.py index 369e367..60a61ef 100644 --- a/tests/apps/organizations/test_test_companies_commands.py +++ b/tests/apps/organizations/test_test_companies_commands.py @@ -64,7 +64,11 @@ class TestCompaniesCommandsTest(TestCase): ) expected_groups = {choice.value for choice in SourceGroup} - expected_sources = {choice.value for choice in ParserLoadLog.Source} + expected_sources = { + choice.value + for choice in ParserLoadLog.Source + if choice is not ParserLoadLog.Source.ROPK_SANCTIONS + } for company in companies: self.assertEqual( set(company.source_extensions.values_list("source_group", flat=True)), diff --git a/tests/apps/parsers/test_ropk_sanctions.py b/tests/apps/parsers/test_ropk_sanctions.py new file mode 100644 index 0000000..0d48b6d --- /dev/null +++ b/tests/apps/parsers/test_ropk_sanctions.py @@ -0,0 +1,500 @@ +from __future__ import annotations + +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.models import ParserLoadLog, ParserSourceArtifact, ParserStagedRecord +from apps.parsers.ropk_sanctions import ( + ROPK_SANCTIONS_MAX_BYTES, + ROPK_SANCTIONS_MIME, + ROPK_SANCTIONS_RECORD_TYPE, + ROPK_SANCTIONS_SOURCE, + RopkSanctionsValidationError, + import_ropk_sanctions, +) +from apps.parsers.tasks import parse_ropk_sanctions +from django.core.files.base import ContentFile +from django.core.files.storage import default_storage +from django.core.files.uploadedfile import SimpleUploadedFile +from django.test import TestCase +from django.urls import reverse +from django.utils import timezone +from openpyxl import Workbook +from organizations.models import ( + Organization, + OrganizationSourceRecord, + SanctionsExtension, +) +from rest_framework import status +from rest_framework.test import APITestCase + +from tests.apps.user.factories import UserFactory + +HEADERS = [ + "rn", + "ogrn", + "inn", + "okpo", + "Санкции - Великобритания HM Treasury", + "Санкции - Евросоюз", + "Санкции - США", + "Санкции - Швейцария", + "Санкции секторальные - США", + "Санкции - Великобритания UKSL", + "Санкции - Украина", +] + + +def _workbook(rows: list[list[object]], *, headers: list[str] | None = None) -> BytesIO: + workbook = Workbook() + sheet = workbook.active + sheet.append(headers or HEADERS) + for row in rows: + sheet.append(row) + output = BytesIO() + workbook.save(output) + workbook.close() + output.seek(0) + return output + + +def _row( + rn: str, + *, + ogrn: str = "1027700132195", + inn: str = "0012345678", + okpo: str = "00123456", + flags: tuple[object, ...] = (1, 0, 1, 0, 1, 0, 1), +) -> list[object]: + return [rn, ogrn, inn, okpo, *flags] + + +class RopkSanctionsImportTest(TestCase): + def setUp(self): + self.organization = Organization.objects.create( + name="АО Санкции", + inn="0012345678", + ogrn="1027700132195", + okpo="00123456", + opk_registry_membership=True, + directory_imported_at=timezone.now(), + ) + + def test_import_enriches_by_identifiers_and_publishes_typed_payload(self): + artifact, result = import_ropk_sanctions( + handle=_workbook([_row("000001")]), + original_name="sanctions.xlsx", + load_batch=1, + uploaded_by_id=None, + ) + + self.assertFalse(result.skipped) + self.assertEqual(result.parsed, 1) + self.assertEqual(result.published, 1) + self.assertEqual(result.quarantined, 0) + self.assertEqual(artifact.status, ParserSourceArtifact.Status.PUBLISHED) + record = OrganizationSourceRecord.objects.get() + self.assertEqual(record.source, ROPK_SANCTIONS_SOURCE) + self.assertEqual(record.record_type, ROPK_SANCTIONS_RECORD_TYPE) + self.assertEqual(record.external_id, "000001") + self.assertEqual(record.title, self.organization.name) + self.assertEqual(record.record_date, "") + self.assertEqual( + record.payload, + { + "rn": "000001", + "uk_hm_treasury": True, + "european_union": False, + "united_states": True, + "switzerland": False, + "united_states_sectoral": True, + "uk_uksl": False, + "ukraine": True, + }, + ) + self.assertIsInstance(record.payload["uk_hm_treasury"], bool) + self.assertEqual(record.extension.organization_id, self.organization.uid) + + def test_import_enriches_missing_inn_and_ogrn_by_unique_okpo(self): + _, result = import_ropk_sanctions( + handle=_workbook([_row("2", ogrn="", inn="")]), + original_name="sanctions.xlsx", + load_batch=2, + uploaded_by_id=None, + ) + + self.assertEqual(result.published, 1) + record = OrganizationSourceRecord.objects.get() + self.assertEqual(record.extension.organization.inn, "0012345678") + self.assertEqual(record.extension.organization.ogrn, "1027700132195") + + def test_import_preserves_full_fourteen_digit_okpo(self): + self.organization.okpo = "00123456789012" + self.organization.save(update_fields=["okpo"]) + + _, result = import_ropk_sanctions( + handle=_workbook([_row("000014", ogrn="", inn="", okpo="00123456789012")]), + original_name="sanctions.xlsx", + load_batch=21, + uploaded_by_id=None, + ) + + self.assertEqual(result.published, 1) + record = OrganizationSourceRecord.objects.get() + self.assertEqual(record.extension.organization.okpo, "00123456789012") + self.assertEqual(record.external_id, "000014") + + def test_invalid_rows_are_quarantined_without_partial_records(self): + _, result = import_ropk_sanctions( + handle=_workbook( + [ + _row("duplicate"), + _row("duplicate"), + _row("bad-flag", flags=(2, 0, 0, 0, 0, 0, 0)), + _row("mismatch", okpo="99999999"), + ] + ), + original_name="sanctions.xlsx", + load_batch=3, + uploaded_by_id=None, + ) + + self.assertEqual(result.published, 0) + self.assertEqual(result.quarantined, 4) + self.assertEqual( + result.reasons, + {"duplicate_rn": 2, "invalid_flag": 1, "identifier_mismatch": 1}, + ) + self.assertFalse(OrganizationSourceRecord.objects.exists()) + self.assertEqual( + ParserStagedRecord.objects.filter( + disposition=ParserStagedRecord.Disposition.QUARANTINED + ).count(), + 4, + ) + + def test_new_snapshot_atomically_replaces_previous_records(self): + import_ropk_sanctions( + handle=_workbook([_row("old")]), + original_name="old.xlsx", + load_batch=4, + uploaded_by_id=None, + ) + import_ropk_sanctions( + handle=_workbook([_row("new")]), + original_name="new.xlsx", + load_batch=5, + uploaded_by_id=None, + ) + + self.assertEqual( + list( + OrganizationSourceRecord.objects.filter( + source=ROPK_SANCTIONS_SOURCE + ).values_list("external_id", flat=True) + ), + ["new"], + ) + + def test_duplicate_checksum_is_skipped_without_changing_snapshot(self): + raw = _workbook([_row("same")]).getvalue() + _, first = import_ropk_sanctions( + handle=BytesIO(raw), + original_name="first.xlsx", + load_batch=6, + uploaded_by_id=None, + ) + artifact, second = import_ropk_sanctions( + handle=BytesIO(raw), + original_name="second.xlsx", + load_batch=7, + uploaded_by_id=None, + ) + + self.assertFalse(first.skipped) + self.assertTrue(second.skipped) + self.assertEqual(artifact.status, ParserSourceArtifact.Status.SKIPPED) + self.assertEqual( + OrganizationSourceRecord.objects.filter( + source=ROPK_SANCTIONS_SOURCE + ).count(), + 1, + ) + + def test_invalid_header_order_rejects_batch_and_preserves_snapshot(self): + import_ropk_sanctions( + handle=_workbook([_row("old")]), + original_name="old.xlsx", + load_batch=8, + uploaded_by_id=None, + ) + reordered = [HEADERS[1], HEADERS[0], *HEADERS[2:]] + + with self.assertRaisesMessage(RopkSanctionsValidationError, "invalid_headers"): + import_ropk_sanctions( + handle=_workbook([_row("new")], headers=reordered), + original_name="invalid.xlsx", + load_batch=9, + uploaded_by_id=None, + ) + + self.assertEqual( + list( + OrganizationSourceRecord.objects.filter( + source=ROPK_SANCTIONS_SOURCE + ).values_list("external_id", flat=True) + ), + ["old"], + ) + self.assertEqual( + ParserSourceArtifact.objects.filter( + source=ROPK_SANCTIONS_SOURCE, + status=ParserSourceArtifact.Status.REJECTED, + ).count(), + 1, + ) + + def test_archive_expansion_limit_rejects_file_but_retains_raw_artifact(self): + with patch( + "apps.parsers.ropk_sanctions.ROPK_SANCTIONS_MAX_UNCOMPRESSED_BYTES", + 1, + ), self.assertRaisesMessage( + RopkSanctionsValidationError, + "xlsx_uncompressed_size_exceeded", + ): + import_ropk_sanctions( + handle=_workbook([_row("zip-limit")]), + original_name="sanctions.xlsx", + load_batch=22, + uploaded_by_id=None, + ) + + artifact = ParserSourceArtifact.objects.get(source=ROPK_SANCTIONS_SOURCE) + self.assertEqual(artifact.status, ParserSourceArtifact.Status.REJECTED) + self.assertTrue(artifact.file.name) + self.assertTrue(artifact.file.storage.exists(artifact.file.name)) + + +class RopkSanctionsApiTest(APITestCase): + def setUp(self): + self.user = UserFactory.create_user() + self.admin = UserFactory.create_user(is_staff=True) + self.upload_url = reverse( + "api_v1:parsers:upload-parser-data", + args=[ROPK_SANCTIONS_SOURCE], + ) + + def _upload( + self, name: str, content: bytes, content_type: str = ROPK_SANCTIONS_MIME + ): + return self.client.post( + self.upload_url, + {"file": SimpleUploadedFile(name, content, content_type=content_type)}, + format="multipart", + ) + + def test_upload_is_admin_only_and_queues_new_source(self): + self.client.force_authenticate(self.user) + self.assertEqual( + self._upload("sanctions.xlsx", b"xlsx").status_code, + status.HTTP_403_FORBIDDEN, + ) + + self.client.force_authenticate(self.admin) + with patch( + "apps.parsers.views._save_uploaded_parser_file", + return_value="parser_uploads/sanctions.xlsx", + ), patch( + "apps.parsers.tasks.parse_ropk_sanctions.apply_async", + side_effect=lambda **kwargs: SimpleNamespace(id=kwargs["task_id"]), + ): + response = self._upload("sanctions.xlsx", b"xlsx") + + self.assertEqual(response.status_code, status.HTTP_202_ACCEPTED) + self.assertEqual(response.data["data"]["source"], ROPK_SANCTIONS_SOURCE) + self.assertEqual(response.data["data"]["status"], JobStatus.PENDING) + self.assertEqual( + response.data["data"]["task_ids"], + [response.data["data"]["task_id"]], + ) + + def test_upload_rejects_extension_mime_and_size_before_queueing(self): + self.client.force_authenticate(self.admin) + with patch("apps.parsers.views._save_uploaded_parser_file") as save_file, patch( + "apps.parsers.tasks.parse_ropk_sanctions.apply_async" + ) as apply_async: + extension = self._upload("sanctions.csv", b"csv", "text/csv") + mime = self._upload("sanctions.xlsx", b"xlsx", "application/zip") + with patch("apps.parsers.views.ROPK_SANCTIONS_MAX_BYTES", 1): + oversized = self._upload("sanctions.xlsx", b"xx") + + self.assertEqual(extension.status_code, status.HTTP_415_UNSUPPORTED_MEDIA_TYPE) + self.assertEqual(extension.data["errors"][0]["code"], "unsupported_media_type") + self.assertEqual(mime.status_code, status.HTTP_415_UNSUPPORTED_MEDIA_TYPE) + self.assertEqual(mime.data["errors"][0]["code"], "unsupported_media_type") + self.assertEqual( + oversized.status_code, status.HTTP_413_REQUEST_ENTITY_TOO_LARGE + ) + self.assertEqual(oversized.data["errors"][0]["code"], "file_too_large") + save_file.assert_not_called() + apply_async.assert_not_called() + self.assertEqual(ROPK_SANCTIONS_MAX_BYTES, 25 * 1024 * 1024) + + def test_upload_requires_file_before_creating_job(self): + self.client.force_authenticate(self.admin) + + response = self.client.post(self.upload_url, {}, format="multipart") + + self.assertEqual(response.status_code, status.HTTP_400_BAD_REQUEST) + self.assertEqual(response.data["errors"][0]["code"], "file_required") + self.assertFalse(BackgroundJob.objects.exists()) + + def test_upload_rejects_concurrent_job_before_saving_file(self): + BackgroundJob.objects.create( + task_id="ropk-running", + task_name="apps.parsers.tasks.parse_ropk_sanctions", + status=JobStatus.STARTED, + ) + self.client.force_authenticate(self.admin) + with patch("apps.parsers.views._save_uploaded_parser_file") as save_file: + response = self._upload("sanctions.xlsx", b"xlsx") + + self.assertEqual(response.status_code, status.HTTP_409_CONFLICT) + self.assertEqual(response.data["errors"][0]["code"], "upload_already_running") + save_file.assert_not_called() + + def test_catalog_records_and_file_only_contract_are_available(self): + organization = Organization.objects.create( + name="АО API Санкции", + inn="7707083801", + ogrn="1027700132101", + okpo="12345671", + opk_registry_membership=True, + directory_imported_at=timezone.now(), + ) + extension = SanctionsExtension.objects.create( + organization=organization, + title="РОПК — Санкции", + ) + OrganizationSourceRecord.objects.create( + extension=extension, + source=ROPK_SANCTIONS_SOURCE, + record_type=ROPK_SANCTIONS_RECORD_TYPE, + external_id="10", + title=organization.name, + status="active", + payload={ + "rn": "10", + "uk_hm_treasury": True, + "european_union": False, + "united_states": False, + "switzerland": False, + "united_states_sectoral": False, + "uk_uksl": False, + "ukraine": False, + }, + ) + self.client.force_authenticate(self.user) + + card = self.client.get( + reverse( + "api_v1:sources:source-cards-detail", + kwargs={"slug": "ropk-sanctions"}, + ) + ) + records = self.client.get( + reverse("api_v2:organizations:organization-source-records-list"), + { + "source_group": "sanctions", + "source": ROPK_SANCTIONS_SOURCE, + "record_type": ROPK_SANCTIONS_RECORD_TYPE, + "ordering": "payload__uk_hm_treasury", + }, + ) + self.client.force_authenticate(self.admin) + run = self.client.post( + reverse( + "api_v1:parsers:run-parser", + args=[ROPK_SANCTIONS_SOURCE], + ), + {}, + format="json", + ) + + self.assertEqual(card.status_code, status.HTTP_200_OK) + self.assertEqual(card.data["data"]["slug"], "ropk-sanctions") + self.assertFalse(card.data["data"]["supports_refresh"]) + self.assertEqual(card.data["data"]["upload_url"], self.upload_url) + self.assertEqual(records.status_code, status.HTTP_200_OK) + self.assertEqual(records.data["meta"]["pagination"]["total_count"], 1) + self.assertEqual(records.data["data"][0]["payload"]["rn"], "10") + self.assertIs(records.data["data"][0]["payload"]["uk_hm_treasury"], True) + self.assertEqual(run.status_code, status.HTTP_405_METHOD_NOT_ALLOWED) + self.assertEqual(run.data["errors"][0]["code"], "parser_run_not_supported") + + refresh = self.client.post( + reverse( + "api_v1:sources:source-cards-refresh", + kwargs={"slug": "ropk-sanctions"}, + ), + {}, + format="json", + ) + self.assertEqual(refresh.status_code, status.HTTP_405_METHOD_NOT_ALLOWED) + self.assertEqual(refresh.data["errors"][0]["code"], "refresh_not_supported") + + def test_openapi_exposes_source_specific_upload_and_boolean_payload(self): + response = self.client.get( + reverse("schema-swagger-ui"), + {"format": "openapi"}, + ) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + schema = response.content.decode("utf-8") + self.assertIn("/api/v1/parsers/upload/ropk_sanctions/", schema) + self.assertIn("XLSX-файл РОПК — Санкции", schema) + self.assertIn("uk_hm_treasury", schema) + self.assertIn("european_union", schema) + + +class RopkSanctionsTaskTest(TestCase): + def test_task_completes_job_and_reports_contract_counts(self): + Organization.objects.create( + name="АО Задача Санкции", + inn="0012345678", + ogrn="1027700132195", + okpo="00123456", + opk_registry_membership=True, + directory_imported_at=timezone.now(), + ) + workbook = _workbook([_row("task-row")]) + + with TemporaryDirectory() as media_root, self.settings(MEDIA_ROOT=media_root): + file_path = default_storage.save( + "parser_uploads/sanctions.xlsx", + ContentFile(workbook.getvalue()), + ) + task_result = parse_ropk_sanctions.apply( + kwargs={ + "file_path": file_path, + "original_name": "sanctions.xlsx", + }, + task_id="ropk-sanctions-task-test", + ) + self.assertFalse(default_storage.exists(file_path)) + + self.assertTrue(task_result.successful()) + self.assertEqual(task_result.result["raw_records_count"], 1) + self.assertEqual(task_result.result["published_records_count"], 1) + self.assertEqual(task_result.result["quarantine_records_count"], 0) + self.assertIn("load_id", task_result.result) + job = BackgroundJob.objects.get(task_id="ropk-sanctions-task-test") + self.assertEqual(job.status, JobStatus.SUCCESS) + self.assertEqual(job.result["published_records_count"], 1) + self.assertEqual( + ParserLoadLog.objects.get(source=ROPK_SANCTIONS_SOURCE).status, + ParserLoadLog.Status.SUCCESS, + ) diff --git a/tests/apps/parsers/test_source_cards_service.py b/tests/apps/parsers/test_source_cards_service.py index 12d9c42..12f5660 100644 --- a/tests/apps/parsers/test_source_cards_service.py +++ b/tests/apps/parsers/test_source_cards_service.py @@ -86,6 +86,7 @@ class SourceCardServiceUnitTest(SimpleTestCase): "planned-inspections", "bankruptcy-procedures", "defense-unreliable-suppliers", + "ropk-sanctions", "arbitration-cases", "information-security-registries", "labor-vacancies", @@ -102,6 +103,7 @@ class SourceCardServiceUnitTest(SimpleTestCase): "Плановые проверки Генпрокуратуры России", "Сведения о процедурах банкротства", "Недобросовестные поставщики ГОЗ", + "РОПК — Санкции", "Арбитражные дела", "Реестры по информационной безопасности", "Вакансии Работа России",