diff --git a/src/apps/parsers/clients/checko/client.py b/src/apps/parsers/clients/checko/client.py index 3f25cf5..d51833b 100644 --- a/src/apps/parsers/clients/checko/client.py +++ b/src/apps/parsers/clients/checko/client.py @@ -193,7 +193,18 @@ RU_FIELD_MAP = { "ЕФРСБ": "bankruptcy", "НомерДела": "case_number", # РНП - "НедобПост": "unfair_supplier", + "НедобПост": "is_unfair_supplier", + "НедобПостЗап": "unfair_supplier", + "РеестрНомер": "registry_number", + "ДатаПуб": "publish_date", + "ДатаУтв": "approval_date", + "ЗаказНаимСокр": "customer_short_name", + "ЗаказНаимПолн": "customer_full_name", + "ЗаказИНН": "customer_inn", + "ЗаказКПП": "customer_kpp", + "ЗакупНомер": "purchase_number", + "ЗакупОпис": "purchase_description", + "ЦенаКонтр": "contract_price", # Численность "СЧР": "employees_count", # Налоги diff --git a/src/apps/parsers/clients/common/__init__.py b/src/apps/parsers/clients/common/__init__.py index f64c3a6..ba255ae 100644 --- a/src/apps/parsers/clients/common/__init__.py +++ b/src/apps/parsers/clients/common/__init__.py @@ -1,5 +1,6 @@ """Общие клиенты и DTO для новых разнородных источников.""" +from apps.parsers.clients.common.eis_registry import EisRegistryProcurementClient from apps.parsers.clients.common.schemas import GenericParserItem from apps.parsers.clients.common.structured import ( StructuredDataClient, @@ -8,6 +9,7 @@ from apps.parsers.clients.common.structured import ( __all__ = [ "GenericParserItem", + "EisRegistryProcurementClient", "StructuredDataClient", "StructuredDataClientError", ] diff --git a/src/apps/parsers/clients/common/eis_registry.py b/src/apps/parsers/clients/common/eis_registry.py new file mode 100644 index 0000000..82746a6 --- /dev/null +++ b/src/apps/parsers/clients/common/eis_registry.py @@ -0,0 +1,142 @@ +"""Addressed EIS procurement lookup for organizations from the OПК registry.""" + +from __future__ import annotations + +import re +from dataclasses import dataclass, field, replace +from datetime import date +from urllib.parse import urlencode + +from apps.parsers.clients.common.schemas import GenericParserItem +from apps.parsers.clients.common.structured import StructuredDataClient +from bs4 import BeautifulSoup +from dateutil.relativedelta import relativedelta + +EIS_PROCUREMENT_SEARCH_URL = ( + "https://zakupki.gov.ru/epz/order/extendedsearch/results.html" +) +EIS_SEARCH_PAGE_MAX_SIZE_BYTES = 4 * 1024 * 1024 +EIS_RECORDS_PER_PAGE = 50 + + +def _digits(value: str) -> str: + return re.sub(r"\D+", "", str(value or "")) + + +@dataclass +class EisRegistryProcurementClient: + """Fetch all recent EIS notices for one exact customer INN and law.""" + + source: str + law: str + proxies: list[str] | None = None + timeout: int = 120 + lookback_months: int = 12 + max_pages: int = 1000 + _structured_client: StructuredDataClient | None = field( + default=None, + repr=False, + ) + + @property + def structured_client(self) -> StructuredDataClient: + if self._structured_client is None: + self._structured_client = StructuredDataClient( + source=self.source, + proxies=self.proxies, + timeout=self.timeout, + ) + return self._structured_client + + def fetch_for_customer( + self, + *, + inn: str, + organization_name: str, + organization_id: str, + today: date | None = None, + ) -> list[GenericParserItem]: + """Fetch every result page and keep only cards with the exact customer INN.""" + customer_inn = _digits(inn) + if not customer_inn: + return [] + + period_end = today or date.today() + period_start = period_end - relativedelta(months=self.lookback_months) + records_by_external_id: dict[str, GenericParserItem] = {} + page_number = 1 + + while page_number <= self.max_pages: + url = self._build_search_url( + inn=customer_inn, + page_number=page_number, + period_start=period_start, + period_end=period_end, + ) + content = self.structured_client.http_client.download_file( + url, + max_size_bytes=EIS_SEARCH_PAGE_MAX_SIZE_BYTES, + ) + page_records = self.structured_client.fetch_records( + content=content, + file_name="results.html", + ) + for record in page_records: + if _digits(record.inn) != customer_inn: + continue + payload = dict(record.payload) + payload.update( + { + "provider": "eis", + "organization_id": organization_id, + "customer_inn": customer_inn, + "customer_name": organization_name, + "law": self.law, + "lookback_months": self.lookback_months, + } + ) + records_by_external_id[record.external_id] = replace( + record, + inn=customer_inn, + organisation_name=organization_name, + payload=payload, + ) + + last_page = self._last_page_number(content) + if page_number >= last_page: + break + page_number += 1 + + return list(records_by_external_id.values()) + + def _build_search_url( + self, + *, + inn: str, + page_number: int, + period_start: date, + period_end: date, + ) -> str: + law_param = "fz44" if self.law == "44" else "fz223" + params = { + law_param: "on", + "searchString": inn, + "strictEqual": "true", + "publishDateFrom": period_start.strftime("%d.%m.%Y"), + "publishDateTo": period_end.strftime("%d.%m.%Y"), + "sortBy": "UPDATE_DATE", + "pageNumber": str(page_number), + "sortDirection": "false", + "recordsPerPage": f"_{EIS_RECORDS_PER_PAGE}", + } + return f"{EIS_PROCUREMENT_SEARCH_URL}?{urlencode(params)}" + + @staticmethod + def _last_page_number(content: bytes) -> int: + soup = BeautifulSoup(content, "html.parser") + page_numbers = [1] + for node in soup.select("[data-pagenumber]"): + value = str(node.get("data-pagenumber") or "") + if value.isdigit(): + page_numbers.append(int(value)) + return max(page_numbers) diff --git a/src/apps/parsers/migrations/0028_registry_procurement_claims.py b/src/apps/parsers/migrations/0028_registry_procurement_claims.py new file mode 100644 index 0000000..683e582 --- /dev/null +++ b/src/apps/parsers/migrations/0028_registry_procurement_claims.py @@ -0,0 +1,43 @@ +from django.db import migrations, models + +LEGACY_WEEKLY_TASK_NAMES = [ + "parser:procurements_44fz:weekly-saturday-msk", + "parser:procurements_223fz:weekly-saturday-msk", + "parser:contracts:weekly-saturday-msk", + "parser:unfair_suppliers:weekly-saturday-msk", +] + + +def disable_legacy_weekly_tasks(apps, schema_editor): + PeriodicTask = apps.get_model("django_celery_beat", "PeriodicTask") + PeriodicTask.objects.filter(name__in=LEGACY_WEEKLY_TASK_NAMES).update(enabled=False) + + +class Migration(migrations.Migration): + dependencies = [ + ("parsers", "0027_checkocollectionattempt"), + ] + + operations = [ + migrations.AlterField( + model_name="checkocollectionattempt", + name="source", + field=models.CharField( + choices=[ + ("arbitration", "Арбитражные дела"), + ("bankruptcy", "Банкротства"), + ("contracts", "Контракты"), + ("inspections", "Проверки"), + ("procurements_44fz", "Закупки 44-ФЗ"), + ("procurements_223fz", "Закупки 223-ФЗ"), + ("unfair_suppliers", "Недобросовестные поставщики"), + ], + max_length=32, + verbose_name="источник", + ), + ), + migrations.RunPython( + disable_legacy_weekly_tasks, + reverse_code=migrations.RunPython.noop, + ), + ] diff --git a/src/apps/parsers/models.py b/src/apps/parsers/models.py index 771f669..b2763d1 100644 --- a/src/apps/parsers/models.py +++ b/src/apps/parsers/models.py @@ -125,6 +125,9 @@ class CheckoCollectionAttempt(TimestampMixin, models.Model): BANKRUPTCY = "bankruptcy", _("Банкротства") CONTRACTS = "contracts", _("Контракты") INSPECTIONS = "inspections", _("Проверки") + PROCUREMENTS_44FZ = "procurements_44fz", _("Закупки 44-ФЗ") + PROCUREMENTS_223FZ = "procurements_223fz", _("Закупки 223-ФЗ") + UNFAIR_SUPPLIERS = "unfair_suppliers", _("Недобросовестные поставщики") class Status(models.TextChoices): IN_PROGRESS = "in_progress", _("В процессе") diff --git a/src/apps/parsers/source_cards.py b/src/apps/parsers/source_cards.py index 93eabd0..29a7ac6 100644 --- a/src/apps/parsers/source_cards.py +++ b/src/apps/parsers/source_cards.py @@ -130,18 +130,16 @@ SOURCE_CARD_DEFINITIONS: tuple[SourceCardDefinition, ...] = ( description="Данные ЕИС закупок по тендерам и заказчикам.", order=20, task_names=( - "apps.parsers.tasks.parse_procurements", - "apps.parsers.tasks.sync_procurements", "apps.parsers.tasks.parse_procurements_44fz", "apps.parsers.tasks.parse_procurements_223fz", - "apps.parsers.tasks.parse_contracts", "apps.parsers.tasks.parse_registry_contracts", + "apps.parsers.tasks.parse_registry_enrichment_sources", ), source_items=( SourceItemDefinition( code="procurements", - title="Единая информационная система закупок", - description=("Закупки и связанные данные из ЕИС по 44-ФЗ и 223-ФЗ."), + title="Общие закупки", + description="Объединение закупок 44-ФЗ и 223-ФЗ без дублирования.", parser_source=ParserLoadLog.Source.PROCUREMENTS, ), SourceItemDefinition( @@ -168,7 +166,7 @@ SOURCE_CARD_DEFINITIONS: tuple[SourceCardDefinition, ...] = ( name="region_code", label="Код региона", description="Код региона ЕИС, например 77 для Москвы.", - required=True, + required=False, ), RefreshParamDefinition( name="law_type", @@ -368,6 +366,10 @@ SOURCE_RECORD_SOURCES_BY_ITEM_CODE = { for item in definition.source_items if item.parser_source } +SOURCE_RECORD_SOURCES_BY_ITEM_CODE["procurements"] = [ + ParserLoadLog.Source.PROCUREMENTS_44FZ, + ParserLoadLog.Source.PROCUREMENTS_223FZ, +] class SourceCardService: @@ -455,7 +457,14 @@ class SourceCardService: source_items = [ cls._build_source_item(item, context) for item in definition.source_items ] - records_count = sum(item["records_count"] for item in source_items) + if definition.slug == "public-procurements": + items_by_code = {item["code"]: item for item in source_items} + records_count = ( + items_by_code["procurements"]["records_count"] + + items_by_code["contracts"]["records_count"] + ) + else: + records_count = sum(item["records_count"] for item in source_items) organizations_count = cls._get_card_organizations_count( definition, source_items, context ) @@ -998,27 +1007,18 @@ class SourceCardService: return [task_info] if definition.slug == "public-procurements": - from apps.parsers.tasks import sync_procurements + from apps.parsers.tasks import parse_registry_enrichment_sources task_info = cls._enqueue_task( - task=sync_procurements, - task_name="apps.parsers.tasks.sync_procurements", + task=parse_registry_enrichment_sources, + task_name="apps.parsers.tasks.parse_registry_enrichment_sources", requested_by_id=requested_by_id, meta={ "source_card": definition.slug, - "source": ParserLoadLog.Source.PROCUREMENTS, - "region_code": params["region_code"], - "law_type": params.get("law_type", "44"), + "source": ParserLoadLog.Source.PROCUREMENTS_44FZ, }, kwargs={ "requested_by_id": requested_by_id, - "region_code": params["region_code"], - "law_type": params.get("law_type", "44"), - **{ - key: value - for key, value in params.items() - if key in {"current_year", "current_month"} - }, }, ) return [task_info] diff --git a/src/apps/parsers/tasks.py b/src/apps/parsers/tasks.py index 4c1dda9..de9ee60 100644 --- a/src/apps/parsers/tasks.py +++ b/src/apps/parsers/tasks.py @@ -37,7 +37,11 @@ from apps.parsers.clients.checko import ( SearchType, ) from apps.parsers.clients.checko.exceptions import CheckoError -from apps.parsers.clients.common import GenericParserItem, StructuredDataClient +from apps.parsers.clients.common import ( + EisRegistryProcurementClient, + GenericParserItem, + StructuredDataClient, +) from apps.parsers.clients.fns import FNSApiClient, FNSApiReport from apps.parsers.clients.gisp import GispProductsClient from apps.parsers.clients.minpromtorg import ( @@ -82,6 +86,7 @@ FEDRESURS_CHECKO_FALLBACK_LIMIT = 100 ARBITRATION_CHECKO_LIMIT = 100 REGISTRY_INSPECTIONS_CHECKO_LIMIT = 1000 REGISTRY_CONTRACTS_CHECKO_LIMIT = 1000 +REGISTRY_ENRICHMENT_BATCH_SIZE = 250 FSTEC_CHECKO_IDENTITY_LOOKUP_LIMIT = 1000 PARSER_STALE_LOAD_MAX_AGE_MINUTES = 90 PARSER_SOFT_TIME_LIMIT_SECONDS = 15 * 60 @@ -177,9 +182,15 @@ def _active_registry_lookup_targets( *, limit: int | None = None, checko_source: str | None = None, + require_inn: bool = False, + organization_ids: list[str] | None = None, ) -> list[RegistryLookupTarget]: """Вернуть организации, которые сейчас состоят хотя бы в одном реестре.""" queryset = SourceOrganization.objects.filter(opk_registry_membership=True) + if organization_ids is not None: + queryset = queryset.filter(uid__in=organization_ids) + if require_inn: + queryset = queryset.exclude(inn="") if checko_source is not None: queryset = queryset.exclude( checko_collection_attempts__source=checko_source, @@ -215,6 +226,17 @@ def _active_registry_lookup_targets( return targets +def _resolve_registry_enrichment_limit(limit: int | None) -> int: + return _resolve_lookup_limit( + limit, + default=getattr( + settings, + "REGISTRY_ENRICHMENT_BATCH_SIZE", + REGISTRY_ENRICHMENT_BATCH_SIZE, + ), + ) + + def _resolve_proxies(proxies: list[str] | None) -> list[str] | None: """ Разрешить итоговый список прокси. @@ -1377,6 +1399,226 @@ def _fetch_checko_registry_inspections( return records +def _fetch_eis_registry_procurement_records( + *, + source: str, + law: str, + limit: int | None, + proxies: list[str] | None, + organization_ids: list[str] | None = None, +) -> list[GenericParserItem]: + """Fetch addressed EIS notices for the next monthly OПК registry batch.""" + resolved_limit = _resolve_registry_enrichment_limit(limit) + if resolved_limit <= 0: + return [] + + claim_source = { + "44": CheckoCollectionAttempt.Source.PROCUREMENTS_44FZ, + "223": CheckoCollectionAttempt.Source.PROCUREMENTS_223FZ, + }[law] + targets = _active_registry_lookup_targets( + limit=resolved_limit, + checko_source=claim_source, + require_inn=True, + organization_ids=organization_ids, + ) + if not targets: + raise ParserSourceSkipped( + f"no OПК registry organizations are due for monthly {law}-FZ lookup" + ) + + client = EisRegistryProcurementClient( + source=source, + law=law, + proxies=proxies, + ) + records: list[GenericParserItem] = [] + failed_lookups = 0 + for target in targets: + attempt = claim_monthly_collection( + organization_id=target.organization_id, + source=claim_source, + ) + if attempt is None: + continue + records_before = len(records) + try: + records.extend( + client.fetch_for_customer( + inn=target.inn, + organization_name=target.name, + organization_id=target.organization_id, + ) + ) + except HTTPClientError as exc: + failed_lookups += 1 + finish_collection(attempt, records_count=0, error=exc) + logger.info( + "EIS %s-FZ lookup failed for customer_inn=%s: %s", + law, + target.inn, + exc, + ) + else: + finish_collection( + attempt, + records_count=len(records) - records_before, + ) + + if failed_lookups == len(targets) and not records: + raise ParserSourceSkipped(f"EIS {law}-FZ lookups failed for all targets") + return records + + +def _checko_unfair_supplier_items( + *, + company, + target: RegistryLookupTarget, +) -> list[GenericParserItem]: + """Convert Checko НедобПостЗап 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: + return [] + if not target.inn and target.ogrn and company_ogrn != target.ogrn: + return [] + + records = [] + for item in getattr(company, "unfair_supplier", ()): + registry_number = str(getattr(item, "registry_number", "") or "").strip() + purchase_number = str(getattr(item, "purchase_number", "") or "").strip() + stable_id = registry_number or purchase_number + if not stable_id: + stable_id = hashlib.sha256( + repr(item).encode("utf-8", errors="replace") + ).hexdigest() + customer = { + "short_name": getattr(item, "customer_short_name", None), + "full_name": getattr(item, "customer_full_name", None), + "inn": getattr(item, "customer_inn", None), + "kpp": getattr(item, "customer_kpp", None), + } + contract_price = getattr(item, "contract_price", None) + records.append( + GenericParserItem( + source=ParserLoadLog.Source.UNFAIR_SUPPLIERS, + external_id=f"checko-unfair-supplier:{stable_id}", + inn=target.inn, + ogrn=target.ogrn, + organisation_name=target.name, + title=str( + getattr(item, "purchase_description", None) + or "Запись реестра недобросовестных поставщиков" + ), + record_date=str( + getattr(item, "publish_date", None) + or getattr(item, "approval_date", None) + or "" + ), + amount=( + Decimal(str(contract_price)) if contract_price is not None else None + ), + status="included", + payload={ + "provider": "checko", + "organization_id": target.organization_id, + "registry_number": registry_number, + "publish_date": getattr(item, "publish_date", None), + "approval_date": getattr(item, "approval_date", None), + "purchase_number": purchase_number, + "purchase_description": getattr( + item, + "purchase_description", + None, + ), + "contract_price": contract_price, + "supplier": { + "inn": target.inn, + "ogrn": target.ogrn, + "name": target.name, + }, + "customer": customer, + }, + ) + ) + return records + + +def _fetch_checko_unfair_supplier_records( # noqa: C901 + *, + limit: int | None, + proxies: list[str] | None, + organization_ids: list[str] | None = None, +) -> list[GenericParserItem]: + """Fetch RNP entries through one Checko /company 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") + + resolved_limit = _resolve_registry_enrichment_limit(limit) + if resolved_limit <= 0: + return [] + targets = _active_registry_lookup_targets( + limit=resolved_limit, + checko_source=CheckoCollectionAttempt.Source.UNFAIR_SUPPLIERS, + organization_ids=organization_ids, + ) + if not targets: + raise ParserSourceSkipped( + "no OПК registry organizations are due for monthly RNP lookup" + ) + + checko_proxies = ( + proxies if getattr(settings, "CHECKO_USE_RUNTIME_PROXIES", False) else None + ) + client = CheckoClient(api_key=api_key, proxies=checko_proxies, timeout=30) + records: list[GenericParserItem] = [] + failed_lookups = 0 + rate_limited = False + for target in targets: + attempt = claim_monthly_collection( + organization_id=target.organization_id, + source=CheckoCollectionAttempt.Source.UNFAIR_SUPPLIERS, + ) + if attempt is None: + continue + records_before = len(records) + try: + response = client.get_company( + CompanyRequest(inn=target.inn or None, ogrn=target.ogrn or None) + ) + except CheckoRateLimitError as exc: + finish_collection(attempt, records_count=0, error=exc) + failed_lookups += 1 + rate_limited = True + logger.warning("Checko 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", + target.inn or target.ogrn, + exc, + ) + else: + if response.data is not None: + records.extend( + _checko_unfair_supplier_items( + company=response.data, + target=target, + ) + ) + finish_collection( + attempt, + records_count=len(records) - records_before, + ) + + if not records and (rate_limited or failed_lookups == len(targets)): + raise ParserSourceSkipped("Checko RNP lookups failed for all targets") + return records + + def _contract_party_payload(party) -> dict: if party is None: return {} @@ -2320,26 +2562,13 @@ def parse_all_sources( generic_results = {} if not getattr(settings, "CELERY_TASK_ALWAYS_EAGER", False): generic_results = { - "procurements_44fz": parse_procurements_44fz.delay(proxies=proxies).id, - "procurements_223fz": parse_procurements_223fz.delay( - proxies=proxies + "registry_enrichment": parse_registry_enrichment_sources.delay( + limit=_resolve_registry_enrichment_limit(None), + proxies=proxies, ).id, - "contracts": parse_contracts.delay(proxies=proxies).id, - "registry_contracts": parse_registry_contracts.delay( - proxies=proxies - ).id, - "unfair_suppliers": parse_unfair_suppliers.delay(proxies=proxies).id, "fas_goz": parse_fas_goz_evasion.delay(proxies=proxies).id, "fns_financial": sync_fns_financial_reports.delay(proxies=proxies).id, - "arbitration": parse_arbitration_cases.delay(proxies=proxies).id, - "fedresurs_bankruptcy": parse_fedresurs_bankruptcy.delay( - proxies=proxies - ).id, - "registry_inspections": parse_registry_inspections.delay( - proxies=proxies - ).id, "fstec": parse_fstec_registers.delay(proxies=proxies).id, - "trudvsem": parse_trudvsem_vacancies.delay(proxies=proxies).id, } results = { @@ -2992,23 +3221,40 @@ def parse_procurements_44fz( *, file_url: str | None = None, file_path: str | None = None, + limit: int | None = None, proxies: list[str] | None = None, + organization_ids: list[str] | None = None, requested_by_id: int | None = None, ) -> dict: - """Парсинг официальной выдачи ЕИС 44-ФЗ в GenericParserRecord.""" + """Addressed OПК lookup by default; explicit files remain a manual tool.""" proxies = _resolve_proxies(proxies) + if file_url or file_path: + + def fetch_records(): + return _fetch_structured_records( + source_key="procurements_44fz", + file_url=file_url, + file_path=file_path, + proxies=proxies, + ) + else: + + def fetch_records(): + return _fetch_eis_registry_procurement_records( + source=ParserLoadLog.Source.PROCUREMENTS_44FZ, + law="44", + limit=limit, + proxies=proxies, + organization_ids=organization_ids, + ) + return _run_generic_parser( self, source_key="procurements_44fz", source=ParserLoadLog.Source.PROCUREMENTS_44FZ, task_name="apps.parsers.tasks.parse_procurements_44fz", requested_by_id=requested_by_id, - fetch_records=lambda: _fetch_structured_records( - source_key="procurements_44fz", - file_url=file_url, - file_path=file_path, - proxies=proxies, - ), + fetch_records=fetch_records, ) @@ -3018,23 +3264,40 @@ def parse_procurements_223fz( *, file_url: str | None = None, file_path: str | None = None, + limit: int | None = None, proxies: list[str] | None = None, + organization_ids: list[str] | None = None, requested_by_id: int | None = None, ) -> dict: - """Парсинг официальной выдачи ЕИС 223-ФЗ в GenericParserRecord.""" + """Addressed OПК lookup by default; explicit files remain a manual tool.""" proxies = _resolve_proxies(proxies) + if file_url or file_path: + + def fetch_records(): + return _fetch_structured_records( + source_key="procurements_223fz", + file_url=file_url, + file_path=file_path, + proxies=proxies, + ) + else: + + def fetch_records(): + return _fetch_eis_registry_procurement_records( + source=ParserLoadLog.Source.PROCUREMENTS_223FZ, + law="223", + limit=limit, + proxies=proxies, + organization_ids=organization_ids, + ) + return _run_generic_parser( self, source_key="procurements_223fz", source=ParserLoadLog.Source.PROCUREMENTS_223FZ, task_name="apps.parsers.tasks.parse_procurements_223fz", requested_by_id=requested_by_id, - fetch_records=lambda: _fetch_structured_records( - source_key="procurements_223fz", - file_url=file_url, - file_path=file_path, - proxies=proxies, - ), + fetch_records=fetch_records, ) @@ -3093,23 +3356,38 @@ def parse_unfair_suppliers( *, file_url: str | None = None, file_path: str | None = None, + limit: int | None = None, proxies: list[str] | None = None, + organization_ids: list[str] | None = None, requested_by_id: int | None = None, ) -> dict: - """Парсинг реестра недобросовестных поставщиков.""" + """Checko RNP lookup by default; explicit files remain a manual tool.""" proxies = _resolve_proxies(proxies) + if file_url or file_path: + + def fetch_records(): + return _fetch_structured_records( + source_key="unfair_suppliers", + file_url=file_url, + file_path=file_path, + proxies=proxies, + ) + else: + + def fetch_records(): + return _fetch_checko_unfair_supplier_records( + limit=limit, + proxies=proxies, + organization_ids=organization_ids, + ) + return _run_generic_parser( self, source_key="unfair_suppliers", source=ParserLoadLog.Source.UNFAIR_SUPPLIERS, task_name="apps.parsers.tasks.parse_unfair_suppliers", requested_by_id=requested_by_id, - fetch_records=lambda: _fetch_structured_records( - source_key="unfair_suppliers", - file_url=file_url, - file_path=file_path, - proxies=proxies, - ), + fetch_records=fetch_records, ) @@ -3241,11 +3519,13 @@ def parse_registry_enrichment_sources( """ Запустить daily-контур обогащения активных организаций из реестров. - Внутри остаются независимые задачи: одни забирают полный официальный реестр, - другие делают lookup по ИНН/ОГРН активных организаций. + Each addressed source advances through the next monthly registry batch. """ proxies = _resolve_proxies(proxies) + resolved_limit = _resolve_registry_enrichment_limit(limit) tasks_to_run = { + "procurements_44fz": parse_procurements_44fz, + "procurements_223fz": parse_procurements_223fz, "contracts": parse_registry_contracts, "unfair_suppliers": parse_unfair_suppliers, "arbitration": parse_arbitration_cases, @@ -3259,8 +3539,15 @@ def parse_registry_enrichment_sources( "proxies": proxies, "requested_by_id": requested_by_id, } - if key in {"contracts", "arbitration", "inspections"}: - kwargs["limit"] = limit + if key in { + "procurements_44fz", + "procurements_223fz", + "contracts", + "unfair_suppliers", + "arbitration", + "inspections", + }: + kwargs["limit"] = resolved_limit result = task.delay(**kwargs) results[key] = result.id return results diff --git a/tests/apps/parsers/test_checko_parsers.py b/tests/apps/parsers/test_checko_parsers.py index 3e215c3..bf01f94 100644 --- a/tests/apps/parsers/test_checko_parsers.py +++ b/tests/apps/parsers/test_checko_parsers.py @@ -37,6 +37,58 @@ class CheckoClientParsingTest(SimpleTestCase): data["\u0417\u0430\u043f\u0438\u0441\u0438"][0]["\u041e\u0413\u0420\u041d"], ) + def test_parse_company_maps_russian_unfair_supplier_records(self): + data = _map_ru_keys( + { + "ОГРН": "1027700000001", + "ИНН": "7701000001", + "НаимСокр": "ООО ОПК", + "НедобПост": True, + "НедобПостЗап": [ + { + "РеестрНомер": "RNP-1", + "ДатаПуб": "2026-07-01", + "ДатаУтв": "2026-06-30", + "ЗаказНаимСокр": "Заказчик", + "ЗаказНаимПолн": "Заказчик полный", + "ЗаказИНН": "7702000002", + "ЗаказКПП": "770201001", + "ЗакупНомер": "PURCHASE-1", + "ЗакупОпис": "Поставка оборудования", + "ЦенаКонтр": 1500000, + }, + { + "РеестрНомер": "RNP-2", + "ЗакупНомер": "PURCHASE-2", + }, + ], + } + ) + + company = self.client._parse_company_data(data) + + self.assertEqual(len(company.unfair_supplier), 2) + record = company.unfair_supplier[0] + self.assertEqual(record.registry_number, "RNP-1") + self.assertEqual(record.customer_inn, "7702000002") + self.assertEqual(record.purchase_description, "Поставка оборудования") + self.assertEqual(record.contract_price, 1500000) + self.assertEqual(company.unfair_supplier[1].registry_number, "RNP-2") + + def test_parse_company_handles_empty_unfair_supplier_records(self): + company = self.client._parse_company_data( + _map_ru_keys( + { + "ОГРН": "1027700000001", + "ИНН": "7701000001", + "НедобПост": False, + "НедобПостЗап": [], + } + ) + ) + + self.assertEqual(company.unfair_supplier, ()) + def test_parse_okved_info_with_additional(self): info = self.client._parse_okved_info( {"code": "62.01", "name": "Development", "version": "2001"}, diff --git a/tests/apps/parsers/test_eis_registry_client.py b/tests/apps/parsers/test_eis_registry_client.py new file mode 100644 index 0000000..a33d4d8 --- /dev/null +++ b/tests/apps/parsers/test_eis_registry_client.py @@ -0,0 +1,124 @@ +from __future__ import annotations + +from datetime import date +from decimal import Decimal +from urllib.parse import parse_qs, urlsplit + +from apps.parsers.clients.base import HTTPClientError +from apps.parsers.clients.common.eis_registry import EisRegistryProcurementClient +from apps.parsers.clients.common.schemas import GenericParserItem +from django.test import SimpleTestCase + + +def _record(*, external_id: str, inn: str, amount: str = "1") -> GenericParserItem: + return GenericParserItem( + source="procurements_44fz", + external_id=external_id, + inn=inn, + organisation_name="Наименование из ЕИС", + title="Закупка", + amount=Decimal(amount), + payload={"raw": external_id}, + ) + + +class _FakeHttpClient: + def __init__(self, pages: dict[int, bytes], *, error_page: int | None = None): + self.pages = pages + self.error_page = error_page + self.urls: list[str] = [] + + def download_file(self, url: str, **_kwargs) -> bytes: + self.urls.append(url) + page = int(parse_qs(urlsplit(url).query)["pageNumber"][0]) + if page == self.error_page: + raise HTTPClientError("page failed") + return self.pages[page] + + +class _FakeStructuredClient: + def __init__(self, pages: dict[int, bytes], records: dict[bytes, list]): + self.http_client = _FakeHttpClient(pages) + self.records = records + + def fetch_records(self, *, content: bytes, file_name: str): + assert file_name == "results.html" + return self.records[content] + + +class EisRegistryProcurementClientTest(SimpleTestCase): + def test_fetches_all_pages_filters_inn_and_deduplicates_external_id(self): + page_1 = b'' + page_2 = b'' + structured = _FakeStructuredClient( + {1: page_1, 2: page_2}, + { + page_1: [ + _record(external_id="same", inn="7701000001", amount="1"), + _record(external_id="foreign", inn="7701000099"), + ], + page_2: [ + _record(external_id="same", inn="7701000001", amount="2"), + _record(external_id="other", inn="7701000001"), + ], + }, + ) + client = EisRegistryProcurementClient( + source="procurements_44fz", + law="44", + _structured_client=structured, + ) + + records = client.fetch_for_customer( + inn="7701000001", + organization_name="ОПК Заказчик", + organization_id="org-1", + today=date(2026, 7, 19), + ) + + self.assertEqual([row.external_id for row in records], ["same", "other"]) + self.assertEqual(records[0].amount, Decimal("2")) + self.assertTrue(all(row.inn == "7701000001" for row in records)) + self.assertEqual(records[0].payload["organization_id"], "org-1") + query = parse_qs(urlsplit(structured.http_client.urls[0]).query) + self.assertEqual(query["fz44"], ["on"]) + self.assertEqual(query["strictEqual"], ["true"]) + self.assertEqual(query["publishDateFrom"], ["19.07.2025"]) + self.assertEqual(query["publishDateTo"], ["19.07.2026"]) + + def test_builds_223fz_query_and_accepts_empty_page(self): + empty_page = b"" + structured = _FakeStructuredClient({1: empty_page}, {empty_page: []}) + client = EisRegistryProcurementClient( + source="procurements_223fz", + law="223", + _structured_client=structured, + ) + + records = client.fetch_for_customer( + inn="7701000001", + organization_name="ОПК Заказчик", + organization_id="org-1", + ) + + self.assertEqual(records, []) + query = parse_qs(urlsplit(structured.http_client.urls[0]).query) + self.assertEqual(query["fz223"], ["on"]) + self.assertNotIn("fz44", query) + + def test_page_error_is_not_silenced(self): + page_1 = b'' + structured = _FakeStructuredClient({1: page_1}, {page_1: []}) + structured.http_client.error_page = 2 + client = EisRegistryProcurementClient( + source="procurements_44fz", + law="44", + _structured_client=structured, + ) + + with self.assertRaises(HTTPClientError): + client.fetch_for_customer( + inn="7701000001", + organization_name="ОПК Заказчик", + organization_id="org-1", + ) diff --git a/tests/apps/parsers/test_registry_procurement_tasks.py b/tests/apps/parsers/test_registry_procurement_tasks.py new file mode 100644 index 0000000..db6ab7a --- /dev/null +++ b/tests/apps/parsers/test_registry_procurement_tasks.py @@ -0,0 +1,252 @@ +from __future__ import annotations + +from contextlib import ExitStack +from types import SimpleNamespace +from unittest.mock import patch + +from apps.parsers import tasks as parser_tasks +from apps.parsers.clients.common.schemas import GenericParserItem +from apps.parsers.models import CheckoCollectionAttempt, ParserLoadLog +from django.test import TestCase, override_settings + +from tests.apps.parsers.organization_helpers import create_directory_organization + + +def _organization(index: int, *, membership: bool = True): + return create_directory_organization( + name=f"Организация {index}", + inn=f"7701000{index:03d}", + ogrn=f"1027700000{index:03d}", + opk_registry_membership=membership, + ) + + +class RegistryProcurementTasksTest(TestCase): + def test_daily_orchestrator_queues_addressed_sources_with_same_batch_limit(self): + task_names = ( + "parse_procurements_44fz", + "parse_procurements_223fz", + "parse_registry_contracts", + "parse_unfair_suppliers", + "parse_arbitration_cases", + "parse_fedresurs_bankruptcy", + "parse_registry_inspections", + "parse_trudvsem_vacancies", + ) + with ExitStack() as stack: + mocks = [ + stack.enter_context( + patch.object( + getattr(parser_tasks, task_name), + "delay", + return_value=SimpleNamespace(id=f"{task_name}-id"), + ) + ) + for task_name in task_names + ] + result = parser_tasks.parse_registry_enrichment_sources( + limit=20, + proxies=[], + ) + + self.assertEqual( + set(result), + { + "procurements_44fz", + "procurements_223fz", + "contracts", + "unfair_suppliers", + "arbitration", + "bankruptcy", + "inspections", + "vacancies", + }, + ) + for task_mock in mocks[:5]: + task_mock.assert_called_once_with( + proxies=[], + requested_by_id=None, + limit=20, + ) + mocks[5].assert_called_once_with(proxies=[], requested_by_id=None) + mocks[6].assert_called_once_with(proxies=[], requested_by_id=None, limit=20) + mocks[7].assert_called_once_with(proxies=[], requested_by_id=None) + + def test_eis_monthly_claims_advance_to_the_next_batch(self): + organizations = [_organization(index) for index in range(1, 4)] + outside_registry = _organization(9, membership=False) + requested_inns: list[str] = [] + + class _Client: + def __init__(self, **_kwargs): + return + + def fetch_for_customer(self, **kwargs): + requested_inns.append(kwargs["inn"]) + return [ + GenericParserItem( + source=ParserLoadLog.Source.PROCUREMENTS_44FZ, + external_id=f"notice-{kwargs['inn']}", + inn=kwargs["inn"], + organisation_name=kwargs["organization_name"], + ) + ] + + with patch.object(parser_tasks, "EisRegistryProcurementClient", _Client): + first = parser_tasks._fetch_eis_registry_procurement_records( + source=ParserLoadLog.Source.PROCUREMENTS_44FZ, + law="44", + limit=2, + proxies=[], + ) + second = parser_tasks._fetch_eis_registry_procurement_records( + source=ParserLoadLog.Source.PROCUREMENTS_44FZ, + law="44", + limit=2, + proxies=[], + ) + + self.assertEqual(len(first), 2) + self.assertEqual(len(second), 1) + self.assertEqual(requested_inns, [item.inn for item in organizations]) + self.assertNotIn(outside_registry.inn, requested_inns) + self.assertEqual( + CheckoCollectionAttempt.objects.filter( + source=CheckoCollectionAttempt.Source.PROCUREMENTS_44FZ, + status=CheckoCollectionAttempt.Status.SUCCESS, + ).count(), + 3, + ) + + def test_eis_request_error_creates_failed_monthly_claim(self): + organization = _organization(1) + + class _Client: + def __init__(self, **_kwargs): + return + + def fetch_for_customer(self, **_kwargs): + raise parser_tasks.HTTPClientError("EIS unavailable") + + with ( + patch.object(parser_tasks, "EisRegistryProcurementClient", _Client), + self.assertRaises(parser_tasks.ParserSourceSkipped), + ): + parser_tasks._fetch_eis_registry_procurement_records( + source=ParserLoadLog.Source.PROCUREMENTS_223FZ, + law="223", + limit=1, + proxies=[], + ) + + attempt = CheckoCollectionAttempt.objects.get( + organization=organization, + source=CheckoCollectionAttempt.Source.PROCUREMENTS_223FZ, + ) + self.assertEqual(attempt.status, CheckoCollectionAttempt.Status.FAILED) + + def test_explicit_organization_ids_do_not_fall_through_to_next_batch(self): + selected = _organization(1) + _organization(2) + calls: list[str] = [] + + class _Client: + def __init__(self, **_kwargs): + return + + def fetch_for_customer(self, **kwargs): + calls.append(kwargs["inn"]) + return [] + + kwargs = { + "source": ParserLoadLog.Source.PROCUREMENTS_44FZ, + "law": "44", + "limit": 20, + "proxies": [], + "organization_ids": [str(selected.uid)], + } + with patch.object(parser_tasks, "EisRegistryProcurementClient", _Client): + self.assertEqual( + parser_tasks._fetch_eis_registry_procurement_records(**kwargs), + [], + ) + with self.assertRaises(parser_tasks.ParserSourceSkipped): + parser_tasks._fetch_eis_registry_procurement_records(**kwargs) + + self.assertEqual(calls, [selected.inn]) + + @override_settings(CHECKO_API_KEY="test-key") + def test_checko_rnp_is_supplier_bound_and_not_requested_twice(self): + organization = _organization(1) + calls: list[str] = [] + + class _Client: + def __init__(self, **_kwargs): + return + + def get_company(self, request): + calls.append(request.inn) + return SimpleNamespace( + data=SimpleNamespace( + inn=organization.inn, + ogrn=organization.ogrn, + unfair_supplier=( + SimpleNamespace( + registry_number="RNP-1", + publish_date="2026-07-01", + approval_date="2026-06-30", + customer_short_name="Заказчик", + customer_full_name="Заказчик полный", + customer_inn="7702000002", + customer_kpp="770201001", + purchase_number="PURCHASE-1", + purchase_description="Поставка оборудования", + contract_price=1500000, + ), + ), + ) + ) + + with patch.object(parser_tasks, "CheckoClient", _Client): + records = parser_tasks._fetch_checko_unfair_supplier_records( + limit=1, + proxies=[], + ) + with self.assertRaises(parser_tasks.ParserSourceSkipped): + parser_tasks._fetch_checko_unfair_supplier_records( + limit=1, + proxies=[], + ) + + self.assertEqual(calls, [organization.inn]) + self.assertEqual(len(records), 1) + self.assertEqual(records[0].inn, organization.inn) + self.assertEqual(records[0].payload["supplier"]["inn"], organization.inn) + self.assertEqual(records[0].payload["customer"]["inn"], "7702000002") + self.assertEqual(records[0].amount, 1500000) + + @override_settings(CHECKO_API_KEY="test-key") + def test_checko_rnp_api_error_creates_failed_monthly_claim(self): + organization = _organization(1) + + class _Client: + def __init__(self, **_kwargs): + return + + def get_company(self, _request): + raise parser_tasks.CheckoError("Checko unavailable") + + with ( + patch.object(parser_tasks, "CheckoClient", _Client), + self.assertRaises(parser_tasks.ParserSourceSkipped), + ): + parser_tasks._fetch_checko_unfair_supplier_records( + limit=1, + proxies=[], + ) + + attempt = CheckoCollectionAttempt.objects.get( + organization=organization, + source=CheckoCollectionAttempt.Source.UNFAIR_SUPPLIERS, + ) + self.assertEqual(attempt.status, CheckoCollectionAttempt.Status.FAILED) diff --git a/tests/apps/parsers/test_source_cards_service.py b/tests/apps/parsers/test_source_cards_service.py index a6049e8..68fb3d5 100644 --- a/tests/apps/parsers/test_source_cards_service.py +++ b/tests/apps/parsers/test_source_cards_service.py @@ -240,10 +240,12 @@ class SourceCardServiceUnitTest(SimpleTestCase): "apps.parsers.source_cards.SourceCardService._enqueue_task", return_value={ "task_id": "task-9", - "task_name": "apps.parsers.tasks.sync_procurements", + "task_name": "apps.parsers.tasks.parse_registry_enrichment_sources", }, ) - def test_refresh_card_for_procurements_uses_default_law_type(self, enqueue_mock): + def test_refresh_card_for_procurements_uses_addressed_registry_task( + self, enqueue_mock + ): result = SourceCardService.refresh_card( slug="public-procurements", requested_by_id=10, @@ -254,12 +256,7 @@ class SourceCardServiceUnitTest(SimpleTestCase): self.assertEqual(result["tasks"][0]["task_id"], "task-9") self.assertEqual( enqueue_mock.call_args.kwargs["kwargs"], - { - "requested_by_id": 10, - "region_code": "77", - "law_type": "44", - "current_year": 2026, - }, + {"requested_by_id": 10}, ) @patch( @@ -540,6 +537,7 @@ class SourceCardServiceDatabaseTest(TestCase): self.assertEqual(card["records_count"], 3) self.assertEqual(card["organizations_count"], 2) source_items = {item["code"]: item for item in card["source_items"]} + self.assertEqual(source_items["procurements"]["records_count"], 2) self.assertEqual(source_items["procurements_44fz"]["organizations_count"], 1) self.assertEqual(source_items["procurements_223fz"]["organizations_count"], 1) self.assertEqual(source_items["contracts"]["organizations_count"], 1) diff --git a/tests/apps/parsers/test_source_cards_views.py b/tests/apps/parsers/test_source_cards_views.py index 7b65485..0d2b98a 100644 --- a/tests/apps/parsers/test_source_cards_views.py +++ b/tests/apps/parsers/test_source_cards_views.py @@ -3,6 +3,8 @@ from __future__ import annotations import hashlib +from types import SimpleNamespace +from unittest.mock import patch from apps.core.models import BackgroundJob, JobStatus from apps.parsers.models import ParserLoadLog @@ -259,19 +261,24 @@ class SourceCardsApiTestCase(APITestCase): ).exists() ) - def test_refresh_procurements_requires_region_code(self): + @override_settings(CELERY_TASK_ALWAYS_EAGER=False) + def test_refresh_procurements_does_not_require_region_code(self): self.client.force_authenticate(self.admin) - response = self.client.post( - reverse( - "api_v1:sources:source-cards-refresh", - kwargs={"slug": "public-procurements"}, - ), - {}, - format="json", - ) + with patch( + "apps.parsers.tasks.parse_registry_enrichment_sources.apply_async", + return_value=SimpleNamespace(id="task-procurements"), + ): + response = self.client.post( + reverse( + "api_v1:sources:source-cards-refresh", + kwargs={"slug": "public-procurements"}, + ), + {}, + format="json", + ) - self.assertEqual(response.status_code, status.HTTP_400_BAD_REQUEST) - self.assertFalse(response.data["success"]) + self.assertEqual(response.status_code, status.HTTP_202_ACCEPTED) + self.assertEqual(response.data["task_id"], "task-procurements") def test_refresh_forbidden_for_regular_user(self): response = self.client.post( diff --git a/tests/apps/parsers/test_sources_api_e2e.py b/tests/apps/parsers/test_sources_api_e2e.py index c873d30..cf74476 100644 --- a/tests/apps/parsers/test_sources_api_e2e.py +++ b/tests/apps/parsers/test_sources_api_e2e.py @@ -6,7 +6,6 @@ from unittest.mock import patch from apps.core.models import BackgroundJob, JobStatus from apps.parsers.models import FinancialReport, FinancialReportLine, ParserLoadLog from django.urls import reverse -from organizations.models import Organization from organizations.source_backfill import OrganizationSourceBackfillService from rest_framework import status from rest_framework.test import APITestCase @@ -206,7 +205,7 @@ class SourcesApiE2ETest(APITestCase): ) with patch( - "apps.parsers.tasks.sync_procurements.apply_async", + "apps.parsers.tasks.parse_registry_enrichment_sources.apply_async", return_value=SimpleNamespace(id="task-procurements"), ): procurements_response = self.client.post( @@ -241,7 +240,7 @@ class SourcesApiE2ETest(APITestCase): self.assertTrue( BackgroundJob.objects.filter( task_id="task-procurements", - task_name="apps.parsers.tasks.sync_procurements", + task_name="apps.parsers.tasks.parse_registry_enrichment_sources", user_id=self.admin.id, ).exists() ) diff --git a/tests/apps/parsers/test_tasks.py b/tests/apps/parsers/test_tasks.py index 38b19d0..dcccb57 100644 --- a/tests/apps/parsers/test_tasks.py +++ b/tests/apps/parsers/test_tasks.py @@ -1855,18 +1855,10 @@ class MinpromtorgTasksTestCase(TestCase): "parse_industrial_products", "parse_manufactures", "sync_inspections", - "parse_procurements_44fz", - "parse_procurements_223fz", - "parse_contracts", - "parse_registry_contracts", - "parse_unfair_suppliers", + "parse_registry_enrichment_sources", "parse_fas_goz_evasion", "sync_fns_financial_reports", - "parse_arbitration_cases", - "parse_fedresurs_bankruptcy", - "parse_registry_inspections", "parse_fstec_registers", - "parse_trudvsem_vacancies", ) with ExitStack() as stack: delays = { @@ -1887,8 +1879,10 @@ class MinpromtorgTasksTestCase(TestCase): "sync_fns_financial_reports-id", ) delays["sync_fns_financial_reports"].assert_called_once_with(proxies=[]) - delays["parse_registry_contracts"].assert_called_once_with(proxies=[]) - delays["parse_registry_inspections"].assert_called_once_with(proxies=[]) + delays["parse_registry_enrichment_sources"].assert_called_once_with( + limit=250, + proxies=[], + ) def test_parse_all_minpromtorg_without_adapter(self): with TestHTTPServer() as server: