feat: add nightly source record exports
All checks were successful
All checks were successful
This commit is contained in:
@@ -89,6 +89,7 @@ services:
|
||||
volumes:
|
||||
- ./src:/app/src
|
||||
- ./logs:/app/logs
|
||||
- ./media:/app/media
|
||||
- ./input:/app/input
|
||||
command: ["celery-worker"]
|
||||
|
||||
|
||||
@@ -53,6 +53,7 @@ services:
|
||||
memswap_limit: 3g
|
||||
volumes:
|
||||
- ./logs:/app/logs
|
||||
- ./media:/app/media
|
||||
- ./input:/app/input
|
||||
command: ["celery-worker"]
|
||||
|
||||
|
||||
@@ -63,7 +63,10 @@ RUN mkdir -p logs media staticfiles input src/static \
|
||||
&& chown -R appuser:appgroup /app
|
||||
|
||||
ENV PATH="/app/.venv/bin:${PATH}" \
|
||||
PYTHONPATH=/app/src
|
||||
PYTHONPATH=/app/src \
|
||||
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
|
||||
|
||||
USER appuser
|
||||
ENTRYPOINT ["/app/docker/scripts/entrypoint.sh"]
|
||||
|
||||
78
docs/source-record-export-matrix-ru.md
Normal file
78
docs/source-record-export-matrix-ru.md
Normal file
@@ -0,0 +1,78 @@
|
||||
# Матрица файловых выгрузок внешних данных State Corp
|
||||
|
||||
## Пользовательский контракт
|
||||
|
||||
Администраторский frontend отправляет
|
||||
`POST /api/v2/organization-source-records/export-ticket/` с массивом `sources`
|
||||
и форматом. Backend проверяет последнее полностью опубликованное поколение и
|
||||
возвращает короткоживущий одноразовый ticket. Затем frontend передаёт ticket в
|
||||
теле обычной HTML-формы на
|
||||
`POST /api/v2/organization-source-records/export-download/`.
|
||||
|
||||
Браузер получает потоковый ZIP напрямую, без многогигабайтного `Blob` в
|
||||
JavaScript. Ticket не попадает в URL и после первого запроса становится
|
||||
недействительным. Совместимый администраторский endpoint
|
||||
`POST /api/v2/organization-source-records/export/` сразу возвращает тот же ZIP
|
||||
для API-клиентов.
|
||||
|
||||
Во время скачивания таблицы `external_data` не читаются: endpoint упаковывает
|
||||
готовые файлы последнего ночного поколения. При отсутствии поколения API
|
||||
возвращает `503` с кодом `source_export_not_ready`.
|
||||
|
||||
## Матрица
|
||||
|
||||
| Группа API | Таблицы State Corp | Файл | CSV | XLSX | JSON |
|
||||
|---|---|---|:---:|:---:|:---:|
|
||||
| `financial_indicators` | `FinancialReport`, `FinancialReportLine` | `financial-indicators` | — | — | да |
|
||||
| `government_procurements` | `PublicProcurement` | `public-procurements` | да | да | да |
|
||||
| `industrial_production` | `IndustrialProduct`, `IndustrialCertificate`, `ManufacturerRegistryEntry` | `manufacturers-and-products` | да | да | да |
|
||||
| `planned_inspections` | `ProsecutorCheck` | `planned-inspections` | да | да | да |
|
||||
| `bankruptcy` | `BankruptcyProcedure` | `bankruptcy-procedures` | да | да | да |
|
||||
| `defense_suppliers` | `DefenseUnreliableSupplier` | `defense-unreliable-suppliers` | да | да | да |
|
||||
| `arbitration` | `ArbitrationCase` | `arbitration-cases` | да | да | да |
|
||||
| `security_registries` | `InformationSecurityRegistryEntry` | `information-security-registries` | да | да | да |
|
||||
| `vacancies` | `LaborVacancy` | `labor-vacancies` | да | да | да |
|
||||
|
||||
Итого формируется 25 логических артефактов. Финансовые показатели всегда
|
||||
выгружаются в JSON с вложенным массивом `financial_lines`. Промышленная группа
|
||||
объединяет три таблицы, а поле `record_type` различает тип строки. Все строки
|
||||
содержат реквизиты организации, включая ОКПО.
|
||||
|
||||
Физических XLSX-файлов может быть больше: по умолчанию один файл содержит не
|
||||
более 100 000 строк данных и получает суффикс `-part-001`, `-part-002` и далее.
|
||||
|
||||
## Ночная генерация
|
||||
|
||||
Celery Beat запускает
|
||||
`apps.external_data.tasks.refresh_source_record_export_artifacts` ежедневно в
|
||||
`05:30 Europe/Moscow`.
|
||||
|
||||
Генератор:
|
||||
|
||||
1. читает каждую нормализованную таблицу один раз без model-level сортировки;
|
||||
2. создаёт компактный JSON-массив и переиспользует его как готовый JSON;
|
||||
3. потоково формирует CSV и write-only XLSX;
|
||||
4. атомарно публикует `current.json` только после готовности всей матрицы;
|
||||
5. при ошибке удаляет staging и продолжает отдавать предыдущее поколение;
|
||||
6. сохраняет текущее и предыдущее поколения по умолчанию.
|
||||
|
||||
Web и Celery worker должны использовать общий read-write volume `/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 распределённой блокировки |
|
||||
| `SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE` | `100000` | Строк данных в одной XLSX-части |
|
||||
| `SOURCE_RECORD_EXPORT_DOWNLOAD_TICKET_TTL_SECONDS` | `300` | Срок действия download-ticket |
|
||||
|
||||
Для атомарной генерации требуется свободное место не меньше
|
||||
`(GENERATIONS_TO_KEEP + 1) * размер поколения` плюс запас файловой системы.
|
||||
37
src/apps/external_data/export_serializers.py
Normal file
37
src/apps/external_data/export_serializers.py
Normal file
@@ -0,0 +1,37 @@
|
||||
"""Request serializers for prepared external-data exports."""
|
||||
|
||||
from apps.external_data.source_record_export import (
|
||||
EXPORT_FORMATS,
|
||||
SOURCE_GROUP_EXPORT_SPECS,
|
||||
)
|
||||
from rest_framework import serializers
|
||||
|
||||
|
||||
class SourceRecordExportRequestSerializer(serializers.Serializer):
|
||||
"""Validate selected source groups and the requested file format."""
|
||||
|
||||
sources = serializers.ListField(
|
||||
child=serializers.ChoiceField(
|
||||
choices=[
|
||||
(source_group, source_group)
|
||||
for source_group in SOURCE_GROUP_EXPORT_SPECS
|
||||
]
|
||||
),
|
||||
allow_empty=False,
|
||||
)
|
||||
format = serializers.ChoiceField(
|
||||
choices=[
|
||||
(export_format, export_format.upper()) for export_format in EXPORT_FORMATS
|
||||
]
|
||||
)
|
||||
|
||||
def validate_sources(self, value: list[str]) -> list[str]:
|
||||
if len(value) != len(set(value)):
|
||||
raise serializers.ValidationError("Источники не должны повторяться.")
|
||||
return value
|
||||
|
||||
|
||||
class SourceRecordExportDownloadSerializer(serializers.Serializer):
|
||||
"""Validate a one-time ticket submitted by a native browser form."""
|
||||
|
||||
ticket = serializers.CharField(max_length=64, trim_whitespace=False)
|
||||
22
src/apps/external_data/export_urls.py
Normal file
22
src/apps/external_data/export_urls.py
Normal file
@@ -0,0 +1,22 @@
|
||||
"""URL routes for prepared external-data exports."""
|
||||
|
||||
from apps.external_data.export_views import (
|
||||
SourceRecordExportDownloadView,
|
||||
SourceRecordExportTicketView,
|
||||
SourceRecordExportView,
|
||||
)
|
||||
from django.urls import path
|
||||
|
||||
app_name = "source_record_exports"
|
||||
|
||||
urlpatterns = [
|
||||
path("export/", SourceRecordExportView.as_view(), name="export"),
|
||||
path(
|
||||
"export-ticket/", SourceRecordExportTicketView.as_view(), name="export-ticket"
|
||||
),
|
||||
path(
|
||||
"export-download/",
|
||||
SourceRecordExportDownloadView.as_view(),
|
||||
name="export-download",
|
||||
),
|
||||
]
|
||||
174
src/apps/external_data/export_views.py
Normal file
174
src/apps/external_data/export_views.py
Normal file
@@ -0,0 +1,174 @@
|
||||
"""HTTP endpoints for prepared external-data exports."""
|
||||
|
||||
from apps.external_data.export_serializers import (
|
||||
SourceRecordExportDownloadSerializer,
|
||||
SourceRecordExportRequestSerializer,
|
||||
)
|
||||
from apps.external_data.source_record_export import (
|
||||
SourceRecordExportArchive,
|
||||
SourceRecordExportArtifactsUnavailable,
|
||||
SourceRecordExportTicketInvalid,
|
||||
build_source_records_export_archive,
|
||||
consume_source_record_export_download_ticket,
|
||||
create_source_record_export_download_ticket,
|
||||
)
|
||||
from django.http import StreamingHttpResponse
|
||||
from drf_yasg import openapi
|
||||
from drf_yasg.utils import swagger_auto_schema
|
||||
from rest_framework import status
|
||||
from rest_framework.permissions import AllowAny, IsAdminUser
|
||||
from rest_framework.request import Request
|
||||
from rest_framework.response import Response
|
||||
from rest_framework.views import APIView
|
||||
|
||||
|
||||
class SourceRecordExportResponseMixin:
|
||||
"""Build the shared streaming ZIP response for export endpoints."""
|
||||
|
||||
@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
|
||||
|
||||
@staticmethod
|
||||
def source_record_export_not_ready_response() -> Response:
|
||||
return Response(
|
||||
{
|
||||
"detail": "Готовая ночная выгрузка ещё не сформирована.",
|
||||
"code": "source_export_not_ready",
|
||||
},
|
||||
status=status.HTTP_503_SERVICE_UNAVAILABLE,
|
||||
headers={"Retry-After": "3600"},
|
||||
)
|
||||
|
||||
|
||||
class SourceRecordExportView(SourceRecordExportResponseMixin, APIView):
|
||||
"""Stream a request-specific ZIP from the current prepared generation."""
|
||||
|
||||
permission_classes = [IsAdminUser]
|
||||
|
||||
@swagger_auto_schema(
|
||||
operation_summary="Выгрузить записи внешних источников",
|
||||
operation_description=(
|
||||
"Упаковывает в ZIP готовые файлы выбранных групп без повторного "
|
||||
"чтения таблиц external_data. Финансовые показатели всегда JSON."
|
||||
),
|
||||
request_body=SourceRecordExportRequestSerializer,
|
||||
responses={
|
||||
200: openapi.Response(
|
||||
description="Потоковый ZIP-архив.",
|
||||
schema=openapi.Schema(type=openapi.TYPE_FILE),
|
||||
),
|
||||
400: "Некорректные параметры.",
|
||||
403: "Доступ разрешён только администраторам.",
|
||||
503: "Ночная выгрузка ещё не сформирована.",
|
||||
},
|
||||
tags=["Внешние данные"],
|
||||
)
|
||||
def post(self, request: Request) -> StreamingHttpResponse | Response:
|
||||
serializer = SourceRecordExportRequestSerializer(data=request.data)
|
||||
serializer.is_valid(raise_exception=True)
|
||||
try:
|
||||
package = build_source_records_export_archive(
|
||||
source_groups=serializer.validated_data["sources"],
|
||||
export_format=serializer.validated_data["format"],
|
||||
)
|
||||
except SourceRecordExportArtifactsUnavailable:
|
||||
return self.source_record_export_not_ready_response()
|
||||
return self.source_record_export_response(package)
|
||||
|
||||
|
||||
class SourceRecordExportTicketView(SourceRecordExportResponseMixin, APIView):
|
||||
"""Issue a short-lived ticket for a native browser download."""
|
||||
|
||||
permission_classes = [IsAdminUser]
|
||||
|
||||
@swagger_auto_schema(
|
||||
operation_summary="Подготовить нативное скачивание внешних данных",
|
||||
request_body=SourceRecordExportRequestSerializer,
|
||||
responses={
|
||||
201: "Одноразовый ticket и имя ZIP-файла.",
|
||||
400: "Некорректные параметры.",
|
||||
403: "Доступ разрешён только администраторам.",
|
||||
503: "Выгрузка или ticket временно недоступны.",
|
||||
},
|
||||
tags=["Внешние данные"],
|
||||
)
|
||||
def post(self, request: Request) -> Response:
|
||||
serializer = SourceRecordExportRequestSerializer(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 self.source_record_export_not_ready_response()
|
||||
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,
|
||||
)
|
||||
|
||||
|
||||
class SourceRecordExportDownloadView(SourceRecordExportResponseMixin, APIView):
|
||||
"""Consume a one-time ticket and stream the prepared ZIP archive."""
|
||||
|
||||
authentication_classes: list = []
|
||||
permission_classes = [AllowAny]
|
||||
|
||||
@swagger_auto_schema(
|
||||
operation_summary="Скачать готовые внешние данные по ticket",
|
||||
request_body=SourceRecordExportDownloadSerializer,
|
||||
responses={
|
||||
200: openapi.Response(
|
||||
description="Потоковый ZIP-архив.",
|
||||
schema=openapi.Schema(type=openapi.TYPE_FILE),
|
||||
),
|
||||
400: "Ticket отсутствует или имеет неверный формат.",
|
||||
410: "Ticket истёк или уже использован.",
|
||||
503: "Опубликованная выгрузка больше недоступна.",
|
||||
},
|
||||
tags=["Внешние данные"],
|
||||
)
|
||||
def post(self, request: Request) -> StreamingHttpResponse | Response:
|
||||
serializer = SourceRecordExportDownloadSerializer(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 self.source_record_export_not_ready_response()
|
||||
return self.source_record_export_response(package)
|
||||
1
src/apps/external_data/management/__init__.py
Normal file
1
src/apps/external_data/management/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""Management package for external-data operations."""
|
||||
1
src/apps/external_data/management/commands/__init__.py
Normal file
1
src/apps/external_data/management/commands/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""Management commands for external-data operations."""
|
||||
@@ -0,0 +1,32 @@
|
||||
"""Build the complete prepared external-data export matrix."""
|
||||
|
||||
import json
|
||||
|
||||
from apps.core.management.commands.base import BaseAppCommand
|
||||
from apps.external_data.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
|
||||
@@ -0,0 +1,61 @@
|
||||
import json
|
||||
|
||||
from django.db import migrations
|
||||
|
||||
NIGHTLY_SOURCE_EXPORT_TASK_NAME = "external-data:source-record-exports:nightly-msk"
|
||||
NIGHTLY_SOURCE_EXPORT_TASK_PATH = (
|
||||
"apps.external_data.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 State Corp external-data export 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"),
|
||||
("external_data", "0006_bankruptcy_procedure_status_length"),
|
||||
]
|
||||
|
||||
operations = [
|
||||
migrations.RunPython(
|
||||
seed_nightly_source_record_export_schedule,
|
||||
reverse_code=remove_nightly_source_record_export_schedule,
|
||||
),
|
||||
]
|
||||
1121
src/apps/external_data/source_record_export.py
Normal file
1121
src/apps/external_data/source_record_export.py
Normal file
File diff suppressed because it is too large
Load Diff
46
src/apps/external_data/tasks.py
Normal file
46
src/apps/external_data/tasks.py
Normal file
@@ -0,0 +1,46 @@
|
||||
"""Celery tasks for prepared external-data exports."""
|
||||
|
||||
import logging
|
||||
|
||||
from apps.core.tasks import PeriodicTask as CorePeriodicTask
|
||||
from apps.external_data.source_record_export import (
|
||||
build_source_record_export_artifacts,
|
||||
)
|
||||
from celery import shared_task
|
||||
from django.conf import settings
|
||||
from django.core.cache import cache
|
||||
|
||||
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 external-data export matrix."""
|
||||
|
||||
lock_key = getattr(
|
||||
settings,
|
||||
"SOURCE_RECORD_EXPORT_LOCK_KEY",
|
||||
"external-data: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)
|
||||
12
src/core/api_v2_urls.py
Normal file
12
src/core/api_v2_urls.py
Normal file
@@ -0,0 +1,12 @@
|
||||
"""API v2 routes shared with the Mostovik administrative frontend contract."""
|
||||
|
||||
from django.urls import include, path
|
||||
|
||||
app_name = "api_v2"
|
||||
|
||||
urlpatterns = [
|
||||
path(
|
||||
"organization-source-records/",
|
||||
include("apps.external_data.export_urls"),
|
||||
),
|
||||
]
|
||||
@@ -40,6 +40,7 @@ urlpatterns = [
|
||||
path("admin/", admin.site.urls),
|
||||
path("health/", include("apps.core.urls")),
|
||||
path("api/v1/", include("core.api_v1_urls", namespace="api_v1")),
|
||||
path("api/v2/", include("core.api_v2_urls", namespace="api_v2")),
|
||||
path("auth/", include("rest_framework.urls")),
|
||||
]
|
||||
|
||||
|
||||
@@ -231,6 +231,26 @@ STATICFILES_STORAGE = "whitenoise.storage.CompressedManifestStaticFilesStorage"
|
||||
|
||||
MEDIA_URL = "/media/"
|
||||
MEDIA_ROOT = PROJECT_ROOT / "media"
|
||||
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",
|
||||
"external-data: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")
|
||||
)
|
||||
|
||||
DEFAULT_AUTO_FIELD = "django.db.models.BigAutoField"
|
||||
AUTH_USER_MODEL = "user.User"
|
||||
|
||||
70
tests/apps/external_data/test_export_tasks.py
Normal file
70
tests/apps/external_data/test_export_tasks.py
Normal file
@@ -0,0 +1,70 @@
|
||||
"""Tests for the external-data export task and schedule."""
|
||||
|
||||
from importlib import import_module
|
||||
from tempfile import TemporaryDirectory
|
||||
|
||||
from apps.external_data.tasks import refresh_source_record_export_artifacts
|
||||
from django.apps import apps as django_apps
|
||||
from django.conf import settings
|
||||
from django.core.cache import cache
|
||||
from django.test import TestCase, override_settings
|
||||
from django_celery_beat.models import PeriodicTask
|
||||
|
||||
|
||||
class SourceRecordExportArtifactsTaskTest(TestCase):
|
||||
"""Check 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:state-corp-source-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_matrix_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_generation_lock_is_held(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):
|
||||
"""Check the nightly Celery Beat schedule for prepared exports."""
|
||||
|
||||
def test_migration_seeds_nightly_export_task_idempotently(self):
|
||||
migration = import_module(
|
||||
"apps.external_data.migrations.0007_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,
|
||||
"apps.external_data.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")
|
||||
225
tests/apps/external_data/test_source_record_export.py
Normal file
225
tests/apps/external_data/test_source_record_export.py
Normal file
@@ -0,0 +1,225 @@
|
||||
"""Tests for prepared State Corp external-data exports."""
|
||||
|
||||
import json
|
||||
import zipfile
|
||||
from io import BytesIO, StringIO
|
||||
from tempfile import TemporaryDirectory
|
||||
|
||||
from apps.external_data.source_record_export import (
|
||||
build_source_record_export_artifacts,
|
||||
load_current_source_record_export_generation,
|
||||
)
|
||||
from django.core.management import call_command
|
||||
from django.test import override_settings
|
||||
from openpyxl import load_workbook
|
||||
from rest_framework import status
|
||||
from rest_framework.test import APITestCase
|
||||
|
||||
from tests.apps.external_data.factories import (
|
||||
FinancialReportFactory,
|
||||
FinancialReportLineFactory,
|
||||
IndustrialCertificateFactory,
|
||||
IndustrialProductFactory,
|
||||
ManufacturerRegistryEntryFactory,
|
||||
ProsecutorCheckFactory,
|
||||
)
|
||||
from tests.apps.organization.factories import OrganizationFactory
|
||||
from tests.apps.user.factories import UserFactory
|
||||
|
||||
|
||||
class SourceRecordExportApiTest(APITestCase):
|
||||
"""Check admin access and zero-query delivery of prepared files."""
|
||||
|
||||
export_url = "/api/v2/organization-source-records/export/"
|
||||
ticket_url = "/api/v2/organization-source-records/export-ticket/"
|
||||
download_url = "/api/v2/organization-source-records/export-download/"
|
||||
|
||||
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()
|
||||
|
||||
def tearDown(self):
|
||||
self.settings_override.disable()
|
||||
self.export_directory.cleanup()
|
||||
super().tearDown()
|
||||
|
||||
@staticmethod
|
||||
def _response_body(response) -> bytes:
|
||||
return b"".join(response.streaming_content)
|
||||
|
||||
def test_export_is_unavailable_before_first_generation(self):
|
||||
self.client.force_authenticate(UserFactory.create_superuser())
|
||||
|
||||
response = self.client.post(
|
||||
self.export_url,
|
||||
{"sources": ["planned_inspections"], "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_generation_builds_full_matrix_from_normalized_tables(self):
|
||||
organization = OrganizationFactory.create(
|
||||
full_name='Акционерное общество "Экспорт"',
|
||||
okpo="12345678",
|
||||
)
|
||||
IndustrialProductFactory.create(organization=organization)
|
||||
IndustrialCertificateFactory.create(organization=organization)
|
||||
ManufacturerRegistryEntryFactory.create(organization=organization)
|
||||
ProsecutorCheckFactory.create(organization=organization)
|
||||
report = FinancialReportFactory.create(organization=organization)
|
||||
FinancialReportLineFactory.create(report=report, line_code="1600")
|
||||
|
||||
generation = build_source_record_export_artifacts()
|
||||
|
||||
self.assertEqual(generation.artifacts_count, 25)
|
||||
self.assertEqual(generation.files_count, 25)
|
||||
self.assertEqual(generation.records_count, 5)
|
||||
industrial_path = next(
|
||||
artifact.path
|
||||
for artifact in generation.artifacts
|
||||
if artifact.source_group == "industrial_production"
|
||||
and artifact.file_format == "json"
|
||||
)
|
||||
industrial_rows = json.loads(industrial_path.read_text(encoding="utf-8"))
|
||||
self.assertEqual(
|
||||
{row["record_type"] for row in industrial_rows},
|
||||
{
|
||||
"industrial_certificate",
|
||||
"industrial_product",
|
||||
"manufacturer_registry_entry",
|
||||
},
|
||||
)
|
||||
self.assertEqual({row["ОКПО"] for row in industrial_rows}, {"12345678"})
|
||||
|
||||
financial_path = next(
|
||||
artifact.path
|
||||
for artifact in generation.artifacts
|
||||
if artifact.source_group == "financial_indicators"
|
||||
)
|
||||
financial_rows = json.loads(financial_path.read_text(encoding="utf-8"))
|
||||
self.assertEqual(financial_rows[0]["financial_lines"][0]["line_code"], "1600")
|
||||
|
||||
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_streams_selected_prepared_files_without_database_queries(self):
|
||||
self.client.force_authenticate(UserFactory.create_superuser())
|
||||
organization = OrganizationFactory.create(okpo="87654321")
|
||||
ProsecutorCheckFactory.create(organization=organization)
|
||||
report = FinancialReportFactory.create(organization=organization)
|
||||
FinancialReportLineFactory.create(report=report)
|
||||
generation = build_source_record_export_artifacts()
|
||||
|
||||
with self.assertNumQueries(0):
|
||||
response = self.client.post(
|
||||
self.export_url,
|
||||
{
|
||||
"sources": ["planned_inspections", "financial_indicators"],
|
||||
"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.assertNotIn("Content-Length", response)
|
||||
|
||||
with zipfile.ZipFile(BytesIO(self._response_body(response))) as archive:
|
||||
self.assertEqual(
|
||||
set(archive.namelist()),
|
||||
{"planned-inspections.xlsx", "financial-indicators.json"},
|
||||
)
|
||||
workbook = load_workbook(
|
||||
BytesIO(archive.read("planned-inspections.xlsx")),
|
||||
read_only=True,
|
||||
)
|
||||
rows = list(workbook["data"].iter_rows(values_only=True))
|
||||
self.assertEqual(
|
||||
rows[0][:6],
|
||||
("Наименование", "ИНН", "ОГРН", "КПП", "ОКПО", "organization"),
|
||||
)
|
||||
self.assertEqual(rows[1][4], "87654321")
|
||||
|
||||
def test_ticket_is_admin_only_and_can_be_consumed_once_without_auth(self):
|
||||
build_source_record_export_artifacts()
|
||||
regular_user = UserFactory.create_user()
|
||||
self.client.force_authenticate(regular_user)
|
||||
forbidden_response = self.client.post(
|
||||
self.ticket_url,
|
||||
{"sources": ["bankruptcy"], "format": "json"},
|
||||
format="json",
|
||||
)
|
||||
self.assertEqual(forbidden_response.status_code, status.HTTP_403_FORBIDDEN)
|
||||
|
||||
self.client.force_authenticate(UserFactory.create_superuser())
|
||||
with self.assertNumQueries(0):
|
||||
ticket_response = self.client.post(
|
||||
self.ticket_url,
|
||||
{"sources": ["bankruptcy"], "format": "json"},
|
||||
format="json",
|
||||
)
|
||||
self.assertEqual(ticket_response.status_code, status.HTTP_201_CREATED)
|
||||
self.assertRegex(ticket_response.data["ticket"], r"^[A-Za-z0-9_-]{43}$")
|
||||
self.assertEqual(ticket_response.data["expires_in"], 300)
|
||||
|
||||
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)
|
||||
with zipfile.ZipFile(
|
||||
BytesIO(self._response_body(download_response))
|
||||
) as archive:
|
||||
self.assertEqual(archive.namelist(), ["bankruptcy-procedures.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")
|
||||
|
||||
@override_settings(SOURCE_RECORD_EXPORT_XLSX_ROWS_PER_FILE=2)
|
||||
def test_xlsx_is_split_into_bounded_workbook_parts(self):
|
||||
organization = OrganizationFactory.create()
|
||||
ProsecutorCheckFactory.create_batch(3, organization=organization)
|
||||
|
||||
generation = build_source_record_export_artifacts()
|
||||
|
||||
inspection_parts = [
|
||||
artifact
|
||||
for artifact in generation.artifacts
|
||||
if artifact.source_group == "planned_inspections"
|
||||
and artifact.file_format == "xlsx"
|
||||
]
|
||||
self.assertEqual(
|
||||
[artifact.part_number for artifact in inspection_parts], [1, 2]
|
||||
)
|
||||
self.assertEqual(
|
||||
[artifact.file_name for artifact in inspection_parts],
|
||||
[
|
||||
"planned-inspections-part-001.xlsx",
|
||||
"planned-inspections-part-002.xlsx",
|
||||
],
|
||||
)
|
||||
Reference in New Issue
Block a user