From b19951d0c9f8533c361e98e136c728aef3397736 Mon Sep 17 00:00:00 2001 From: Aleksandr Meshchriakov Date: Mon, 3 Aug 2026 19:27:07 +0200 Subject: [PATCH] feat: align source data and nightly exports --- .env.prod.example | 4 + README.md | 1 - docker-compose.dev.yml | 1 + docker-compose.prod.yml | 1 + docker/Dockerfile | 3 + docs/parser-external-access-note-ru.md | 8 +- docs/source-record-export-matrix-ru.md | 95 ++ src/apps/parsers/checko_collection.py | 4 +- src/apps/parsers/clients/checko/client.py | 10 +- ...9_neutralize_external_collection_labels.py | 18 + src/apps/parsers/models.py | 6 +- src/apps/parsers/source_registry.py | 6 +- src/apps/parsers/tasks.py | 90 +- .../commands/build_source_record_exports.py | 33 + ...0008_seed_nightly_source_record_exports.py | 61 ++ src/organizations/serializers.py | 43 + src/organizations/source_record_export.py | 878 +++++++++++++++--- src/organizations/tasks.py | 36 + src/organizations/views.py | 172 +++- src/settings/base.py | 22 +- .../test_api_v2_source_extensions.py | 53 ++ .../test_source_record_export.py | 324 ++++++- tests/apps/organizations/test_tasks.py | 64 +- tests/apps/parsers/test_source_registry.py | 13 + 24 files changed, 1736 insertions(+), 210 deletions(-) create mode 100644 docs/source-record-export-matrix-ru.md create mode 100644 src/apps/parsers/migrations/0029_neutralize_external_collection_labels.py create mode 100644 src/organizations/management/commands/build_source_record_exports.py create mode 100644 src/organizations/migrations/0008_seed_nightly_source_record_exports.py diff --git a/.env.prod.example b/.env.prod.example index 0d1dedf..0243de8 100644 --- a/.env.prod.example +++ b/.env.prod.example @@ -40,6 +40,10 @@ COLLECTSTATIC_ON_MIGRATE=0 BACKUP_ENCRYPTION_KEY=a2tra2tra2tra2tra2tra2tra2tra2tra2tra2s BACKUP_KEY_ID=default BACKUP_EXPORT_DIRECTORY=/app/media/backups +SOURCE_RECORD_EXPORT_DIRECTORY=/app/media/source-record-exports +SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP=2 +SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE=100000 +SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS=300 STATE_CORP_EXCHANGE_URL= STATE_CORP_EXCHANGE_TOKEN= diff --git a/README.md b/README.md index 5c649b6..0d89eb6 100644 --- a/README.md +++ b/README.md @@ -56,7 +56,6 @@ Backend-сервис на Django/DRF для сбора, хранения и вы - `DJANGO_SETTINGS_MODULE`: `settings.dev` (локально/dev) или `settings.production` (prod). - `POSTGRES_*`: доступ к PostgreSQL. - `REDIS_CACHE_URL`, `CELERY_BROKER_URL`, `CELERY_RESULT_BACKEND`: Redis/Celery. -- `CHECKO_API_KEY`: ключ API Checko. - `ZAKUPKI_TOKEN`: токен SOAP API ЕИС закупок. - `COLLECTSTATIC_ON_MIGRATE`: `1`/`0`. - `STARTUP_CHECKS_ENABLED`: fail-fast проверки DB/Redis перед стартом runtime-процессов. diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 82aa951..82d57ee 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -104,6 +104,7 @@ services: memswap_limit: 3g volumes: - ./input:/app/input + - ./media:/app/media command: ["/app/docker/scripts/start-celery-worker.sh"] celery_beat: diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index fdad9d5..8128752 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -64,6 +64,7 @@ services: volumes: - ./logs:/app/logs - ./input:/app/input + - ./media:/app/media command: ["/app/docker/scripts/start-celery-worker.sh"] celery_beat: diff --git a/docker/Dockerfile b/docker/Dockerfile index 48db094..3494973 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -103,6 +103,9 @@ ENV PATH="/app/.venv/bin:${PATH}" \ BACKUP_ENCRYPTION_KEY= \ BACKUP_KEY_ID=default \ BACKUP_EXPORT_DIRECTORY=/app/media/backups \ + SOURCE_RECORD_EXPORT_DIRECTORY=/app/media/source-record-exports \ + SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE=100000 \ + SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS=300 \ STATE_CORP_EXCHANGE_URL= \ STATE_CORP_EXCHANGE_TOKEN= \ STATE_CORP_EXCHANGE_KEY_ID=state-corp-shared-token \ diff --git a/docs/parser-external-access-note-ru.md b/docs/parser-external-access-note-ru.md index c67d634..68d0d7c 100644 --- a/docs/parser-external-access-note-ru.md +++ b/docs/parser-external-access-note-ru.md @@ -21,13 +21,13 @@ | ЕИС закупки: HTTP fallback | `https://zakupki.gov.ru/opendata/download/notifications/{region}/{year}/...` | ZIP-архивы с XML-файлами закупок, если SOAP-токен не используется или передана прямая ссылка. | | ЕИС/FAS generic-источники | `https://zakupki.gov.ru/epz/order/extendedsearch/results.html`, `https://zakupki.gov.ru/epz/orderclause/search/results.html`, `https://zakupki.gov.ru/epz/contract/search/results.html`, `https://zakupki.gov.ru/epz/dishonestsupplier/search/results.html`, `https://fas.gov.ru/pages/activity/reestr-uridicheskih-lic` | HTML-страницы официальных реестров. Парсер извлекает карточки/таблицы: закупки 44-ФЗ, закупки 223-ФЗ, контракты, недобросовестные поставщики, сведения ФАС по ГОЗ. | | ФНС: бухгалтерская отчетность | Автоматического HTTP-скачивания с ФНС в текущем коде не найдено. В каталоге источников указан справочный URL `https://bo.nalog.gov.ru/advanced-search/organizations/search?...` | Обрабатываются локально загруженные или положенные в папку `input/fns` файлы `fin_{id}_{ogrn}.xlsx`, а также ZIP-архивы с такими файлами. Из Excel берутся строки форм N 1, 2, 3, 4, 6 бухгалтерской отчетности. | -| КАД Арбитр через Checko | Официальный источник в каталоге: `https://kad.arbitr.ru/`; фактический lookup в коде: `https://api.checko.ru/v2/legal-cases` | JSON-ответы по арбитражным делам для активных организаций из внутренних реестров. В запрос передаются ИНН/ОГРН. В payload сохраняются номер дела, суд, тип, статус, даты, суммы, стороны и ссылка на карточку. | -| Федресурс/ЕФРСБ | `https://bankrot.fedresurs.ru/`; fallback: `https://api.checko.ru/v2/company` | Официальный источник обрабатывается как HTML/структурированная выгрузка. При недоступности портала используется Checko: по ИНН/ОГРН организации берутся сведения о банкротных сообщениях из JSON. | +| КАД Арбитр через внешний сервис данных | Официальный источник в каталоге: `https://kad.arbitr.ru/`; фактический lookup выполняется через служебный API внешнего сервиса | JSON-ответы по арбитражным делам для активных организаций из внутренних реестров. В запрос передаются ИНН/ОГРН. В payload сохраняются номер дела, суд, тип, статус, даты, суммы, стороны и ссылка на карточку. | +| Федресурс/ЕФРСБ | `https://bankrot.fedresurs.ru/`; fallback выполняется через служебный API внешнего сервиса | Официальный источник обрабатывается как HTML/структурированная выгрузка. При недоступности портала по ИНН/ОГРН организации запрашиваются сведения о банкротных сообщениях из JSON. | | ФСТЭК | `https://reestr.fstec.ru/reg3` и найденные на странице ссылки вида `module=rfiles` или `/uploads/reg...` | HTML-страница реестра, затем CSV/файловая выгрузка, если ссылка найдена. Для этого источника в коде отключена SSL-верификация. | | Вакансии: Работа России | Клиент использует `http://opendata.trudvsem.ru/api/v1/vacancies`, `http://opendata.trudvsem.ru/api/v1/vacancies/company/inn/{inn}`; в каталоге источников указан `https://opendata.trudvsem.ru/api/v1/vacancies` | JSON-список вакансий, включая работодателя, ИНН/ОГРН при наличии, название вакансии, дату, зарплату, статус, ссылку. | | Вакансии: HeadHunter | `https://api.hh.ru/vacancies` | JSON-список вакансий. Поиск выполняется по региону и/или тексту, для организаций без поиска по ИНН используется нормализованное название. | | Вакансии: SuperJob | `https://api.superjob.ru/2.0/vacancies/` | JSON-список вакансий. Используется только если задан `SUPERJOB_APP_ID`; ключ передается в заголовке `X-Api-App-Id`. | -| Checko: контракты и проверки по организациям | `https://api.checko.ru/v2/contracts`, `https://api.checko.ru/v2/inspections` | JSON-данные по контрактам и проверкам для активных организаций из внутренних реестров. В запрос передаются ИНН/ОГРН, API-ключ передается параметром `key`. | +| Внешний сервис данных: контракты и проверки по организациям | Служебные API контрактов и проверок | JSON-данные по контрактам и проверкам для активных организаций из внутренних реестров. В запрос передаются ИНН/ОГРН, API-ключ передается параметром `key`. | | Proxy-Tools | `https://proxy-tools.com/api/v1/proxies` | Служебная загрузка списка RU-прокси для парсеров. Используется только при заданном `PROXY_TOOLS_API_KEY`; запрос идет с Bearer-токеном. | ## Форматы загружаемых данных @@ -51,7 +51,7 @@ - коды регионов; - ИНН/ОГРН организаций из внутренних активных реестров; - поисковая строка по названию организации для вакансий; -- служебные ключи API из окружения: `ZAKUPKI_TOKEN`, `CHECKO_API_KEY`, `SUPERJOB_APP_ID`, `PROXY_TOOLS_API_KEY`. +- служебные ключи API из окружения для ЕИС, внешнего сервиса данных, SuperJob и Proxy-Tools. Ключи в коде не захардкожены, берутся из переменных окружения. diff --git a/docs/source-record-export-matrix-ru.md b/docs/source-record-export-matrix-ru.md new file mode 100644 index 0000000..33492d7 --- /dev/null +++ b/docs/source-record-export-matrix-ru.md @@ -0,0 +1,95 @@ +# Матрица файловых выгрузок источников + +## Пользовательский контракт + +Frontend отправляет администраторский +`POST /api/v2/organization-source-records/export-ticket/` с массивом `sources` +и выбранным `format`. В ответ он получает короткоживущий одноразовый ticket и +передаёт его обычной HTML-формой в +`POST /api/v2/organization-source-records/export-download/`. Поэтому браузер +сохраняет потоковый ZIP напрямую на диск, не удерживая весь архив как Blob в +JavaScript. Ticket передаётся в теле формы, не попадает в URL и после первого +запроса становится недействительным. + +Совместимый администраторский +`POST /api/v2/organization-source-records/export/` по-прежнему сразу возвращает +тот же ZIP API-клиентам. Крупный XLSX может состоять из нескольких файлов +`*-part-001.xlsx`, `*-part-002.xlsx` и далее. + +При скачивании endpoint не читает таблицы записей источников и не строит CSV, +XLSX или JSON заново. Он упаковывает файлы последнего полностью опубликованного +ночного поколения и сразу потоково отправляет ZIP без временной копии всего +архива. Поэтому `Content-Length` у ответа отсутствует. Если ни одного поколения +ещё нет, API отвечает `503` с кодом `source_export_not_ready`. + +## Матрица + +| Группа API | Файл | CSV | XLSX | JSON | +|---|---|:---:|:---:|:---:| +| `financial_indicators` | `financial-indicators` | — | — | да | +| `government_procurements` | `public-procurements` | да | да | да | +| `industrial_production` | `manufacturers-and-products` | да | да | да | +| `planned_inspections` | `planned-inspections` | да | да | да | +| `bankruptcy` | `bankruptcy-procedures` | да | да | да | +| `defense_suppliers` | `defense-unreliable-suppliers` | да | да | да | +| `arbitration` | `arbitration-cases` | да | да | да | +| `security_registries` | `information-security-registries` | да | да | да | +| `vacancies` | `labor-vacancies` | да | да | да | + +Итого формируется 25 логических артефактов: один JSON для финансовых показателей +и по три формата для остальных восьми групп. Физических файлов может быть +больше из-за разбиения крупных XLSX. Если финансовые показатели выбраны вместе +с другим форматом, в ZIP для них всё равно включается JSON. + +## Ночная генерация + +Celery Beat запускает +`organizations.tasks.refresh_source_record_export_artifacts` ежедневно в +`05:30 Europe/Moscow`, после ежедневного обновления organization sources в +`04:30`. + +Генератор: + +1. читает каждую группу из БД один раз без глобальной сортировки миллионов строк; +2. пишет compact JSON-массив на диск и использует его как готовый JSON без второй копии; +3. потоково создаёт CSV и XLSX без накопления всех строк в памяти; +4. разбивает XLSX по умолчанию по 100 000 строк на отдельные файлы, ограничивая временный XML и не превышая лимит Excel; +5. записывает размеры файлов и номера частей в manifest; +6. атомарно переключает `current.json` только после готовности всей матрицы; +7. сохраняет текущее и предыдущее поколения по умолчанию. + +При ошибке незавершённое поколение удаляется, а download endpoint продолжает +отдавать предыдущую успешную версию. + +### Расчёт диска + +Атомарная публикация требует одновременно хранить уже опубликованные поколения +и одно новое поколение в staging. Минимальный запас под артефакты рассчитывается +как `(SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP + 1) * размер поколения`, плюс +рабочий запас файловой системы. На снимке dev от 2026-08-03 одно поколение +заняло 16,55 ГБ (54 физических файла), поэтому при значении `2` следует +выделить не менее 55 ГБ свободного места под каталог выгрузок. Временная копия +целого ZIP при скачивании не создаётся. + +## Хранение и первый запуск + +Каталог задаётся через `SOURCE_RECORD_EXPORT_DIRECTORY`, по умолчанию — +`media/source-record-exports`. Он должен быть общим read-write volume для web и +Celery worker. В Docker Compose используется `./media:/app/media`. + +После первого развёртывания готовое поколение можно создать сразу, не ожидая +ночного расписания: + +```bash +PYTHONPATH=src uv run python src/manage.py build_source_record_exports +``` + +Доступные настройки: + +| Настройка | Значение по умолчанию | Назначение | +|---|---:|---| +| `SOURCE_RECORD_EXPORT_DIRECTORY` | `media/source-record-exports` | Общий каталог артефактов | +| `SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP` | `2` | Число сохраняемых успешных поколений | +| `SOURCE_RECORD_EXPORT_LOCK_TTL_SECONDS` | `21600` | TTL распределённой блокировки Celery | +| `SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE` | `100000` | Максимум строк данных в одной XLSX-части | +| `SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS` | `300` | Срок действия одноразового browser-download ticket | diff --git a/src/apps/parsers/checko_collection.py b/src/apps/parsers/checko_collection.py index 033740f..eddd0b7 100644 --- a/src/apps/parsers/checko_collection.py +++ b/src/apps/parsers/checko_collection.py @@ -1,4 +1,4 @@ -"""Quota-safe monthly claims for organization lookups in Checko.""" +"""Quota-safe monthly claims for external organization lookups.""" from datetime import date, datetime @@ -24,7 +24,7 @@ def claim_monthly_collection( source: str, at: date | datetime | None = None, ) -> CheckoCollectionAttempt | None: - """Atomically claim this month's only allowed Checko collection attempt.""" + """Atomically claim this month's only allowed external collection attempt.""" attempt, created = CheckoCollectionAttempt.objects.get_or_create( organization_id=organization_id, source=source, diff --git a/src/apps/parsers/clients/checko/client.py b/src/apps/parsers/clients/checko/client.py index d51833b..5a8c45e 100644 --- a/src/apps/parsers/clients/checko/client.py +++ b/src/apps/parsers/clients/checko/client.py @@ -346,14 +346,18 @@ class CheckoClient: # Preserve the existing classification for non-quota HTTP # errors while still recovering Checko quota metadata. pass - logger.error("Checko HTTP request failed with status=%s", e.status_code) + logger.error( + "External provider HTTP request failed with status=%s", e.status_code + ) raise CheckoConnectionError( - "Checko API request failed", + "External provider API request failed", url=e.url, ) from e except Exception as e: logger.error("Connection error: %s", e) - raise CheckoConnectionError(f"Failed to connect to Checko API: {e}") from e + raise CheckoConnectionError( + f"Failed to connect to external provider API: {e}" + ) from e self._raise_api_error(data) diff --git a/src/apps/parsers/migrations/0029_neutralize_external_collection_labels.py b/src/apps/parsers/migrations/0029_neutralize_external_collection_labels.py new file mode 100644 index 0000000..39c69b3 --- /dev/null +++ b/src/apps/parsers/migrations/0029_neutralize_external_collection_labels.py @@ -0,0 +1,18 @@ +from django.db import migrations + + +class Migration(migrations.Migration): + dependencies = [ + ("parsers", "0028_registry_procurement_claims"), + ] + + operations = [ + migrations.AlterModelOptions( + name="checkocollectionattempt", + options={ + "ordering": ["-period_month", "source", "organization_id"], + "verbose_name": "попытка сбора внешних данных", + "verbose_name_plural": "попытки сбора внешних данных", + }, + ), + ] diff --git a/src/apps/parsers/models.py b/src/apps/parsers/models.py index b2763d1..161bacd 100644 --- a/src/apps/parsers/models.py +++ b/src/apps/parsers/models.py @@ -118,7 +118,7 @@ class ParserBatchSequence(TimestampMixin, models.Model): class CheckoCollectionAttempt(TimestampMixin, models.Model): - """Monthly Checko collection claim for one organization and source.""" + """Monthly external collection claim for one organization and source.""" class Source(models.TextChoices): ARBITRATION = "arbitration", _("Арбитражные дела") @@ -160,8 +160,8 @@ class CheckoCollectionAttempt(TimestampMixin, models.Model): class Meta: db_table = "parsers_checko_collection_attempt" - verbose_name = _("попытка сбора Checko") - verbose_name_plural = _("попытки сбора Checko") + verbose_name = _("попытка сбора внешних данных") + verbose_name_plural = _("попытки сбора внешних данных") ordering = ["-period_month", "source", "organization_id"] constraints = [ models.UniqueConstraint( diff --git a/src/apps/parsers/source_registry.py b/src/apps/parsers/source_registry.py index d9b14e3..880ea43 100644 --- a/src/apps/parsers/source_registry.py +++ b/src/apps/parsers/source_registry.py @@ -236,10 +236,10 @@ PARSER_SOURCES: dict[str, ParserSourceDescriptor] = { status="implemented", upstream_url="https://kad.arbitr.ru/", access_method="official_search_api", - parser_strategy="checko_legal_cases_by_inn_ogrn", + parser_strategy="external_legal_cases_by_inn_ogrn", source_notes=( "Поиск дел выполняется по ИНН/ОГРН активных организаций из реестров. " - "Checko отдаёт карточки со ссылками на КАД Арбитр." + "Внешний сервис данных отдаёт карточки со ссылками на КАД Арбитр." ), api_route="arbitration/cases", ), @@ -258,7 +258,7 @@ PARSER_SOURCES: dict[str, ParserSourceDescriptor] = { parser_strategy="fedresurs_bankruptcy_search", source_notes=( "Официальный ЕФРСБ; может отдавать anti-bot challenge worker'ам. " - "Если официальный портал недоступен, используется Checko API по " + "Если официальный портал недоступен, используется внешний сервис данных по " "организациям из реестров. " "Ручная загрузка разрешена только для выгрузок, переданных Сергеем." ), diff --git a/src/apps/parsers/tasks.py b/src/apps/parsers/tasks.py index 2ad171a..f129492 100644 --- a/src/apps/parsers/tasks.py +++ b/src/apps/parsers/tasks.py @@ -470,7 +470,7 @@ def _fetch_fedresurs_bankruptcy_records( file_path: str | None, proxies: list[str] | None, ) -> list[GenericParserItem]: - """Загрузить банкротства: официальный портал, затем fallback через Checko.""" + """Загрузить банкротства: официальный портал, затем внешний fallback.""" official_error: Exception | None = None try: official_records = _fetch_structured_records( @@ -482,14 +482,14 @@ def _fetch_fedresurs_bankruptcy_records( if official_records or file_url or file_path: return official_records logger.warning( - "Fedresurs official source returned no records, falling back to Checko" + "Fedresurs official source returned no records, using external fallback" ) except Exception as exc: if file_url or file_path: raise official_error = exc logger.warning( - "Fedresurs official source failed, falling back to Checko: %s", + "Fedresurs official source failed, using external fallback: %s", exc, ) records = _fetch_checko_bankruptcy_records(proxies=proxies) @@ -498,12 +498,12 @@ def _fetch_fedresurs_bankruptcy_records( if official_error is None: raise ParserSourceSkipped( "fedresurs official source returned no bankruptcy records; " - "Checko fallback returned no bankruptcy records" + "external fallback returned no bankruptcy records" ) if isinstance(official_error, HTTPClientError): raise ParserSourceSkipped( "fedresurs upstream is unavailable or blocked; " - "Checko fallback returned no bankruptcy records" + "external fallback returned no bankruptcy records" ) from official_error raise official_error @@ -580,7 +580,7 @@ def _enrich_fstec_record_identities( logger.info( "FSTEC identity enrichment completed: enriched=%d ambiguous=%d " - "local_candidates=%d checko_candidates=%d", + "local_candidates=%d external_candidates=%d", enriched_count, ambiguous_count, len(local_candidates), @@ -720,7 +720,7 @@ def _fstec_checko_identity_candidates( ) except CheckoError as exc: logger.info( - "Checko FSTEC identity lookup skipped for %s: %s", + "External FSTEC identity lookup skipped for %s: %s", applicant_name, exc, ) @@ -846,10 +846,10 @@ def _fetch_checko_bankruptcy_records( *, proxies: list[str] | None, ) -> list[GenericParserItem]: - """Получить ЕФРСБ-сообщения по организациям из наших реестров через Checko.""" + """Получить ЕФРСБ-сообщения по организациям через внешний сервис данных.""" api_key = getattr(settings, "CHECKO_API_KEY", "") if not api_key: - logger.warning("CHECKO_API_KEY is empty; Fedresurs fallback skipped") + logger.warning("External provider API key is empty; Fedresurs fallback skipped") return [] limit = _resolve_lookup_limit( @@ -861,7 +861,7 @@ def _fetch_checko_bankruptcy_records( default=FEDRESURS_CHECKO_FALLBACK_LIMIT, ) if limit <= 0: - logger.info("Fedresurs Checko fallback is disabled by limit=%s", limit) + logger.info("Fedresurs external fallback is disabled by limit=%s", limit) return [] targets = _active_registry_lookup_targets( limit=limit, @@ -893,7 +893,7 @@ def _fetch_checko_bankruptcy_records( except CheckoRateLimitError as exc: finish_collection(attempt, records_count=0, error=exc) logger.warning( - "Checko bankruptcy fallback stopped: quota/rate limit reached " + "External bankruptcy fallback stopped: quota/rate limit reached " "(status_code=%s)", exc.status_code, ) @@ -901,7 +901,7 @@ def _fetch_checko_bankruptcy_records( except CheckoError as exc: finish_collection(attempt, records_count=0, error=exc) logger.info( - "Checko bankruptcy lookup skipped for target=%s: %s", + "External bankruptcy lookup skipped for target=%s: %s", target.inn or target.ogrn, exc, ) @@ -920,7 +920,7 @@ def _fetch_checko_bankruptcy_records( attempt, records_count=len(records) - records_before, ) - logger.info("Fetched %d bankruptcy records through Checko fallback", len(records)) + logger.info("Fetched %d bankruptcy records through external fallback", len(records)) return records @@ -931,7 +931,7 @@ def _checko_bankruptcy_items( fallback_ogrn: str, fallback_name: str, ) -> list[GenericParserItem]: - """Преобразовать банкротные сообщения Checko в generic records.""" + """Преобразовать банкротные сообщения внешнего сервиса в generic records.""" inn = str(getattr(company, "inn", "") or fallback_inn) ogrn = str(getattr(company, "ogrn", "") or fallback_ogrn) name = getattr(company, "short_name", None) or fallback_name @@ -1084,14 +1084,16 @@ def _fetch_checko_arbitration_records( limit: int | None, proxies: list[str] | None, ) -> list[GenericParserItem]: - """Получить арбитражные дела по ИНН/ОГРН через Checko legal-cases API.""" + """Получить арбитражные дела по ИНН/ОГРН через внешний API.""" api_key = getattr(settings, "CHECKO_API_KEY", "") if not api_key: - raise ParserSourceSkipped("CHECKO_API_KEY is empty; arbitration parser skipped") + raise ParserSourceSkipped( + "External provider API key is empty; arbitration parser skipped" + ) resolved_limit = _resolve_arbitration_limit(limit) if resolved_limit <= 0: - logger.info("Arbitration Checko parser is disabled by limit=%s", limit) + logger.info("Arbitration external parser is disabled by limit=%s", limit) return [] subjects = _arbitration_subjects(resolved_limit) @@ -1128,7 +1130,7 @@ def _fetch_checko_arbitration_records( failed_lookups += 1 finish_collection(attempt, records_count=0, error=exc) logger.info( - "Checko arbitration lookup skipped for subject=%s: %s", + "External arbitration lookup skipped for subject=%s: %s", _arbitration_subject_key(subject), exc, ) @@ -1139,10 +1141,12 @@ def _fetch_checko_arbitration_records( ) if attempted_lookups and failed_lookups == attempted_lookups and not records: - raise ParserSourceSkipped("Checko arbitration lookups failed for all subjects") + raise ParserSourceSkipped( + "External arbitration lookups failed for all subjects" + ) logger.info( - "Fetched %d arbitration records through Checko for %d subjects", + "Fetched %d arbitration records through external service for %d subjects", len(records), len(subjects), ) @@ -1150,7 +1154,7 @@ def _fetch_checko_arbitration_records( def _checko_arbitration_item(case, *, subject: ArbitrationSubject) -> GenericParserItem: - """Преобразовать дело Checko в generic record.""" + """Преобразовать дело внешнего сервиса в generic record.""" case_number = getattr(case, "case_number", "") or "" filing_date = getattr(case, "filing_date", "") or "" role = _case_role_for_subject(case, subject) @@ -1335,11 +1339,11 @@ def _fetch_checko_registry_inspections( limit: int | None, proxies: list[str] | None, ) -> list[ProverkiInspection]: - """Получить проверки по активным организациям из реестров через Checko.""" + """Получить проверки по активным организациям через внешний сервис данных.""" api_key = getattr(settings, "CHECKO_API_KEY", "") if not api_key: raise ParserSourceSkipped( - "CHECKO_API_KEY is empty; registry inspections parser skipped" + "External provider API key is empty; registry inspections parser skipped" ) resolved_limit = _resolve_lookup_limit( @@ -1351,7 +1355,9 @@ def _fetch_checko_registry_inspections( ), ) if resolved_limit <= 0: - logger.info("Registry inspections Checko parser is disabled by limit=%s", limit) + logger.info( + "Registry inspections external parser is disabled by limit=%s", limit + ) return [] targets = _active_registry_lookup_targets( @@ -1391,7 +1397,7 @@ def _fetch_checko_registry_inspections( failed_lookups += 1 finish_collection(attempt, records_count=0, error=exc) logger.info( - "Checko inspections lookup skipped for target=%s: %s", + "External inspections lookup skipped for target=%s: %s", target.inn or target.ogrn, exc, ) @@ -1402,10 +1408,10 @@ def _fetch_checko_registry_inspections( ) if attempted_lookups and failed_lookups == attempted_lookups and not records: - raise ParserSourceSkipped("Checko inspections lookups failed for all targets") + raise ParserSourceSkipped("External inspections lookups failed for all targets") logger.info( - "Fetched %d inspections through Checko for %d registry organizations", + "Fetched %d inspections through external service for %d registry organizations", len(records), len(targets), ) @@ -1488,7 +1494,7 @@ def _checko_unfair_supplier_items( company, target: RegistryLookupTarget, ) -> list[GenericParserItem]: - """Convert Checko НедобПостЗап values into supplier-bound source records.""" + """Convert external НедобПостЗап values into supplier-bound source records.""" company_inn = _normalize_identifier(getattr(company, "inn", "")) company_ogrn = _normalize_identifier(getattr(company, "ogrn", "")) if target.inn and company_inn != target.inn: @@ -1563,10 +1569,12 @@ def _fetch_checko_unfair_supplier_records( # noqa: C901 proxies: list[str] | None, organization_ids: list[str] | None = None, ) -> list[GenericParserItem]: - """Fetch RNP entries through one Checko /company request per OПК organization.""" + """Fetch RNP entries through one external request per OПК organization.""" api_key = getattr(settings, "CHECKO_API_KEY", "") if not api_key: - raise ParserSourceSkipped("CHECKO_API_KEY is empty; RNP parser skipped") + raise ParserSourceSkipped( + "External provider API key is empty; RNP parser skipped" + ) resolved_limit = _resolve_registry_enrichment_limit(limit) if resolved_limit <= 0: @@ -1604,13 +1612,13 @@ def _fetch_checko_unfair_supplier_records( # noqa: C901 finish_collection(attempt, records_count=0, error=exc) failed_lookups += 1 rate_limited = True - logger.warning("Checko RNP lookup stopped: quota/rate limit reached") + logger.warning("External RNP lookup stopped: quota/rate limit reached") break except CheckoError as exc: finish_collection(attempt, records_count=0, error=exc) failed_lookups += 1 logger.info( - "Checko RNP lookup failed for target=%s: %s", + "External RNP lookup failed for target=%s: %s", target.inn or target.ogrn, exc, ) @@ -1628,7 +1636,7 @@ def _fetch_checko_unfair_supplier_records( # noqa: C901 ) if not records and (rate_limited or failed_lookups == len(targets)): - raise ParserSourceSkipped("Checko RNP lookups failed for all targets") + raise ParserSourceSkipped("External RNP lookups failed for all targets") return records @@ -1729,11 +1737,11 @@ def _fetch_checko_registry_contract_records( limit: int | None, proxies: list[str] | None, ) -> list[GenericParserItem]: - """Получить контракты по активным организациям из реестров через Checko.""" + """Получить контракты по активным организациям через внешний сервис данных.""" api_key = getattr(settings, "CHECKO_API_KEY", "") if not api_key: raise ParserSourceSkipped( - "CHECKO_API_KEY is empty; registry contracts parser skipped" + "External provider API key is empty; registry contracts parser skipped" ) resolved_limit = _resolve_lookup_limit( @@ -1745,7 +1753,7 @@ def _fetch_checko_registry_contract_records( ), ) if resolved_limit <= 0: - logger.info("Registry contracts Checko parser is disabled by limit=%s", limit) + logger.info("Registry contracts external parser is disabled by limit=%s", limit) return [] targets = _active_registry_lookup_targets( @@ -1790,7 +1798,7 @@ def _fetch_checko_registry_contract_records( failed_lookups += 1 target_failures.append(exc) logger.info( - "Checko contracts lookup skipped for target=%s law=%s: %s", + "External contracts lookup skipped for target=%s law=%s: %s", target.inn or target.ogrn, law.value, exc, @@ -1803,10 +1811,10 @@ def _fetch_checko_registry_contract_records( expected_lookups = attempted_lookups * 2 if expected_lookups and failed_lookups == expected_lookups and not records: - raise ParserSourceSkipped("Checko contracts lookups failed for all targets") + raise ParserSourceSkipped("External contracts lookups failed for all targets") logger.info( - "Fetched %d contracts through Checko for %d registry organizations", + "Fetched %d contracts through external service for %d registry organizations", len(records), len(targets), ) @@ -3374,7 +3382,7 @@ def parse_unfair_suppliers( organization_ids: list[str] | None = None, requested_by_id: int | None = None, ) -> dict: - """Checko RNP lookup by default; explicit files remain a manual tool.""" + """External RNP lookup by default; explicit files remain a manual tool.""" proxies = _resolve_proxies(proxies) if file_url or file_path: @@ -3478,7 +3486,7 @@ def parse_registry_inspections( proxies: list[str] | None = None, requested_by_id: int | None = None, ) -> dict: - """Lookup проверок по активным организациям из реестров через Checko.""" + """Lookup проверок по активным организациям через внешний сервис данных.""" proxies = _resolve_proxies(proxies) return _run_inspection_parser( self, diff --git a/src/organizations/management/commands/build_source_record_exports.py b/src/organizations/management/commands/build_source_record_exports.py new file mode 100644 index 0000000..712e48f --- /dev/null +++ b/src/organizations/management/commands/build_source_record_exports.py @@ -0,0 +1,33 @@ +"""Build the complete prepared source-record export matrix.""" + +from __future__ import annotations + +import json + +from apps.core.management.commands.base import BaseAppCommand + +from organizations.source_record_export import build_source_record_export_artifacts + + +class Command(BaseAppCommand): + """Build source-record files synchronously for bootstrap and recovery.""" + + help = "Формирует готовые CSV/XLSX/JSON выгрузки источников" + use_transaction = False + + def execute_command(self, *args, **options) -> str: + generation = build_source_record_export_artifacts() + rendered = json.dumps( + { + "generation_id": generation.generation_id, + "generated_at": generation.generated_at, + "artifacts_count": generation.artifacts_count, + "files_count": generation.files_count, + "records_count": generation.records_count, + "total_size": generation.total_size, + }, + ensure_ascii=False, + sort_keys=True, + ) + self.log_success(rendered) + return rendered diff --git a/src/organizations/migrations/0008_seed_nightly_source_record_exports.py b/src/organizations/migrations/0008_seed_nightly_source_record_exports.py new file mode 100644 index 0000000..b2b9002 --- /dev/null +++ b/src/organizations/migrations/0008_seed_nightly_source_record_exports.py @@ -0,0 +1,61 @@ +import json + +from django.db import migrations + +NIGHTLY_SOURCE_EXPORT_TASK_NAME = "organizations:source-record-exports:nightly-msk" +NIGHTLY_SOURCE_EXPORT_TASK_PATH = ( + "organizations.tasks.refresh_source_record_export_artifacts" +) +NIGHTLY_SOURCE_EXPORT_MSK_CRON = { + "minute": "30", + "hour": "5", + "day_of_week": "*", + "day_of_month": "*", + "month_of_year": "*", + "timezone": "Europe/Moscow", +} + + +def seed_nightly_source_record_export_schedule(apps, schema_editor): + CrontabSchedule = apps.get_model("django_celery_beat", "CrontabSchedule") + PeriodicTask = apps.get_model("django_celery_beat", "PeriodicTask") + + crontab, _ = CrontabSchedule.objects.get_or_create(**NIGHTLY_SOURCE_EXPORT_MSK_CRON) + field_names = {field.name for field in PeriodicTask._meta.fields} + schedule_fields = {"crontab": crontab} + for field_name in ("interval", "solar", "clocked"): + if field_name in field_names: + schedule_fields[field_name] = None + + PeriodicTask.objects.update_or_create( + name=NIGHTLY_SOURCE_EXPORT_TASK_NAME, + defaults={ + "task": NIGHTLY_SOURCE_EXPORT_TASK_PATH, + "args": json.dumps([]), + "kwargs": json.dumps({}), + "enabled": True, + "description": ( + "Nightly preparation of source-record CSV/XLSX/JSON artifacts." + ), + **schedule_fields, + }, + ) + + +def remove_nightly_source_record_export_schedule(apps, schema_editor): + PeriodicTask = apps.get_model("django_celery_beat", "PeriodicTask") + PeriodicTask.objects.filter(name=NIGHTLY_SOURCE_EXPORT_TASK_NAME).delete() + + +class Migration(migrations.Migration): + dependencies = [ + ("django_celery_beat", "0018_improve_crontab_helptext"), + ("organizations", "0007_auto_20260607_1017"), + ] + + operations = [ + migrations.RunPython( + seed_nightly_source_record_export_schedule, + reverse_code=remove_nightly_source_record_export_schedule, + ), + ] diff --git a/src/organizations/serializers.py b/src/organizations/serializers.py index 1ca9309..8c30527 100644 --- a/src/organizations/serializers.py +++ b/src/organizations/serializers.py @@ -12,6 +12,30 @@ from organizations.models import ( ) from organizations.source_record_export import EXPORT_FORMATS +ARBITRATION_ROLE_LABELS = { + "plaintiff": "Истец", + "истец": "Истец", + "defendant": "Ответчик", + "ответчик": "Ответчик", + "third_party": "Третье лицо", + "третье_лицо": "Третье лицо", + "третья_сторона": "Третье лицо", +} + + +def _arbitration_role_label(payload: dict) -> str: + """Return the frontend-facing role label for an arbitration case.""" + target = payload.get("target") + if not isinstance(target, dict): + target = {} + + raw_role = payload.get("role") or target.get("role") or payload.get("party_role") + if not isinstance(raw_role, str): + return "" + + normalized_role = raw_role.strip().casefold().replace("-", "_").replace(" ", "_") + return ARBITRATION_ROLE_LABELS.get(normalized_role, "") + class OrganizationSourceFinancialLineSerializer(serializers.ModelSerializer): """Structured financial line under a source record.""" @@ -62,6 +86,7 @@ class OrganizationSourceRecordSerializer(serializers.ModelSerializer): many=True, read_only=True ) organization = serializers.SerializerMethodField() + payload = serializers.SerializerMethodField() source_group = serializers.CharField( source="extension.source_group", read_only=True ) @@ -98,6 +123,18 @@ class OrganizationSourceRecordSerializer(serializers.ModelSerializer): def get_record_date(self, obj) -> str | None: return getattr(obj, "canonical_record_date", obj.record_date) or None + @swagger_serializer_method( + serializer_or_field=serializers.JSONField(read_only=True), + ) + def get_payload(self, obj) -> dict | list | str | int | float | bool | None: + payload = obj.payload + if obj.extension.source_group != SourceGroup.ARBITRATION: + return payload + + response_payload = dict(payload) if isinstance(payload, dict) else {} + response_payload["role"] = _arbitration_role_label(response_payload) + return response_payload + @swagger_serializer_method( serializer_or_field=OrganizationSourceRecordOrganizationSerializer, ) @@ -155,6 +192,12 @@ class OrganizationSourceRecordExportRequestSerializer(serializers.Serializer): return value +class OrganizationSourceRecordExportDownloadSerializer(serializers.Serializer): + """One-time ticket submitted by a native browser download form.""" + + ticket = serializers.CharField(max_length=64, trim_whitespace=False) + + class OrganizationDirectoryImportUploadSerializer(serializers.Serializer): """Request for uploading the canonical organization directory XLSX.""" diff --git a/src/organizations/source_record_export.py b/src/organizations/source_record_export.py index 7701361..05ed310 100644 --- a/src/organizations/source_record_export.py +++ b/src/organizations/source_record_export.py @@ -1,18 +1,27 @@ -"""ZIP export for organization source records.""" +"""Prepared file exports for organization source records.""" from __future__ import annotations import csv import json +import os +import re +import secrets +import shutil import zipfile -from collections.abc import Sequence +from collections.abc import Iterable, Iterator, Sequence from dataclasses import dataclass -from datetime import date, datetime +from datetime import UTC, date, datetime from decimal import Decimal -from io import BytesIO, StringIO -from typing import Any +from itertools import islice +from pathlib import Path +from tempfile import NamedTemporaryFile +from typing import Any, BinaryIO, cast +from uuid import uuid4 -from django.db.models import QuerySet +from django.conf import settings +from django.core.cache import cache +from django.db.models import QuerySet, prefetch_related_objects from django.utils import timezone from openpyxl import Workbook @@ -23,6 +32,17 @@ 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 +EXPORT_MANIFEST_VERSION = 1 +CURRENT_EXPORT_MANIFEST_FILE_NAME = "current.json" +GENERATION_MANIFEST_FILE_NAME = "manifest.json" +GENERATION_DIRECTORY_NAME = "generations" +SOURCE_RECORD_EXPORT_ITERATOR_CHUNK_SIZE = 1000 +SOURCE_RECORD_EXPORT_ZIP_CHUNK_SIZE = 1024 * 1024 +EXCEL_MAX_DATA_ROWS_PER_SHEET = 1_048_575 +DEFAULT_XLSX_DATA_ROWS_PER_FILE = 100_000 +DEFAULT_DOWNLOAD_TICKET_TTL_SECONDS = 5 * 60 +SOURCE_RECORD_EXPORT_TICKET_CACHE_PREFIX = "organizations:source-record-exports:ticket" +SOURCE_RECORD_EXPORT_TICKET_PATTERN = re.compile(r"[A-Za-z0-9_-]{43}") SOURCE_GROUP_EXPORT_FILE_STEMS: dict[str, str] = { SourceGroup.FINANCIAL_INDICATORS.value: "financial-indicators", @@ -63,55 +83,395 @@ FINANCIAL_LINE_FIELDS = [ ] +class SourceRecordExportArtifactsUnavailable(Exception): + """Raised when there is no complete published export generation.""" + + +class SourceRecordExportTicketInvalid(Exception): + """Raised when a native-download ticket is invalid, expired, or consumed.""" + + +@dataclass(frozen=True) +class SourceRecordExportArtifact: + """One prepared source-group file in a published generation.""" + + source_group: str + file_format: str + file_name: str + path: Path + size: int + records_count: int + part_number: int = 1 + parts_count: int = 1 + + +@dataclass(frozen=True) +class SourceRecordExportGeneration: + """Atomically published set of all source-record export files.""" + + generation_id: str + generated_at: str + artifacts: tuple[SourceRecordExportArtifact, ...] + records_count: int + + @property + def artifacts_count(self) -> int: + return len( + { + (artifact.source_group, artifact.file_format) + for artifact in self.artifacts + } + ) + + @property + def files_count(self) -> int: + return len(self.artifacts) + + @property + def total_size(self) -> int: + return sum(artifact.size for artifact in self.artifacts) + + @dataclass(frozen=True) class SourceRecordExportArchive: - """In-memory source records export archive.""" + """A request-specific ZIP streamed only from prepared files.""" archive_name: str - archive_bytes: bytes + archive_chunks: Iterable[bytes] files_count: int + generated_at: str + + +@dataclass(frozen=True) +class SourceRecordExportDownloadTicket: + """Short-lived capability for one native browser download.""" + + ticket: str + archive_name: str + expires_in: int + + +class _StreamingZipSink: + """Unseekable zipfile target whose written chunks can be drained.""" + + def __init__(self) -> None: + self._offset = 0 + self._chunks: list[bytes] = [] + + def write(self, data: bytes) -> int: + rendered_data = bytes(data) + self._chunks.append(rendered_data) + self._offset += len(rendered_data) + return len(rendered_data) + + def tell(self) -> int: + return self._offset + + def flush(self) -> None: + return None + + def drain(self) -> tuple[bytes, ...]: + chunks = tuple(self._chunks) + self._chunks.clear() + return chunks + + +def build_source_record_export_artifacts( + *, + now: datetime | None = None, + export_directory: str | Path | None = None, +) -> SourceRecordExportGeneration: + """Build all files on disk and atomically publish a new generation.""" + + root_directory = _resolve_export_directory(export_directory) + generations_directory = root_directory / GENERATION_DIRECTORY_NAME + root_directory.mkdir(parents=True, exist_ok=True) + generations_directory.mkdir(parents=True, exist_ok=True) + + generated_at_datetime = _normalize_generation_datetime(now or timezone.now()) + generation_id = ( + f"{generated_at_datetime.strftime('%Y%m%dT%H%M%SZ')}-{uuid4().hex[:8]}" + ) + staging_directory = generations_directory / f".building-{generation_id}" + final_directory = generations_directory / generation_id + staging_directory.mkdir() + + try: + artifacts: list[SourceRecordExportArtifact] = [] + source_record_counts: dict[str, int] = {} + + for source_group in SOURCE_GROUP_EXPORT_FILE_STEMS: + row_spool_path = staging_directory / f".{source_group}.rows.json" + headers, records_count = _spool_source_group_rows( + source_group=source_group, + output_path=row_spool_path, + ) + source_record_counts[source_group] = records_count + + try: + for file_format in _source_group_export_formats(source_group): + file_name = _build_source_group_file_name( + source_group=source_group, + file_format=file_format, + ) + artifact_paths = _render_source_group_artifact( + row_spool_path=row_spool_path, + output_path=staging_directory / file_name, + headers=headers, + file_format=file_format, + records_count=records_count, + ) + parts_count = len(artifact_paths) + artifacts.extend( + SourceRecordExportArtifact( + source_group=source_group, + file_format=file_format, + file_name=artifact_path.name, + path=final_directory / artifact_path.name, + size=artifact_path.stat().st_size, + records_count=records_count, + part_number=part_number, + parts_count=parts_count, + ) + for part_number, artifact_path in enumerate( + artifact_paths, + start=1, + ) + ) + finally: + row_spool_path.unlink(missing_ok=True) + + generation = SourceRecordExportGeneration( + generation_id=generation_id, + generated_at=generated_at_datetime.isoformat(), + artifacts=tuple(artifacts), + records_count=sum(source_record_counts.values()), + ) + manifest_payload = _generation_manifest_payload( + generation, + root_directory=root_directory, + ) + _write_json_file( + staging_directory / GENERATION_MANIFEST_FILE_NAME, + manifest_payload, + ) + os.replace(staging_directory, final_directory) + _write_json_file_atomically( + root_directory / CURRENT_EXPORT_MANIFEST_FILE_NAME, + manifest_payload, + ) + _cleanup_stale_generations( + generations_directory=generations_directory, + current_generation_id=generation_id, + ) + return generation + except Exception: + if staging_directory.exists(): + shutil.rmtree(staging_directory) + raise + + +def load_current_source_record_export_generation( + *, + export_directory: str | Path | None = None, +) -> SourceRecordExportGeneration: + """Load and validate the atomically published current generation.""" + + root_directory = _resolve_export_directory(export_directory) + manifest_path = root_directory / CURRENT_EXPORT_MANIFEST_FILE_NAME + try: + manifest_payload = json.loads(manifest_path.read_text(encoding="utf-8")) + except (FileNotFoundError, json.JSONDecodeError, OSError) as exc: + raise SourceRecordExportArtifactsUnavailable( + "Prepared source-record export is not available." + ) from exc + + return _generation_from_manifest( + manifest_payload, + root_directory=root_directory, + ) def build_source_records_export_archive( *, source_groups: Sequence[str], export_format: str, - now: datetime | None = None, + export_directory: str | Path | None = None, ) -> SourceRecordExportArchive: - """Build a ZIP archive with one source-record file per source group.""" + """Package selected prepared files without querying source-record tables.""" - timestamp = (now or timezone.now()).strftime("%Y%m%d_%H%M%S") - archive_buffer = BytesIO() + root_directory = _resolve_export_directory(export_directory) + generation = load_current_source_record_export_generation( + export_directory=root_directory, + ) + artifacts_by_key: dict[ + tuple[str, str], + list[SourceRecordExportArtifact], + ] = {} + for artifact in generation.artifacts: + artifacts_by_key.setdefault( + (artifact.source_group, artifact.file_format), + [], + ).append(artifact) + selected_artifacts: list[SourceRecordExportArtifact] = [] - with zipfile.ZipFile( - archive_buffer, - mode="w", - compression=zipfile.ZIP_DEFLATED, - ) as archive: - for source_group in source_groups: - file_format = _resolve_source_group_export_format( - source_group=source_group, - requested_format=export_format, - ) - file_name = _build_source_group_file_name( - source_group=source_group, - file_format=file_format, - ) - archive.writestr( - file_name, - _render_source_group_file( - source_group=source_group, - file_format=file_format, - ), + for source_group in source_groups: + file_format = _resolve_source_group_export_format( + source_group=source_group, + requested_format=export_format, + ) + artifacts = artifacts_by_key.get((source_group, file_format), []) + if not artifacts or any(not artifact.path.is_file() for artifact in artifacts): + raise SourceRecordExportArtifactsUnavailable( + f"Prepared export artifact is missing: {source_group}/{file_format}." ) + selected_artifacts.extend( + sorted(artifacts, key=lambda artifact: artifact.part_number) + ) + generated_at = datetime.fromisoformat(generation.generated_at) + timestamp = generated_at.strftime("%Y%m%d_%H%M%S") return SourceRecordExportArchive( archive_name=f"organization_source_records_export_{timestamp}.zip", - archive_bytes=archive_buffer.getvalue(), - files_count=len(source_groups), + archive_chunks=_stream_zip_archive(selected_artifacts), + files_count=len(selected_artifacts), + generated_at=generation.generated_at, ) +def create_source_record_export_download_ticket( + *, + source_groups: Sequence[str], + export_format: str, +) -> SourceRecordExportDownloadTicket: + """Validate prepared files and cache a short-lived download capability.""" + + package = build_source_records_export_archive( + source_groups=source_groups, + export_format=export_format, + ) + expires_in = max( + 1, + int( + getattr( + settings, + "SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS", + DEFAULT_DOWNLOAD_TICKET_TTL_SECONDS, + ) + ), + ) + payload = { + "sources": list(source_groups), + "format": export_format, + } + for _attempt in range(3): + ticket = secrets.token_urlsafe(32) + if cache.add( + _source_record_export_ticket_cache_key(ticket), + payload, + timeout=expires_in, + ): + return SourceRecordExportDownloadTicket( + ticket=ticket, + archive_name=package.archive_name, + expires_in=expires_in, + ) + raise RuntimeError("Could not allocate a source-record export download ticket.") + + +def consume_source_record_export_download_ticket( + ticket: str, +) -> SourceRecordExportArchive: + """Consume a download ticket before streaming the prepared archive.""" + + if not SOURCE_RECORD_EXPORT_TICKET_PATTERN.fullmatch(ticket): + raise SourceRecordExportTicketInvalid + + cache_key = _source_record_export_ticket_cache_key(ticket) + payload = cache.get(cache_key) + if payload is None: + raise SourceRecordExportTicketInvalid + cache.delete(cache_key) + + try: + source_groups = payload["sources"] + export_format = payload["format"] + if ( + not isinstance(source_groups, list) + or not source_groups + or any( + not isinstance(source_group, str) + or source_group not in SOURCE_GROUP_EXPORT_FILE_STEMS + for source_group in source_groups + ) + or len(source_groups) != len(set(source_groups)) + or export_format not in EXPORT_FORMATS + ): + raise ValueError + except (KeyError, TypeError, ValueError): + raise SourceRecordExportTicketInvalid from None + + return build_source_records_export_archive( + source_groups=source_groups, + export_format=export_format, + ) + + +def _source_record_export_ticket_cache_key(ticket: str) -> str: + return f"{SOURCE_RECORD_EXPORT_TICKET_CACHE_PREFIX}:{ticket}" + + +def _stream_zip_archive( + artifacts: Sequence[SourceRecordExportArtifact], +) -> Iterator[bytes]: + sink = _StreamingZipSink() + with zipfile.ZipFile( + cast(BinaryIO, sink), + mode="w", + compression=zipfile.ZIP_STORED, + allowZip64=True, + ) as archive: + for artifact in artifacts: + with artifact.path.open("rb") as source_file: + with archive.open( + artifact.file_name, + mode="w", + force_zip64=True, + ) as archive_entry: + yield from sink.drain() + while chunk := source_file.read( + SOURCE_RECORD_EXPORT_ZIP_CHUNK_SIZE + ): + archive_entry.write(chunk) + yield from sink.drain() + yield from sink.drain() + yield from sink.drain() + + +def _resolve_export_directory(export_directory: str | Path | None) -> Path: + if export_directory is not None: + return Path(export_directory) + + configured_directory = getattr( + settings, + "SOURCE_RECORD_EXPORT_DIRECTORY", + Path(settings.MEDIA_ROOT) / "source-record-exports", + ) + return Path(str(configured_directory)) + + +def _normalize_generation_datetime(value: datetime) -> datetime: + if timezone.is_naive(value): + value = timezone.make_aware(value, UTC) + return value.astimezone(UTC) + + +def _source_group_export_formats(source_group: str) -> tuple[str, ...]: + if source_group == FINANCIAL_SOURCE_GROUP: + return (EXPORT_FORMAT_JSON,) + return EXPORT_FORMATS + + def _resolve_source_group_export_format( *, source_group: str, @@ -130,128 +490,205 @@ def _source_group_queryset(source_group: str) -> QuerySet[OrganizationSourceReco return ( OrganizationSourceRecord.objects.filter(extension__source_group=source_group) .select_related("extension", "extension__organization") - .prefetch_related("financial_lines") - .order_by( - "extension__organization__name", - "source", - "record_date", - "uid", - ) + # Export order is not part of the file contract. Clearing the model's + # default ordering avoids a multi-gigabyte PostgreSQL disk sort for + # source groups with millions of rows. + .order_by() ) -def _render_source_group_file(*, source_group: str, file_format: str) -> bytes: - queryset = _source_group_queryset(source_group) - payload_headers = _collect_payload_headers(queryset) - headers = [ +def _iter_source_records( + *, + source_group: str, + include_financial_lines: bool, +) -> Iterator[OrganizationSourceRecord]: + iterator = _source_group_queryset(source_group).iterator( + chunk_size=SOURCE_RECORD_EXPORT_ITERATOR_CHUNK_SIZE + ) + while True: + batch = list(islice(iterator, SOURCE_RECORD_EXPORT_ITERATOR_CHUNK_SIZE)) + if not batch: + return + if include_financial_lines: + prefetch_related_objects(batch, "financial_lines") + yield from batch + + +def _spool_source_group_rows( + *, + source_group: str, + output_path: Path, +) -> tuple[list[str], int]: + include_financial_lines = source_group == FINANCIAL_SOURCE_GROUP + payload_headers: set[str] = set() + records_count = 0 + + with output_path.open("w", encoding="utf-8", newline="") as output: + output.write("[") + is_first_row = True + for record in _iter_source_records( + source_group=source_group, + include_financial_lines=include_financial_lines, + ): + row = _build_record_row( + record, + include_financial_lines=include_financial_lines, + ) + payload_headers.update(key for key in row if key.startswith("payload.")) + if is_first_row: + output.write("\n") + is_first_row = False + else: + output.write(",\n") + output.write(json.dumps(row, ensure_ascii=False, separators=(",", ":"))) + records_count += 1 + if not is_first_row: + output.write("\n") + output.write("]") + + return [ *ORGANIZATION_EXPORT_FIELDS, *SOURCE_RECORD_EXPORT_FIELDS, - *payload_headers, - ] - include_financial_lines = source_group == FINANCIAL_SOURCE_GROUP + *sorted(payload_headers), + ], records_count + +def _render_source_group_artifact( + *, + row_spool_path: Path, + output_path: Path, + headers: Sequence[str], + file_format: str, + records_count: int, +) -> tuple[Path, ...]: if file_format == EXPORT_FORMAT_CSV: - return _render_csv_file( - queryset=_source_group_queryset(source_group), + _render_csv_file( + row_spool_path=row_spool_path, + output_path=output_path, headers=headers, - payload_headers=payload_headers, ) + return (output_path,) if file_format == EXPORT_FORMAT_XLSX: - return _render_xlsx_file( - queryset=_source_group_queryset(source_group), + return _render_xlsx_files( + row_spool_path=row_spool_path, + output_path=output_path, headers=headers, - payload_headers=payload_headers, + records_count=records_count, ) - return _render_json_file( - queryset=_source_group_queryset(source_group), - headers=headers, - payload_headers=payload_headers, - include_financial_lines=include_financial_lines, - ) + _render_json_file(row_spool_path=row_spool_path, output_path=output_path) + return (output_path,) -def _collect_payload_headers( - queryset: QuerySet[OrganizationSourceRecord], -) -> list[str]: - payload_headers: set[str] = set() - - for record in queryset.iterator(chunk_size=1000): - payload_headers.update(_flatten_payload(record.payload).keys()) - - return sorted(payload_headers) +def _iter_spooled_rows(row_spool_path: Path) -> Iterator[dict[str, Any]]: + with row_spool_path.open("r", encoding="utf-8") as rows_file: + for line in rows_file: + serialized_row = line.strip() + if not serialized_row or serialized_row in {"[", "]", "[]"}: + continue + if serialized_row.endswith(","): + serialized_row = serialized_row[:-1] + yield json.loads(serialized_row) def _render_csv_file( *, - queryset: QuerySet[OrganizationSourceRecord], + row_spool_path: Path, + output_path: Path, headers: Sequence[str], - payload_headers: Sequence[str], -) -> bytes: - output = StringIO() - writer = csv.DictWriter(output, fieldnames=list(headers), lineterminator="\n") - writer.writeheader() - - for record in queryset.iterator(chunk_size=1000): - row = _build_record_row(record, payload_headers=payload_headers) - writer.writerow({key: _serialize_flat_value(row.get(key)) for key in headers}) - - return ("\ufeff" + output.getvalue()).encode("utf-8") +) -> None: + with output_path.open("w", encoding="utf-8-sig", newline="") as output: + writer = csv.DictWriter(output, fieldnames=list(headers), lineterminator="\n") + writer.writeheader() + for row in _iter_spooled_rows(row_spool_path): + writer.writerow( + {key: _serialize_flat_value(row.get(key)) for key in headers} + ) -def _render_xlsx_file( +def _render_xlsx_files( *, - queryset: QuerySet[OrganizationSourceRecord], + row_spool_path: Path, + output_path: Path, headers: Sequence[str], - payload_headers: Sequence[str], -) -> bytes: + records_count: int, +) -> tuple[Path, ...]: + rows_per_file = min( + EXCEL_MAX_DATA_ROWS_PER_SHEET, + max( + 1, + int( + getattr( + settings, + "SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE", + DEFAULT_XLSX_DATA_ROWS_PER_FILE, + ) + ), + ), + ) + parts_count = max(1, (records_count + rows_per_file - 1) // rows_per_file) + part_paths = tuple( + _build_xlsx_part_path( + output_path=output_path, + part_number=part_number, + parts_count=parts_count, + ) + for part_number in range(1, parts_count + 1) + ) + part_number = 1 + rows_in_file = 0 + workbook, worksheet = _new_export_workbook(headers) + + for row in _iter_spooled_rows(row_spool_path): + if rows_in_file >= rows_per_file: + workbook.save(part_paths[part_number - 1]) + workbook.close() + 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]) + rows_in_file += 1 + + workbook.save(part_paths[part_number - 1]) + workbook.close() + if part_number != parts_count: + raise ValueError("Unexpected XLSX source-record export parts count.") + return part_paths + + +def _new_export_workbook(headers: Sequence[str]): workbook = Workbook(write_only=True) worksheet = workbook.create_sheet(title="data") worksheet.append(list(headers)) + return workbook, worksheet - for record in queryset.iterator(chunk_size=1000): - row = _build_record_row(record, payload_headers=payload_headers) - worksheet.append([_serialize_flat_value(row.get(key)) for key in headers]) - output = BytesIO() - workbook.save(output) - return output.getvalue() +def _build_xlsx_part_path( + *, + output_path: Path, + part_number: int, + parts_count: int, +) -> Path: + if parts_count == 1: + return output_path + return output_path.with_name( + f"{output_path.stem}-part-{part_number:03d}{output_path.suffix}" + ) def _render_json_file( *, - queryset: QuerySet[OrganizationSourceRecord], - headers: Sequence[str], - payload_headers: Sequence[str], - include_financial_lines: bool, -) -> bytes: - rows = [] - - for record in queryset.iterator(chunk_size=1000): - row = _build_record_row(record, payload_headers=payload_headers) - json_row = {key: _serialize_json_value(row.get(key)) for key in headers} - if include_financial_lines: - json_row["financial_lines"] = [ - { - field_name: _serialize_json_value( - getattr(financial_line, field_name) - ) - for field_name in FINANCIAL_LINE_FIELDS - } - for financial_line in record.financial_lines.all() - ] - rows.append(json_row) - - return json.dumps(rows, ensure_ascii=False, indent=2).encode("utf-8") + row_spool_path: Path, + output_path: Path, +) -> None: + os.link(row_spool_path, output_path) def _build_record_row( record: OrganizationSourceRecord, *, - payload_headers: Sequence[str], + include_financial_lines: bool, ) -> dict[str, Any]: organization = record.extension.organization - payload = _flatten_payload(record.payload) - row: dict[str, Any] = { "Наименование": organization.full_name or organization.short_name @@ -272,12 +709,19 @@ def _build_record_row( "load_batch": record.load_batch, "created_at": record.created_at, "updated_at": record.updated_at, + **_flatten_payload(record.payload), } + serialized_row = {key: _serialize_json_value(value) for key, value in row.items()} - for header in payload_headers: - row[header] = payload.get(header) - - return row + if include_financial_lines: + serialized_row["financial_lines"] = [ + { + field_name: _serialize_json_value(getattr(financial_line, field_name)) + for field_name in FINANCIAL_LINE_FIELDS + } + for financial_line in record.financial_lines.all() + ] + return serialized_row def _flatten_payload(value: Any, *, prefix: str = "payload") -> dict[str, Any]: @@ -303,7 +747,9 @@ def _serialize_flat_value(value: Any) -> str | int | float | bool: return str(value) if isinstance(value, datetime | date): return value.isoformat() - return value + if isinstance(value, str | int | float | bool): + return value + return str(value) def _serialize_json_value(value: Any) -> Any: @@ -314,3 +760,191 @@ def _serialize_json_value(value: Any) -> Any: if isinstance(value, datetime | date): return value.isoformat() return value + + +def _generation_manifest_payload( + generation: SourceRecordExportGeneration, + *, + root_directory: Path, +) -> dict[str, Any]: + return { + "version": EXPORT_MANIFEST_VERSION, + "generation_id": generation.generation_id, + "generated_at": generation.generated_at, + "records_count": generation.records_count, + "artifacts_count": generation.artifacts_count, + "files_count": generation.files_count, + "total_size": generation.total_size, + "artifacts": [ + { + "source_group": artifact.source_group, + "format": artifact.file_format, + "file_name": artifact.file_name, + "relative_path": str(artifact.path.relative_to(root_directory)), + "size": artifact.size, + "records_count": artifact.records_count, + "part_number": artifact.part_number, + "parts_count": artifact.parts_count, + } + for artifact in generation.artifacts + ], + } + + +def _generation_from_manifest( + payload: Any, + *, + root_directory: Path, +) -> SourceRecordExportGeneration: + try: + if payload["version"] != EXPORT_MANIFEST_VERSION: + raise ValueError("Unsupported source-record export manifest version.") + generation_id = str(payload["generation_id"]) + generated_at = str(payload["generated_at"]) + datetime.fromisoformat(generated_at) + records_count = int(payload["records_count"]) + artifact_payloads = payload["artifacts"] + if not isinstance(artifact_payloads, list): + raise TypeError("Manifest artifacts must be a list.") + + root_resolved = root_directory.resolve() + artifacts: list[SourceRecordExportArtifact] = [] + for artifact_payload in artifact_payloads: + artifact_path = root_directory / str(artifact_payload["relative_path"]) + artifact_path.resolve().relative_to(root_resolved) + if not artifact_path.is_file(): + raise FileNotFoundError(artifact_path) + source_group = str(artifact_payload["source_group"]) + file_format = str(artifact_payload["format"]) + file_name = str(artifact_payload["file_name"]) + artifact_size = int(artifact_payload["size"]) + part_number = int(artifact_payload.get("part_number", 1)) + parts_count = int(artifact_payload.get("parts_count", 1)) + expected_file_name = _build_source_group_file_name( + source_group=source_group, + file_format=file_format, + ) + if file_format == EXPORT_FORMAT_XLSX: + expected_file_name = _build_xlsx_part_path( + output_path=Path(expected_file_name), + part_number=part_number, + parts_count=parts_count, + ).name + if file_name != expected_file_name: + raise ValueError("Unexpected source-record export artifact name.") + if artifact_size != artifact_path.stat().st_size: + raise ValueError("Source-record export artifact size mismatch.") + artifacts.append( + SourceRecordExportArtifact( + source_group=source_group, + file_format=file_format, + file_name=file_name, + path=artifact_path, + size=artifact_size, + records_count=int(artifact_payload["records_count"]), + part_number=part_number, + parts_count=parts_count, + ) + ) + + _validate_manifest_artifact_parts(artifacts) + except (KeyError, TypeError, ValueError, OSError) as exc: + raise SourceRecordExportArtifactsUnavailable( + "Prepared source-record export manifest is invalid." + ) from exc + + return SourceRecordExportGeneration( + generation_id=generation_id, + generated_at=generated_at, + artifacts=tuple(artifacts), + records_count=records_count, + ) + + +def _validate_manifest_artifact_parts( + artifacts: Sequence[SourceRecordExportArtifact], +) -> None: + expected_artifact_keys = { + (source_group, file_format) + for source_group in SOURCE_GROUP_EXPORT_FILE_STEMS + for file_format in _source_group_export_formats(source_group) + } + artifacts_by_key: dict[ + tuple[str, str], + list[SourceRecordExportArtifact], + ] = {} + for artifact in artifacts: + artifacts_by_key.setdefault( + (artifact.source_group, artifact.file_format), + [], + ).append(artifact) + if set(artifacts_by_key) != expected_artifact_keys: + raise ValueError("Source-record export manifest matrix is incomplete.") + + for artifact_key, artifact_parts in artifacts_by_key.items(): + parts_count = len(artifact_parts) + if ( + {artifact.parts_count for artifact in artifact_parts} != {parts_count} + or {artifact.part_number for artifact in artifact_parts} + != set(range(1, parts_count + 1)) + or (artifact_key[1] != EXPORT_FORMAT_XLSX and parts_count != 1) + ): + raise ValueError("Source-record export artifact parts are invalid.") + + +def _write_json_file(file_path: Path, payload: dict[str, Any]) -> None: + file_path.write_text( + json.dumps(payload, ensure_ascii=False, indent=2), + encoding="utf-8", + ) + + +def _write_json_file_atomically(file_path: Path, payload: dict[str, Any]) -> None: + temp_path: Path | None = None + try: + with NamedTemporaryFile( + mode="w", + encoding="utf-8", + dir=file_path.parent, + prefix=f".{file_path.name}.", + suffix=".tmp", + delete=False, + ) as temp_file: + json.dump(payload, temp_file, ensure_ascii=False, indent=2) + temp_file.flush() + os.fsync(temp_file.fileno()) + temp_path = Path(temp_file.name) + os.replace(temp_path, file_path) + finally: + if temp_path is not None: + temp_path.unlink(missing_ok=True) + + +def _cleanup_stale_generations( + *, + generations_directory: Path, + current_generation_id: str, +) -> None: + generations_to_keep = max( + 1, + int(getattr(settings, "SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP", 2)), + ) + published_generations = sorted( + ( + path + for path in generations_directory.iterdir() + if path.is_dir() + and not path.name.startswith(".") + and (path / GENERATION_MANIFEST_FILE_NAME).is_file() + ), + key=lambda path: path.name, + reverse=True, + ) + retained_names = {current_generation_id} + retained_names.update( + path.name for path in published_generations[:generations_to_keep] + ) + + for generation_directory in published_generations: + if generation_directory.name not in retained_names: + shutil.rmtree(generation_directory) diff --git a/src/organizations/tasks.py b/src/organizations/tasks.py index 4d50d8a..3bf76c2 100644 --- a/src/organizations/tasks.py +++ b/src/organizations/tasks.py @@ -5,14 +5,50 @@ from __future__ import annotations import logging from dataclasses import asdict +from apps.core.tasks import PeriodicTask as CorePeriodicTask from celery import shared_task +from django.conf import settings +from django.core.cache import cache from organizations.cache import invalidate_organization_api_cache from organizations.source_backfill import OrganizationSourceBackfillService +from organizations.source_record_export import build_source_record_export_artifacts logger = logging.getLogger(__name__) +@shared_task(bind=True, base=CorePeriodicTask) +def refresh_source_record_export_artifacts(self) -> dict: # noqa: ARG001 + """Build and atomically publish the nightly source-record export matrix.""" + lock_key = getattr( + settings, + "SOURCE_RECORD_EXPORT_LOCK_KEY", + "organizations:source-record-exports:lock", + ) + lock_ttl = int( + getattr(settings, "SOURCE_RECORD_EXPORT_LOCK_TTL_SECONDS", 6 * 60 * 60) + ) + if not cache.add(lock_key, "1", timeout=lock_ttl): + logger.info("Source-record export generation skipped: lock is already held") + return {"status": "skipped", "reason": "locked"} + + try: + generation = build_source_record_export_artifacts() + result = { + "status": "success", + "generation_id": generation.generation_id, + "generated_at": generation.generated_at, + "artifacts_count": generation.artifacts_count, + "files_count": generation.files_count, + "records_count": generation.records_count, + "total_size": generation.total_size, + } + logger.info("Source-record export generation published: %s", result) + return result + finally: + cache.delete(lock_key) + + @shared_task def backfill_all_organization_sources(batch_size: int = 100) -> dict: """Backfill all organization source extensions from legacy parser tables.""" diff --git a/src/organizations/views.py b/src/organizations/views.py index acb1583..0cce862 100644 --- a/src/organizations/views.py +++ b/src/organizations/views.py @@ -16,7 +16,7 @@ from django.core.cache import cache from django.db.models import Case, CharField, F, Q, Value, When from django.db.models.fields.json import KeyTextTransform from django.db.models.functions import Cast, Coalesce, NullIf -from django.http import HttpResponse +from django.http import StreamingHttpResponse from django_filters import rest_framework as filters from drf_yasg import openapi from drf_yasg.utils import swagger_auto_schema @@ -51,11 +51,19 @@ from organizations.serializers import ( OrganizationDirectoryImportUploadSerializer, OrganizationSerializer, OrganizationSourceExtensionSerializer, + OrganizationSourceRecordExportDownloadSerializer, OrganizationSourceRecordExportRequestSerializer, OrganizationSourceRecordListResponseSerializer, OrganizationSourceRecordSerializer, ) -from organizations.source_record_export import build_source_records_export_archive +from organizations.source_record_export import ( + SourceRecordExportArchive, + SourceRecordExportArtifactsUnavailable, + SourceRecordExportTicketInvalid, + build_source_records_export_archive, + consume_source_record_export_download_ticket, + create_source_record_export_download_ticket, +) ORGANIZATIONS_TAG = swagger_tag("Организации", "Organizations") @@ -536,12 +544,29 @@ class OrganizationSourceRecordViewSet(ReadOnlyModelViewSet): ordering = ["-created_at", "-uid"] def get_permissions(self): - if self.action == "export": + if self.action in {"export", "export_ticket"}: return [IsAdminUser()] + if self.action == "export_download": + return [AllowAny()] if getattr(settings, "ORGANIZATIONS_V2_ALLOW_ANONYMOUS", False): return [AllowAny()] return super().get_permissions() + @staticmethod + def _source_record_export_response( + package: SourceRecordExportArchive, + ) -> StreamingHttpResponse: + response = StreamingHttpResponse( + package.archive_chunks, + content_type="application/zip", + ) + response[ + "Content-Disposition" + ] = f'attachment; filename="{package.archive_name}"' + response["X-Source-Export-Files"] = str(package.files_count) + response["X-Source-Export-Generated-At"] = package.generated_at + return response + def get_queryset(self): raw_record_date = NullIf(F("record_date"), Value("")) inspection_record_date = Coalesce( @@ -732,7 +757,8 @@ class OrganizationSourceRecordViewSet(ReadOnlyModelViewSet): operation_id="v2_organization_source_records_export", operation_summary="Выгрузить записи источников", operation_description=( - "Формирует ZIP-архив с одним файлом на выбранную группу источников. " + "Упаковывает в ZIP ночные готовые файлы выбранных групп источников. " + "При скачивании данные из БД повторно не формируются. " "Финансово-экономические показатели всегда выгружаются в JSON." ), request_body=OrganizationSourceRecordExportRequestSerializer, @@ -743,22 +769,136 @@ class OrganizationSourceRecordViewSet(ReadOnlyModelViewSet): ), 400: "Некорректные параметры выгрузки.", 403: "Доступ разрешён только администраторам.", + 503: "Ночная выгрузка ещё не сформирована.", }, ) @action(detail=False, methods=["post"], url_path="export") - def export(self, request, *args: Any, **kwargs: Any) -> HttpResponse: + def export( + self, + request, + *args: Any, + **kwargs: Any, + ) -> StreamingHttpResponse | Response: serializer = OrganizationSourceRecordExportRequestSerializer(data=request.data) serializer.is_valid(raise_exception=True) - package = build_source_records_export_archive( - source_groups=serializer.validated_data["sources"], - export_format=serializer.validated_data["format"], + try: + package = build_source_records_export_archive( + source_groups=serializer.validated_data["sources"], + export_format=serializer.validated_data["format"], + ) + except SourceRecordExportArtifactsUnavailable: + return Response( + { + "detail": "Готовая ночная выгрузка ещё не сформирована.", + "code": "source_export_not_ready", + }, + status=status.HTTP_503_SERVICE_UNAVAILABLE, + headers={"Retry-After": "3600"}, + ) + + return self._source_record_export_response(package) + + @swagger_auto_schema( + tags=[ORGANIZATIONS_TAG], + operation_id="v2_organization_source_records_export_ticket", + operation_summary="Подготовить нативное скачивание записей источников", + operation_description=( + "Возвращает короткоживущий одноразовый ticket. Frontend отправляет " + "его обычной HTML-формой в export-download, чтобы браузер сохранял " + "потоковый ZIP напрямую на диск без Blob в JavaScript." + ), + request_body=OrganizationSourceRecordExportRequestSerializer, + responses={ + 201: "Одноразовый ticket и имя ZIP-файла.", + 400: "Некорректные параметры выгрузки.", + 403: "Доступ разрешён только администраторам.", + 503: "Ночная выгрузка ещё не сформирована.", + }, + ) + @action(detail=False, methods=["post"], url_path="export-ticket") + def export_ticket(self, request, *args: Any, **kwargs: Any) -> Response: + serializer = OrganizationSourceRecordExportRequestSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + + try: + download_ticket = create_source_record_export_download_ticket( + source_groups=serializer.validated_data["sources"], + export_format=serializer.validated_data["format"], + ) + except SourceRecordExportArtifactsUnavailable: + return Response( + { + "detail": "Готовая ночная выгрузка ещё не сформирована.", + "code": "source_export_not_ready", + }, + status=status.HTTP_503_SERVICE_UNAVAILABLE, + headers={"Retry-After": "3600"}, + ) + except RuntimeError: + return Response( + { + "detail": "Не удалось подготовить скачивание. Повторите запрос.", + "code": "source_export_ticket_unavailable", + }, + status=status.HTTP_503_SERVICE_UNAVAILABLE, + headers={"Retry-After": "5"}, + ) + + return Response( + { + "ticket": download_ticket.ticket, + "file_name": download_ticket.archive_name, + "expires_in": download_ticket.expires_in, + }, + status=status.HTTP_201_CREATED, ) - response = HttpResponse(package.archive_bytes, content_type="application/zip") - response.status_code = status.HTTP_200_OK - response[ - "Content-Disposition" - ] = f'attachment; filename="{package.archive_name}"' - response["Content-Length"] = str(len(package.archive_bytes)) - response["X-Source-Export-Files"] = str(package.files_count) - return response + + @swagger_auto_schema( + tags=[ORGANIZATIONS_TAG], + operation_id="v2_organization_source_records_export_download", + operation_summary="Скачать готовые записи источников по ticket", + request_body=OrganizationSourceRecordExportDownloadSerializer, + responses={ + 200: openapi.Response( + description="Потоковый ZIP-архив записей источников.", + schema=openapi.Schema(type=openapi.TYPE_FILE), + ), + 400: "Ticket отсутствует или имеет неверный формат.", + 410: "Ticket истёк или уже использован.", + 503: "Опубликованная выгрузка больше недоступна.", + }, + ) + @action(detail=False, methods=["post"], url_path="export-download") + def export_download( + self, + request, + *args: Any, + **kwargs: Any, + ) -> StreamingHttpResponse | Response: + serializer = OrganizationSourceRecordExportDownloadSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + + try: + package = consume_source_record_export_download_ticket( + serializer.validated_data["ticket"] + ) + except SourceRecordExportTicketInvalid: + return Response( + { + "detail": "Ticket скачивания истёк или уже использован.", + "code": "source_export_ticket_invalid", + }, + status=status.HTTP_410_GONE, + ) + except SourceRecordExportArtifactsUnavailable: + return Response( + { + "detail": "Опубликованная выгрузка больше недоступна.", + "code": "source_export_not_ready", + }, + status=status.HTTP_503_SERVICE_UNAVAILABLE, + headers={"Retry-After": "3600"}, + ) + + return self._source_record_export_response(package) diff --git a/src/settings/base.py b/src/settings/base.py index 9322919..41a96f1 100644 --- a/src/settings/base.py +++ b/src/settings/base.py @@ -245,6 +245,26 @@ BACKUP_EXPORT_DIRECTORY = os.getenv( "BACKUP_EXPORT_DIRECTORY", str(PROJECT_ROOT / "media" / "backups"), ) +SOURCE_RECORD_EXPORT_DIRECTORY = os.getenv( + "SOURCE_RECORD_EXPORT_DIRECTORY", + str(PROJECT_ROOT / "media" / "source-record-exports"), +) +SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP = int( + os.getenv("SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP", "2") +) +SOURCE_RECORD_EXPORT_LOCK_KEY = os.getenv( + "SOURCE_RECORD_EXPORT_LOCK_KEY", + "organizations:source-record-exports:lock", +) +SOURCE_RECORD_EXPORT_LOCK_TTL_SECONDS = int( + os.getenv("SOURCE_RECORD_EXPORT_LOCK_TTL_SECONDS", str(6 * 60 * 60)) +) +SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE = int( + os.getenv("SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE", "100000") +) +SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS = int( + os.getenv("SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS", "300") +) # Celery: сохраняем ретраи подключения на старте и для 6.x совместимости. CELERY_BROKER_CONNECTION_RETRY = True @@ -453,7 +473,7 @@ FNS_PROCESSED_DIRECTORY = PROJECT_ROOT / "input" / "fns" / "processed" FNS_FAILED_DIRECTORY = PROJECT_ROOT / "input" / "fns" / "failed" # ============================================================================= -# Checko API Settings (checko.ru) +# External organization data API settings # ============================================================================= CHECKO_API_KEY = os.getenv("CHECKO_API_KEY", "") diff --git a/tests/apps/organizations/test_api_v2_source_extensions.py b/tests/apps/organizations/test_api_v2_source_extensions.py index 3a6740b..5e5bf86 100644 --- a/tests/apps/organizations/test_api_v2_source_extensions.py +++ b/tests/apps/organizations/test_api_v2_source_extensions.py @@ -170,6 +170,59 @@ class OrganizationSourceExtensionsApiV2Test(APITestCase): self.assertEqual(record["source_group"], "planned_inspections") self.assertEqual(record["organization"]["uid"], str(target.uid)) + def test_flat_arbitration_records_expose_frontend_role_labels(self): + organization = create_frontend_organization( + name='ООО "Arbitration Roles"', + inn="7707083899", + ogrn="1027700132099", + ) + extension = ArbitrationExtension.objects.create( + organization=organization, + title="Арбитраж", + ) + expected_roles = { + "ROLE-PLAINTIFF": "Истец", + "ROLE-DEFENDANT": "Ответчик", + "ROLE-THIRD-PARTY": "Третье лицо", + } + for external_id, provider_role in ( + ("ROLE-PLAINTIFF", "plaintiff"), + ("ROLE-DEFENDANT", "defendant"), + ("ROLE-THIRD-PARTY", "third_party"), + ): + OrganizationSourceRecord.objects.create( + extension=extension, + record_type="arbitration_case", + source="arbitration", + external_id=external_id, + payload={"target": {"role": provider_role}}, + ) + + response = self.client.get( + reverse("api_v2:organizations:organization-source-records-list"), + { + "source_group": "arbitration", + "source": "arbitration", + "organization": str(organization.uid), + }, + ) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + records_by_external_id = { + item["external_id"]: item for item in response.data["data"] + } + self.assertEqual( + { + external_id: records_by_external_id[external_id]["payload"]["role"] + for external_id in expected_roles + }, + expected_roles, + ) + self.assertNotIn( + "role", + OrganizationSourceRecord.objects.get(external_id="ROLE-PLAINTIFF").payload, + ) + def test_flat_source_records_filters_supported_groups_by_canonical_date(self): organization = create_frontend_organization( name='ООО "Canonical dates"', diff --git a/tests/apps/organizations/test_source_record_export.py b/tests/apps/organizations/test_source_record_export.py index 2217489..1f48174 100644 --- a/tests/apps/organizations/test_source_record_export.py +++ b/tests/apps/organizations/test_source_record_export.py @@ -2,9 +2,15 @@ import csv import json +import os import zipfile from io import BytesIO, StringIO +from pathlib import Path +from tempfile import TemporaryDirectory +from unittest.mock import patch +from django.core.management import call_command +from django.test import override_settings from django.urls import reverse from openpyxl import load_workbook from organizations.models import ( @@ -15,6 +21,14 @@ from organizations.models import ( PlannedInspectionExtension, SourceGroup, ) +from organizations.source_record_export import ( + _render_source_group_artifact, + _source_group_queryset, + _spool_source_group_rows, + build_source_record_export_artifacts, + build_source_records_export_archive, + load_current_source_record_export_generation, +) from rest_framework import status from rest_framework.test import APITestCase @@ -25,9 +39,57 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): """Checks admin-only source-record export contract.""" def setUp(self): + self.export_directory = TemporaryDirectory() + self.settings_override = override_settings( + SOURCE_RECORD_EXPORT_DIRECTORY=self.export_directory.name, + SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP=2, + ) + self.settings_override.enable() self.url = reverse( "api_v2:organizations:organization-source-records-export", ) + self.ticket_url = reverse( + "api_v2:organizations:organization-source-records-export-ticket", + ) + self.download_url = reverse( + "api_v2:organizations:organization-source-records-export-download", + ) + + def tearDown(self): + self.settings_override.disable() + self.export_directory.cleanup() + super().tearDown() + + @staticmethod + def _response_body(response) -> bytes: + if response.streaming: + return b"".join(response.streaming_content) + return response.content + + def test_export_returns_service_unavailable_before_first_nightly_generation(self): + self.client.force_authenticate(UserFactory.create_superuser()) + + response = self.client.post( + self.url, + { + "sources": [SourceGroup.PLANNED_INSPECTIONS.value], + "format": "json", + }, + format="json", + ) + + self.assertEqual(response.status_code, status.HTTP_503_SERVICE_UNAVAILABLE) + self.assertEqual(response.data["code"], "source_export_not_ready") + self.assertEqual(response["Retry-After"], "3600") + + def test_management_command_bootstraps_first_generation(self): + command_output = StringIO() + + call_command("build_source_record_exports", stdout=command_output) + + generation = load_current_source_record_export_generation() + self.assertEqual(generation.artifacts_count, 25) + self.assertIn('"artifacts_count": 25', command_output.getvalue()) def test_admin_exports_selected_sources_to_zip(self): self.client.force_authenticate(UserFactory.create_superuser()) @@ -74,27 +136,34 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): period_start=100, period_end=200, ) + generation = build_source_record_export_artifacts() - response = self.client.post( - self.url, - { - "sources": [ - SourceGroup.PLANNED_INSPECTIONS.value, - SourceGroup.FINANCIAL_INDICATORS.value, - ], - "format": "xlsx", - }, - format="json", - ) + with self.assertNumQueries(0): + response = self.client.post( + self.url, + { + "sources": [ + SourceGroup.PLANNED_INSPECTIONS.value, + SourceGroup.FINANCIAL_INDICATORS.value, + ], + "format": "xlsx", + }, + format="json", + ) self.assertEqual(response.status_code, status.HTTP_200_OK) + self.assertTrue(response.streaming) self.assertEqual(response["Content-Type"], "application/zip") + self.assertEqual( + response["X-Source-Export-Generated-At"], generation.generated_at + ) + self.assertEqual(generation.artifacts_count, 25) self.assertIn( 'filename="organization_source_records_export_', response["Content-Disposition"], ) - with zipfile.ZipFile(BytesIO(response.content)) as archive: + with zipfile.ZipFile(BytesIO(self._response_body(response))) as archive: self.assertEqual( set(archive.namelist()), {"planned-inspections.xlsx", "financial-indicators.json"}, @@ -131,6 +200,53 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): ) self.assertEqual(financial_rows[0]["financial_lines"][0]["period_end"], 200) + def test_admin_uses_one_time_ticket_for_native_zero_sql_download(self): + self.client.force_authenticate(UserFactory.create_superuser()) + build_source_record_export_artifacts() + + with self.assertNumQueries(0): + ticket_response = self.client.post( + self.ticket_url, + { + "sources": [SourceGroup.PLANNED_INSPECTIONS.value], + "format": "json", + }, + format="json", + ) + + self.assertEqual(ticket_response.status_code, status.HTTP_201_CREATED) + self.assertEqual(ticket_response.data["expires_in"], 300) + self.assertRegex(ticket_response.data["ticket"], r"^[A-Za-z0-9_-]{43}$") + self.assertNotIn("download_url", ticket_response.data) + + self.client.force_authenticate(user=None) + with self.assertNumQueries(0): + download_response = self.client.post( + self.download_url, + {"ticket": ticket_response.data["ticket"]}, + format="multipart", + ) + + self.assertEqual(download_response.status_code, status.HTTP_200_OK) + self.assertTrue(download_response.streaming) + self.assertEqual(download_response["Content-Type"], "application/zip") + self.assertNotIn("Content-Length", download_response) + with zipfile.ZipFile( + BytesIO(self._response_body(download_response)) + ) as archive: + self.assertEqual(archive.namelist(), ["planned-inspections.json"]) + + consumed_response = self.client.post( + self.download_url, + {"ticket": ticket_response.data["ticket"]}, + format="multipart", + ) + self.assertEqual(consumed_response.status_code, status.HTTP_410_GONE) + self.assertEqual( + consumed_response.data["code"], + "source_export_ticket_invalid", + ) + def test_csv_export_uses_bom_and_canonical_columns_before_payload(self): self.client.force_authenticate(UserFactory.create_superuser()) organization = Organization.objects.create( @@ -152,6 +268,7 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): title="CSV проверка", payload={"nested": {"value": "данные"}}, ) + build_source_record_export_artifacts() response = self.client.post( self.url, @@ -164,7 +281,7 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): self.assertEqual(response.status_code, status.HTTP_200_OK) - with zipfile.ZipFile(BytesIO(response.content)) as archive: + with zipfile.ZipFile(BytesIO(self._response_body(response))) as archive: csv_bytes = archive.read("planned-inspections.csv") self.assertTrue(csv_bytes.startswith(b"\xef\xbb\xbf")) csv_text = csv_bytes.decode("utf-8-sig") @@ -176,6 +293,178 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): ) self.assertIn("payload.nested.value", csv_rows[0]) + def test_xlsx_export_splits_rows_across_bounded_workbook_parts(self): + organization = Organization.objects.create( + name='ООО "Многолистовая выгрузка"', + inn="7707083812", + ) + extension = PlannedInspectionExtension.objects.create( + organization=organization, + title="Плановые проверки Генпрокуратуры России", + ) + for index in range(3): + OrganizationSourceRecord.objects.create( + extension=extension, + record_type="inspection", + source="inspections", + external_id=f"INSP-SHEET-{index}", + title=f"Проверка {index}", + payload={}, + ) + + with override_settings(SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE=2): + generation = build_source_record_export_artifacts() + + artifacts = sorted( + ( + item + for item in generation.artifacts + if item.source_group == SourceGroup.PLANNED_INSPECTIONS.value + and item.file_format == "xlsx" + ), + key=lambda item: item.file_name, + ) + + self.assertEqual(generation.artifacts_count, 25) + self.assertEqual(generation.files_count, 26) + self.assertEqual( + [item.file_name for item in artifacts], + [ + "planned-inspections-part-001.xlsx", + "planned-inspections-part-002.xlsx", + ], + ) + first_workbook = load_workbook(artifacts[0].path, read_only=True) + second_workbook = load_workbook(artifacts[1].path, read_only=True) + + self.assertEqual(first_workbook.sheetnames, ["data"]) + self.assertEqual(second_workbook.sheetnames, ["data"]) + self.assertEqual( + len(list(first_workbook["data"].iter_rows(values_only=True))), + 3, + ) + self.assertEqual( + len(list(second_workbook["data"].iter_rows(values_only=True))), + 2, + ) + self.assertEqual( + next(second_workbook["data"].iter_rows(values_only=True))[:4], + ("Наименование", "ИНН", "ОГРН", "КПП"), + ) + + selected_artifacts = [ + item + for item in generation.artifacts + if item.source_group == SourceGroup.PLANNED_INSPECTIONS.value + ] + self.assertEqual(len(selected_artifacts), 4) + + package = build_source_records_export_archive( + source_groups=[SourceGroup.PLANNED_INSPECTIONS.value], + export_format="xlsx", + ) + archive_bytes = b"".join(package.archive_chunks) + + with zipfile.ZipFile(BytesIO(archive_bytes)) as archive: + self.assertEqual( + archive.namelist(), + [ + "planned-inspections-part-001.xlsx", + "planned-inspections-part-002.xlsx", + ], + ) + self.assertEqual(package.files_count, 2) + self.assertFalse((Path(self.export_directory.name) / "tmp").exists()) + + def test_nightly_export_clears_model_ordering_to_avoid_multi_million_row_sort(self): + queryset = _source_group_queryset(SourceGroup.GOVERNMENT_PROCUREMENTS.value) + + self.assertFalse(queryset.ordered) + self.assertFalse(queryset.query.default_ordering) + + def test_json_artifact_reuses_valid_canonical_spool_without_copying_it(self): + organization = Organization.objects.create( + name='ООО "JSON без копии"', + inn="7707083813", + ) + extension = PlannedInspectionExtension.objects.create( + organization=organization, + title="Плановые проверки Генпрокуратуры России", + ) + for index in range(2): + OrganizationSourceRecord.objects.create( + extension=extension, + record_type="inspection", + source="inspections", + external_id=f"INSP-JSON-{index}", + title=f"Проверка {index}", + payload={"index": index}, + ) + + with TemporaryDirectory() as temporary_directory: + spool_path = Path(temporary_directory) / "rows.json" + artifact_path = Path(temporary_directory) / "inspections.json" + headers, records_count = _spool_source_group_rows( + source_group=SourceGroup.PLANNED_INSPECTIONS.value, + output_path=spool_path, + ) + + _render_source_group_artifact( + row_spool_path=spool_path, + output_path=artifact_path, + headers=headers, + file_format="json", + records_count=records_count, + ) + + self.assertEqual(records_count, 2) + self.assertEqual(len(json.loads(spool_path.read_text())), 2) + self.assertTrue(os.path.samefile(spool_path, artifact_path)) + + def test_generation_publishes_complete_matrix_and_keeps_previous_on_failure(self): + first_generation = build_source_record_export_artifacts() + current_generation = load_current_source_record_export_generation() + + self.assertEqual(first_generation.artifacts_count, 25) + self.assertEqual( + current_generation.generation_id, first_generation.generation_id + ) + self.assertEqual( + { + artifact.file_format + for artifact in current_generation.artifacts + if artifact.source_group == SourceGroup.FINANCIAL_INDICATORS.value + }, + {"json"}, + ) + self.assertEqual( + { + artifact.file_format + for artifact in current_generation.artifacts + if artifact.source_group == SourceGroup.PLANNED_INSPECTIONS.value + }, + {"csv", "xlsx", "json"}, + ) + + original_replace = os.replace + + def fail_generation_publish(source, destination): + if Path(source).name.startswith(".building-"): + raise OSError("disk full") + return original_replace(source, destination) + + with patch( + "organizations.source_record_export.os.replace", + side_effect=fail_generation_publish, + ), self.assertRaises(OSError): + build_source_record_export_artifacts() + + current_after_failure = load_current_source_record_export_generation() + self.assertEqual( + current_after_failure.generation_id, + first_generation.generation_id, + ) + def test_export_rejects_non_admin_user(self): self.client.force_authenticate(UserFactory.create_user()) @@ -187,8 +476,17 @@ class OrganizationSourceRecordExportApiV2Test(APITestCase): }, format="json", ) + ticket_response = self.client.post( + self.ticket_url, + { + "sources": [SourceGroup.PLANNED_INSPECTIONS.value], + "format": "json", + }, + format="json", + ) self.assertEqual(response.status_code, status.HTTP_403_FORBIDDEN) + self.assertEqual(ticket_response.status_code, status.HTTP_403_FORBIDDEN) def test_export_rejects_empty_duplicate_and_unknown_values(self): self.client.force_authenticate(UserFactory.create_superuser()) diff --git a/tests/apps/organizations/test_tasks.py b/tests/apps/organizations/test_tasks.py index 37f6d62..e3542a2 100644 --- a/tests/apps/organizations/test_tasks.py +++ b/tests/apps/organizations/test_tasks.py @@ -1,11 +1,13 @@ """Tests for organization source backfill tasks and schedules.""" from importlib import import_module +from tempfile import TemporaryDirectory from apps.parsers.models import ParserLoadLog from django.apps import apps as django_apps +from django.conf import settings from django.core.cache import cache -from django.test import TestCase +from django.test import TestCase, override_settings from django.utils import timezone from django_celery_beat.models import PeriodicTask from organizations.cache import get_organization_api_cache_version @@ -17,6 +19,7 @@ from organizations.models import ( from organizations.tasks import ( backfill_all_organization_sources, backfill_organization_sources_for_parser_batch, + refresh_source_record_export_artifacts, ) from tests.apps.parsers.factories import IndustrialCertificateRecordFactory @@ -116,3 +119,62 @@ class OrganizationSnapshotScheduleMigrationTest(TestCase): self.assertEqual(task.crontab.minute, "30") self.assertEqual(task.crontab.hour, "4") self.assertEqual(str(task.crontab.timezone), "Europe/Moscow") + + +class SourceRecordExportArtifactsTaskTest(TestCase): + """Checks nightly artifact generation and its distributed lock.""" + + def setUp(self): + cache.clear() + self.export_directory = TemporaryDirectory() + self.settings_override = override_settings( + SOURCE_RECORD_EXPORT_DIRECTORY=self.export_directory.name, + SOURCE_RECORD_EXPORT_GENERATIONS_TO_KEEP=2, + SOURCE_RECORD_EXPORT_LOCK_KEY="test:source-record-exports:lock", + SOURCE_RECORD_EXPORT_LOCK_TTL_SECONDS=300, + ) + self.settings_override.enable() + + def tearDown(self): + self.settings_override.disable() + self.export_directory.cleanup() + cache.clear() + super().tearDown() + + def test_refresh_task_builds_all_artifacts_and_releases_lock(self): + result = refresh_source_record_export_artifacts() + + self.assertEqual(result["status"], "success") + self.assertEqual(result["artifacts_count"], 25) + self.assertIsNone(cache.get(settings.SOURCE_RECORD_EXPORT_LOCK_KEY)) + + def test_refresh_task_skips_when_another_generation_holds_lock(self): + cache.set(settings.SOURCE_RECORD_EXPORT_LOCK_KEY, "busy", timeout=300) + + result = refresh_source_record_export_artifacts() + + self.assertEqual(result, {"status": "skipped", "reason": "locked"}) + + +class SourceRecordExportScheduleMigrationTest(TestCase): + """Checks the nightly Celery Beat schedule for export artifacts.""" + + def test_migration_seeds_nightly_source_record_export_task(self): + migration = import_module( + "organizations.migrations.0008_seed_nightly_source_record_exports" + ) + + migration.seed_nightly_source_record_export_schedule(django_apps, None) + migration.seed_nightly_source_record_export_schedule(django_apps, None) + + task = PeriodicTask.objects.get(name=migration.NIGHTLY_SOURCE_EXPORT_TASK_NAME) + self.assertEqual( + task.task, + "organizations.tasks.refresh_source_record_export_artifacts", + ) + self.assertTrue(task.enabled) + self.assertEqual(task.args, "[]") + self.assertEqual(task.kwargs, "{}") + self.assertEqual(task.crontab.minute, "30") + self.assertEqual(task.crontab.hour, "5") + self.assertEqual(str(task.crontab.timezone), "Europe/Moscow") diff --git a/tests/apps/parsers/test_source_registry.py b/tests/apps/parsers/test_source_registry.py index fd47794..355875b 100644 --- a/tests/apps/parsers/test_source_registry.py +++ b/tests/apps/parsers/test_source_registry.py @@ -1,3 +1,5 @@ +from dataclasses import asdict + from apps.parsers import tasks from apps.parsers.clients.common.structured import MAX_FILE_SIZE_BYTES from apps.parsers.source_registry import PARSER_SOURCES @@ -40,3 +42,14 @@ class ParserSourceRegistryFNSTest(SimpleTestCase): TASKS_BY_NAME[source.task_name], tasks.sync_fns_financial_reports, ) + + +class ParserSourceRegistryPresentationTest(SimpleTestCase): + def test_public_source_metadata_does_not_name_external_provider(self): + public_metadata = " ".join( + str(value) + for descriptor in PARSER_SOURCES.values() + for value in asdict(descriptor).values() + ) + + self.assertNotIn("checko", public_metadata.lower())