From 31eaa7e9071b7e07c2e5ec59337259fcdc64bbe8 Mon Sep 17 00:00:00 2001 From: Aleksandr Meshchriakov Date: Fri, 14 Aug 2026 17:02:56 +0200 Subject: [PATCH] fix(parsers): make source jobs resumable --- docs/parser-external-access-note-ru.md | 7 +- src/apps/parsers/clients/vacancies.py | 18 +- src/apps/parsers/serializers.py | 2 +- src/apps/parsers/source_cards.py | 18 +- src/apps/parsers/source_registry.py | 10 +- src/apps/parsers/tasks.py | 710 ++++++++++++------ src/organizations/views.py | 8 +- src/settings/base.py | 2 - .../test_api_v2_source_extensions.py | 53 ++ tests/apps/parsers/test_gosedo_media_news.py | 76 +- .../apps/parsers/test_source_cards_service.py | 25 + tests/apps/parsers/test_source_registry.py | 17 + tests/apps/parsers/test_tasks.py | 401 +++++----- tests/apps/parsers/test_vacancy_clients.py | 13 +- 14 files changed, 884 insertions(+), 476 deletions(-) diff --git a/docs/parser-external-access-note-ru.md b/docs/parser-external-access-note-ru.md index 68d0d7c..89aca05 100644 --- a/docs/parser-external-access-note-ru.md +++ b/docs/parser-external-access-note-ru.md @@ -24,9 +24,7 @@ | КАД Арбитр через внешний сервис данных | Официальный источник в каталоге: `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`. | +| Вакансии: Работа России | `https://opendata.trudvsem.ru/api/v1/vacancies`, `https://opendata.trudvsem.ru/api/v1/vacancies/company/inn/{inn}` | JSON-список вакансий, включая работодателя, ИНН/ОГРН при наличии, название вакансии, дату, зарплату, статус, ссылку. Организации обрабатываются короткими пакетами, результат сохраняется после каждой организации. | | Внешний сервис данных: контракты и проверки по организациям | Служебные API контрактов и проверок | JSON-данные по контрактам и проверкам для активных организаций из внутренних реестров. В запрос передаются ИНН/ОГРН, API-ключ передается параметром `key`. | | Proxy-Tools | `https://proxy-tools.com/api/v1/proxies` | Служебная загрузка списка RU-прокси для парсеров. Используется только при заданном `PROXY_TOOLS_API_KEY`; запрос идет с Bearer-токеном. | @@ -51,7 +49,7 @@ - коды регионов; - ИНН/ОГРН организаций из внутренних активных реестров; - поисковая строка по названию организации для вакансий; -- служебные ключи API из окружения для ЕИС, внешнего сервиса данных, SuperJob и Proxy-Tools. +- служебные ключи API из окружения для ЕИС, внешнего сервиса данных и Proxy-Tools. Ключи в коде не захардкожены, берутся из переменных окружения. @@ -61,4 +59,5 @@ - Для ФНС текущая реализация backend не скачивает файлы автоматически с сайта ФНС, а обрабатывает уже полученные Excel/ZIP-файлы через папку наблюдения или API-загрузку. - Для `proverki.gov.ru` возможен запуск headless Chromium через Playwright, потому что часть загрузок доступна через JS-интерфейс портала. - Для ФСТЭК SSL-верификация отключена настройкой клиента источника. +- Источники вакансий HeadHunter и SuperJob отключены; сбор выполняется только из «Работы России». - Runtime-прокси из БД используются только при включенном `PARSER_USE_RUNTIME_PROXIES=true`; отдельная задача синхронизации прокси обращается к Proxy-Tools только при наличии `PROXY_TOOLS_API_KEY`. diff --git a/src/apps/parsers/clients/vacancies.py b/src/apps/parsers/clients/vacancies.py index c9da5b8..6c92382 100644 --- a/src/apps/parsers/clients/vacancies.py +++ b/src/apps/parsers/clients/vacancies.py @@ -25,6 +25,8 @@ TRUDVSEM_SOURCE = "trudvsem" HH_SOURCE = "hh" SUPERJOB_SOURCE = "superjob" SUPPORTED_VACANCY_SOURCES = (TRUDVSEM_SOURCE, HH_SOURCE, SUPERJOB_SOURCE) +ACTIVE_VACANCY_SOURCES = (TRUDVSEM_SOURCE,) +DISABLED_VACANCY_SOURCES = (HH_SOURCE, SUPERJOB_SOURCE) class VacanciesClientError(HTTPClientError): @@ -297,21 +299,17 @@ class VacanciesClient: clients: dict[str, VacancyProvider] = { TRUDVSEM_SOURCE: TrudvsemClient(proxies=self.proxies), - HH_SOURCE: HHVacanciesClient( - user_agent=self.hh_user_agent, - proxies=self.proxies, - ), } - if self.superjob_app_id: - clients[SUPERJOB_SOURCE] = SuperJobVacanciesClient( - app_id=self.superjob_app_id, - proxies=self.proxies, - ) self._source_clients_cache = clients return clients def _selected_sources(self) -> list[str]: - selected = self.sources or list(SUPPORTED_VACANCY_SOURCES) + selected = self.sources or list(ACTIVE_VACANCY_SOURCES) + disabled = sorted(set(selected) & set(DISABLED_VACANCY_SOURCES)) + if disabled: + raise VacanciesClientError( + f"Disabled vacancy sources: {', '.join(disabled)}" + ) unknown = sorted(set(selected) - set(SUPPORTED_VACANCY_SOURCES)) if unknown: raise VacanciesClientError( diff --git a/src/apps/parsers/serializers.py b/src/apps/parsers/serializers.py index 7a195b2..c33f30e 100644 --- a/src/apps/parsers/serializers.py +++ b/src/apps/parsers/serializers.py @@ -540,7 +540,7 @@ class ParserRunRequestSerializer(serializers.Serializer): allow_empty=True, ) vacancy_sources = serializers.ListField( - child=serializers.ChoiceField(choices=["trudvsem", "hh", "superjob"]), + child=serializers.ChoiceField(choices=["trudvsem"]), required=False, allow_empty=False, ) diff --git a/src/apps/parsers/source_cards.py b/src/apps/parsers/source_cards.py index 44f74ac..b36f415 100644 --- a/src/apps/parsers/source_cards.py +++ b/src/apps/parsers/source_cards.py @@ -335,7 +335,7 @@ SOURCE_CARD_DEFINITIONS: tuple[SourceCardDefinition, ...] = ( SourceItemDefinition( code="trudvsem", title="Вакансии Работа России", - description="Вакансии работодателей из Работа России, HH и SuperJob.", + description="Вакансии работодателей из Работа России.", parser_source=ParserLoadLog.Source.TRUDVSEM, ), ), @@ -365,7 +365,7 @@ SOURCE_CARD_DEFINITIONS: tuple[SourceCardDefinition, ...] = ( title="Новости СМИ", description="Загруженные упоминания организаций в СМИ с оценкой тональности.", order=110, - task_names=(), + task_names=("apps.parsers.tasks.parse_media_news",), source_items=( SourceItemDefinition( code="media_mentions", @@ -1160,6 +1160,20 @@ class SourceCardService: meta: dict[str, Any], kwargs: dict[str, Any], ) -> dict[str, str]: + if task_name == "apps.parsers.tasks.parse_trudvsem_vacancies": + existing_job = ( + BackgroundJobService.get_queryset() + .filter(task_name=task_name) + .filter(cls._fresh_active_job_filter()) + .order_by("-created_at") + .first() + ) + if existing_job is not None: + return { + "task_id": existing_job.task_id, + "task_name": existing_job.task_name, + } + task_id = str(uuid.uuid4()) BackgroundJobService.create_job( task_id=task_id, diff --git a/src/apps/parsers/source_registry.py b/src/apps/parsers/source_registry.py index b47a262..d7a533f 100644 --- a/src/apps/parsers/source_registry.py +++ b/src/apps/parsers/source_registry.py @@ -285,15 +285,15 @@ PARSER_SOURCES: dict[str, ParserSourceDescriptor] = { key="trudvsem", source=ParserLoadLog.Source.TRUDVSEM, title="Вакансии", - agency="Работа России / HH / SuperJob", - data_scope="Вакансии работодателей из нескольких job-board источников", + agency="Работа России", + data_scope="Вакансии работодателей из официального источника Работа России", task_name="apps.parsers.tasks.parse_trudvsem_vacancies", upstream_url="https://opendata.trudvsem.ru/api/v1/vacancies", access_method="public_api", - parser_strategy="multi_source_vacancies_api", + parser_strategy="incremental_trudvsem_api", source_notes=( - "Internal source key remains trudvsem for backward compatibility; " - "payload.vacancy_source distinguishes trudvsem, hh and superjob." + "Пакетная обработка организаций с промежуточным сохранением, " + "прогрессом и возобновлением после остановки." ), api_route="trudvsem/vacancies", ), diff --git a/src/apps/parsers/tasks.py b/src/apps/parsers/tasks.py index 27337f9..729c1bb 100644 --- a/src/apps/parsers/tasks.py +++ b/src/apps/parsers/tasks.py @@ -16,6 +16,7 @@ from datetime import datetime from decimal import Decimal from pathlib import Path +from apps.core.models import BackgroundJob, JobStatus from apps.core.services import BackgroundJobService from apps.core.tasks import PeriodicTask as CorePeriodicTask from apps.parsers.checko_collection import ( @@ -73,7 +74,12 @@ from apps.parsers.source_registry import PARSER_SOURCES from celery import shared_task from django.conf import settings from django.db.models import Q -from organizations.models import Organization as SourceOrganization +from organizations.models import ( + Organization as SourceOrganization, +) +from organizations.models import ( + OrganizationSourceRecord, +) from organizations.services import ( normalize_organization_name as normalize_identity_name, ) @@ -141,30 +147,9 @@ class FNSApiFetchResult: VACANCY_REGISTRY_MAX_PAGES_PER_ORGANIZATION = 100 -VACANCY_REGISTRY_TEXT_SEARCH_MAX_PAGES_PER_ORGANIZATION = 1 -VACANCY_EMPLOYER_WORD_RE = re.compile(r"[0-9A-Za-zА-Яа-яЁё]+") -VACANCY_EMPLOYER_IGNORED_WORDS = { - "ао", - "акционерное", - "государственное", - "зао", - "индивидуальный", - "ип", - "муниципальное", - "нао", - "некоммерческая", - "оао", - "общество", - "ограниченной", - "ооо", - "ответственностью", - "пао", - "предприниматель", - "публичное", - "с", - "унитарное", - "фгуп", -} +VACANCY_REGISTRY_ORGANIZATIONS_PER_TASK = 25 +VACANCY_TASK_NAME = "apps.parsers.tasks.parse_trudvsem_vacancies" +TRUDVSEM_VACANCY_SOURCE = "trudvsem" def _resolve_lookup_limit( @@ -3787,7 +3772,427 @@ def cleanup_stale_parser_loads( } -@shared_task(bind=True) +def _normalize_trudvsem_vacancy_sources( + vacancy_sources: list[str] | None, +) -> list[str]: + """Оставить единственный включённый источник вакансий.""" + selected = vacancy_sources or [TRUDVSEM_VACANCY_SOURCE] + disabled = sorted(set(selected) - {TRUDVSEM_VACANCY_SOURCE}) + if disabled: + raise ValueError("Поддерживается только источник trudvsem") + return [TRUDVSEM_VACANCY_SOURCE] + + +def _vacancy_targets_signature(targets: list[RegistryLookupTarget]) -> str: + """Вернуть безопасный отпечаток порядка организаций для возобновления.""" + digest = hashlib.sha256() + for target in targets: + digest.update(target.organization_id.encode("utf-8")) + digest.update(b"\n") + return digest.hexdigest() + + +def _vacancy_batch_records_count(batch_id: int) -> int: + return OrganizationSourceRecord.objects.filter( + source=TRUDVSEM_VACANCY_SOURCE, + load_batch=batch_id, + ).count() + + +def _find_resumable_vacancy_job( + *, + exclude_task_id: str, + targets_signature: str, + targets_count: int, +): + """Найти последний совместимый остановленный проход вакансий.""" + candidates = ( + BackgroundJobService.get_queryset() + .filter( + task_name=VACANCY_TASK_NAME, + status__in=[JobStatus.REVOKED, JobStatus.FAILURE], + ) + .exclude(task_id=exclude_task_id) + .order_by("-updated_at")[:20] + ) + for candidate in candidates: + meta = candidate.meta or {} + next_offset = int(meta.get("next_offset") or 0) + batch_id = meta.get("batch_id") + if ( + meta.get("targets_signature") != targets_signature + or int(meta.get("total_organizations") or 0) != targets_count + or next_offset <= 0 + or next_offset >= targets_count + or batch_id is None + ): + continue + load = ParserLoadLog.objects.filter( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=batch_id, + ).first() + if load is not None: + return candidate, load + return None, None + + +def _update_vacancy_registry_checkpoint( + *, + job, + load_log: ParserLoadLog, + batch_id: int, + processed: int, + total: int, + failed: int, + targets_signature: str, +) -> int: + """Зафиксировать прогресс после полностью обработанной организации.""" + saved_count = _vacancy_batch_records_count(batch_id) + progress = 100 if total == 0 else min(99, round(processed * 100 / total)) + job.meta = { + **(job.meta or {}), + "batch_id": batch_id, + "next_offset": processed, + "total_organizations": total, + "failed_organizations": failed, + "saved_records": saved_count, + "targets_signature": targets_signature, + } + job.progress = progress + job.progress_message = ( + f"Обработано организаций: {processed} из {total}; " + f"сохранено вакансий: {saved_count}; ошибок: {failed}" + ) + job.save(update_fields=["meta", "progress", "progress_message", "updated_at"]) + ParserLoadLogService.update( + load_log, + status=ParserLoadLog.Status.IN_PROGRESS, + records_count=saved_count, + error_message="", + ) + return saved_count + + +def _finish_revoked_vacancy_registry_run( + *, + job, + load_log: ParserLoadLog, + batch_id: int, + processed: int, + total: int, + failed: int, +) -> dict: + """Закрыть отменённый проход, сохранив уже записанный результат.""" + saved_count = _vacancy_batch_records_count(batch_id) + ParserLoadLogService.update( + load_log, + status=ParserLoadLog.Status.SKIPPED, + records_count=saved_count, + error_message=( + f"Остановлено после {processed} из {total} организаций; " + "сохранённые записи не удалены" + ), + ) + return { + "batch_id": batch_id, + "saved": saved_count, + "processed_organizations": processed, + "failed_organizations": failed, + "status": "revoked", + "resumed": bool((job.meta or {}).get("resumed")), + } + + +@dataclass +class VacancyRegistryRunState: + """Состояние одного возобновляемого прохода вакансий.""" + + task_id: str + targets: list[RegistryLookupTarget] + targets_signature: str + job: BackgroundJob + load_log: ParserLoadLog + batch_id: int + next_offset: int + failed: int + resumed: bool + + @property + def total(self) -> int: + return len(self.targets) + + +def _prepare_vacancy_registry_run( + self, + *, + registry_organization_limit: int | None, + requested_by_id: int | None, + run_task_id: str | None, +) -> VacancyRegistryRunState: + """Создать новый проход или восстановить совместимый checkpoint.""" + task_id = run_task_id or self.request.id or str(uuid.uuid4()) + targets = _active_registry_vacancy_targets(limit=registry_organization_limit) + total = len(targets) + targets_signature = _vacancy_targets_signature(targets) + resumed = False + + job = BackgroundJobService.get_by_task_id_or_none(task_id) + if run_task_id is not None: + if job is None: + raise RuntimeError("Не найдена родительская задача вакансий") + meta = job.meta or {} + batch_id = int(meta["batch_id"]) + load_log = ParserLoadLog.objects.get( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=batch_id, + ) + if meta.get("targets_signature") != targets_signature: + meta = { + **meta, + "next_offset": 0, + "failed_organizations": 0, + "total_organizations": total, + "targets_signature": targets_signature, + } + job.meta = meta + job.save(update_fields=["meta", "updated_at"]) + resumed = bool(meta.get("resumed")) + else: + resumable_job, resumable_load = _find_resumable_vacancy_job( + exclude_task_id=task_id, + targets_signature=targets_signature, + targets_count=total, + ) + if resumable_job is not None and resumable_load is not None: + resumable_meta = resumable_job.meta or {} + batch_id = int(resumable_meta["batch_id"]) + load_log = resumable_load + next_offset = int(resumable_meta["next_offset"]) + failed = int(resumable_meta.get("failed_organizations") or 0) + resumed = True + ParserLoadLogService.update( + load_log, + status=ParserLoadLog.Status.IN_PROGRESS, + records_count=_vacancy_batch_records_count(batch_id), + error_message="", + ) + else: + ( + load_log, + batch_id, + ) = ParserLoadLogService.create_load_log_with_next_batch_id( + source=ParserLoadLog.Source.TRUDVSEM, + status=ParserLoadLog.Status.IN_PROGRESS, + ) + next_offset = 0 + failed = 0 + + job = _get_or_create_background_job( + task_id=task_id, + task_name=VACANCY_TASK_NAME, + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=batch_id, + requested_by_id=requested_by_id, + meta={ + "source_key": TRUDVSEM_VACANCY_SOURCE, + "next_offset": next_offset, + "total_organizations": total, + "failed_organizations": failed, + "saved_records": _vacancy_batch_records_count(batch_id), + "targets_signature": targets_signature, + "resumed": resumed, + "resumed_from_task_id": ( + resumable_job.task_id if resumable_job is not None else None + ), + }, + ) + job.mark_started() + _update_vacancy_registry_checkpoint( + job=job, + load_log=load_log, + batch_id=batch_id, + processed=next_offset, + total=total, + failed=failed, + targets_signature=targets_signature, + ) + + meta = job.meta or {} + return VacancyRegistryRunState( + task_id=task_id, + targets=targets, + targets_signature=targets_signature, + job=job, + load_log=load_log, + batch_id=batch_id, + next_offset=int(meta.get("next_offset") or 0), + failed=int(meta.get("failed_organizations") or 0), + resumed=resumed, + ) + + +def _vacancy_registry_run_is_revoked(state: VacancyRegistryRunState) -> bool: + state.job.refresh_from_db(fields=["status", "meta"]) + return state.job.status == JobStatus.REVOKED + + +def _revoked_vacancy_registry_result(state: VacancyRegistryRunState) -> dict: + return _finish_revoked_vacancy_registry_run( + job=state.job, + load_log=state.load_log, + batch_id=state.batch_id, + processed=state.next_offset, + total=state.total, + failed=state.failed, + ) + + +def _process_vacancy_registry_chunk( + state: VacancyRegistryRunState, + *, + limit: int, + proxies: list[str] | None, +) -> bool: + """Обработать ограниченный блок; вернуть True при отмене.""" + chunk_end = min( + state.total, + state.next_offset + max(1, VACANCY_REGISTRY_ORGANIZATIONS_PER_TASK), + ) + + with VacanciesClient( + proxies=proxies, + sources=[TRUDVSEM_VACANCY_SOURCE], + ) as client: + for index in range(state.next_offset, chunk_end): + if _vacancy_registry_run_is_revoked(state): + return True + + target = state.targets[index] + try: + records = _fetch_registry_target_vacancy_records( + client, + target, + page_size=max(1, limit), + ) + if records: + GenericParserRecordService.save_records( + records, + batch_id=state.batch_id, + source=ParserLoadLog.Source.TRUDVSEM, + ) + except Exception as exc: + state.failed += 1 + logger.warning( + "Vacancy fetch failed for registry target %d of %d: %s", + index + 1, + state.total, + exc, + ) + + state.next_offset = index + 1 + _update_vacancy_registry_checkpoint( + job=state.job, + load_log=state.load_log, + batch_id=state.batch_id, + processed=state.next_offset, + total=state.total, + failed=state.failed, + targets_signature=state.targets_signature, + ) + return _vacancy_registry_run_is_revoked(state) + + +def _queue_or_finish_vacancy_registry_run( + state: VacancyRegistryRunState, + *, + limit: int, + registry_organization_limit: int | None, + proxies: list[str] | None, + requested_by_id: int | None, +) -> dict: + """Передать следующий блок в очередь или закрыть проход.""" + saved_count = _vacancy_batch_records_count(state.batch_id) + if state.next_offset < state.total: + try: + parse_trudvsem_vacancies.apply_async( + kwargs={ + "limit": limit, + "vacancy_sources": [TRUDVSEM_VACANCY_SOURCE], + "registry_organizations_only": True, + "registry_organization_limit": registry_organization_limit, + "proxies": proxies, + "requested_by_id": requested_by_id, + "_vacancy_run_task_id": state.task_id, + } + ) + except Exception as exc: + ParserLoadLogService.mark_failed(state.load_log, str(exc)) + state.job.fail(error=str(exc)) + raise + return { + "batch_id": state.batch_id, + "saved": saved_count, + "processed_organizations": state.next_offset, + "failed_organizations": state.failed, + "status": "in_progress", + "resumed": state.resumed, + } + + result = { + "batch_id": state.batch_id, + "saved": saved_count, + "processed_organizations": state.next_offset, + "failed_organizations": state.failed, + "status": "success", + "resumed": state.resumed, + } + if state.total > 0 and state.failed >= state.total: + message = "Не удалось обработать ни одной организации в Работа России" + result["status"] = "failure" + ParserLoadLogService.mark_failed(state.load_log, message) + state.job.fail(error=message) + return result + + ParserLoadLogService.update( + state.load_log, + status=ParserLoadLog.Status.SUCCESS, + records_count=saved_count, + error_message="", + ) + state.job.complete(result=result) + return result + + +def _run_incremental_registry_vacancies( + self, + *, + limit: int, + registry_organization_limit: int | None, + proxies: list[str] | None, + requested_by_id: int | None, + run_task_id: str | None, +) -> dict: + """Обрабатывать реестр короткими возобновляемыми Celery-проходами.""" + state = _prepare_vacancy_registry_run( + self, + registry_organization_limit=registry_organization_limit, + requested_by_id=requested_by_id, + run_task_id=run_task_id, + ) + if _vacancy_registry_run_is_revoked(state): + return _revoked_vacancy_registry_result(state) + if _process_vacancy_registry_chunk(state, limit=limit, proxies=proxies): + return _revoked_vacancy_registry_result(state) + return _queue_or_finish_vacancy_registry_run( + state, + limit=limit, + registry_organization_limit=registry_organization_limit, + proxies=proxies, + requested_by_id=requested_by_id, + ) + + +@shared_task(bind=True, acks_late=True, reject_on_worker_lost=True) def parse_trudvsem_vacancies( self, *, @@ -3801,14 +4206,50 @@ def parse_trudvsem_vacancies( registry_organization_limit: int | None = None, proxies: list[str] | None = None, requested_by_id: int | None = None, + _vacancy_run_task_id: str | None = None, ) -> dict: """Парсинг вакансий по активным организациям реестров.""" + vacancy_sources = _normalize_trudvsem_vacancy_sources(vacancy_sources) proxies = _resolve_proxies(proxies) + if _should_fetch_registry_organization_vacancies( + registry_organizations_only=registry_organizations_only, + region_code=region_code, + company_inn=company_inn, + text=text, + ): + try: + return _run_incremental_registry_vacancies( + self, + limit=limit, + registry_organization_limit=registry_organization_limit, + proxies=proxies, + requested_by_id=requested_by_id, + run_task_id=_vacancy_run_task_id, + ) + except Exception as exc: + task_id = _vacancy_run_task_id or self.request.id + job = ( + BackgroundJobService.get_by_task_id_or_none(task_id) + if task_id + else None + ) + if job is not None and not job.is_finished: + batch_id = (job.meta or {}).get("batch_id") + if batch_id is not None: + load_log = ParserLoadLog.objects.filter( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=batch_id, + ).first() + if load_log is not None: + ParserLoadLogService.mark_failed(load_log, str(exc)) + job.fail(error=str(exc)) + raise + return _run_generic_parser( self, source_key="trudvsem", source=ParserLoadLog.Source.TRUDVSEM, - task_name="apps.parsers.tasks.parse_trudvsem_vacancies", + task_name=VACANCY_TASK_NAME, requested_by_id=requested_by_id, fetch_records=lambda: _fetch_vacancy_records( proxies=proxies, @@ -3818,8 +4259,6 @@ def parse_trudvsem_vacancies( company_inn=company_inn, text=text, vacancy_sources=vacancy_sources, - registry_organizations_only=registry_organizations_only, - registry_organization_limit=registry_organization_limit, ), ) @@ -3833,26 +4272,9 @@ def _fetch_vacancy_records( company_inn: str | None, text: str | None, vacancy_sources: list[str] | None, - registry_organizations_only: bool, - registry_organization_limit: int | None, ) -> list[GenericParserItem]: - if _should_fetch_registry_organization_vacancies( - registry_organizations_only=registry_organizations_only, - region_code=region_code, - company_inn=company_inn, - text=text, - ): - return _fetch_registry_organization_vacancy_records( - proxies=proxies, - limit=limit, - vacancy_sources=vacancy_sources, - registry_organization_limit=registry_organization_limit, - ) - with VacanciesClient( proxies=proxies, - superjob_app_id=getattr(settings, "SUPERJOB_APP_ID", ""), - hh_user_agent=getattr(settings, "HH_USER_AGENT", ""), sources=vacancy_sources, ) as client: return client.fetch_vacancies( @@ -3914,60 +4336,11 @@ def _fetch_registry_target_vacancy_records( *, page_size: int, ) -> list[GenericParserItem]: - iter_source_clients = getattr(client, "iter_source_clients", None) - if iter_source_clients is None: - return _fetch_registry_target_source_vacancy_records( - client, - target, - page_size=page_size, - company_inn=target.inn, - ) - - records: list[GenericParserItem] = [] - errors: list[str] = [] - attempts = 0 - - for source, source_client in iter_source_clients(): - if getattr(source_client, "supports_company_inn", False): - kwargs = {"company_inn": target.inn} - else: - if not target.name: - logger.info( - "Vacancy source %s is skipped for registry organization %s: " - "empty organization name", - source, - target.organization_id, - ) - continue - kwargs = {"text": _vacancy_registry_text_query(target)} - - attempts += 1 - try: - source_records = _fetch_registry_target_source_vacancy_records( - source_client, - target, - page_size=page_size, - **kwargs, - ) - except Exception as exc: - logger.warning( - "Vacancy source %s failed for registry organization %s (%s): %s", - source, - target.organization_id, - target.inn, - exc, - ) - errors.append(f"{source}: {exc}") - continue - - records.extend(source_records) - - if errors and not records and attempts: - raise RuntimeError( - "All vacancy sources failed for registry organization " - f"{target.organization_id} ({target.inn}); first error: {errors[0]}" - ) - return records + return _fetch_registry_target_source_vacancy_records( + client, + target, + page_size=page_size, + ) def _fetch_registry_target_source_vacancy_records( @@ -3975,155 +4348,28 @@ def _fetch_registry_target_source_vacancy_records( target: RegistryLookupTarget, *, page_size: int, - company_inn: str | None = None, - text: str | None = None, ) -> list[GenericParserItem]: records: list[GenericParserItem] = [] offset = 0 - filter_by_employer_name = company_inn is None - max_pages = ( - VACANCY_REGISTRY_TEXT_SEARCH_MAX_PAGES_PER_ORGANIZATION - if filter_by_employer_name - else VACANCY_REGISTRY_MAX_PAGES_PER_ORGANIZATION - ) - for _ in range(max_pages): + for _ in range(VACANCY_REGISTRY_MAX_PAGES_PER_ORGANIZATION): page_records = source_client.fetch_vacancies( limit=page_size, offset=offset, - company_inn=company_inn, - text=text, + company_inn=target.inn, + ) + records.extend( + _attach_registry_vacancy_target(record, target) for record in page_records ) - if filter_by_employer_name: - matched_records = [ - record - for record in page_records - if _vacancy_record_matches_registry_target(record, target) - ] - else: - matched_records = page_records - records.extend(matched_records) if len(page_records) < page_size: return records offset += page_size - if filter_by_employer_name: - return records - raise RuntimeError( "Vacancy registry organization page limit reached " f"for organization {target.organization_id} ({target.inn})" ) -def _vacancy_record_matches_registry_target( - record: GenericParserItem, - target: RegistryLookupTarget, -) -> bool: - target_key = _vacancy_employer_match_key(target.name) - employer_key = _vacancy_employer_match_key(_vacancy_record_employer_name(record)) - if not target_key or not employer_key: - return False - if target_key == employer_key: - return True - if min(len(target_key), len(employer_key)) < 8: - return False - return target_key in employer_key or employer_key in target_key - - -def _vacancy_registry_text_query(target: RegistryLookupTarget) -> str: - return _vacancy_employer_match_key(target.name) or target.name - - -def _vacancy_record_employer_name(record: GenericParserItem) -> str: - if record.organisation_name: - return record.organisation_name - - payload = record.payload if isinstance(record.payload, dict) else {} - for key in ("employer", "company"): - nested = payload.get(key) - if isinstance(nested, dict) and nested.get("name"): - return str(nested["name"]) - for key in ("firm_name", "company_name", "organisation_name"): - value = payload.get(key) - if value: - return str(value) - return "" - - -def _vacancy_employer_match_key(name: str) -> str: - words = [] - for match in VACANCY_EMPLOYER_WORD_RE.finditer(name.casefold().replace("ё", "е")): - word = match.group(0) - if word not in VACANCY_EMPLOYER_IGNORED_WORDS: - words.append(word) - return " ".join(words) - - -def _fetch_registry_organization_vacancy_records( - *, - proxies: list[str] | None, - limit: int, - vacancy_sources: list[str] | None, - registry_organization_limit: int | None, -) -> list[GenericParserItem]: - targets = _active_registry_vacancy_targets(limit=registry_organization_limit) - if not targets: - logger.info("No active registry organizations for vacancies sync") - return [] - - page_size = max(1, limit) - records: list[GenericParserItem] = [] - errors: list[str] = [] - successful_fetches = 0 - with VacanciesClient( - proxies=proxies, - superjob_app_id=getattr(settings, "SUPERJOB_APP_ID", ""), - hh_user_agent=getattr(settings, "HH_USER_AGENT", ""), - sources=vacancy_sources, - ) as client: - for target in targets: - try: - organization_records = _fetch_registry_target_vacancy_records( - client, - target, - page_size=page_size, - ) - except Exception as exc: - logger.warning( - "Vacancy fetch failed for registry organization %s (%s): %s", - target.organization_id, - target.inn, - exc, - ) - errors.append(f"{target.inn}: {exc}") - continue - - successful_fetches += 1 - records.extend( - _attach_registry_vacancy_target(record, target) - for record in organization_records - ) - - if errors and successful_fetches == 0: - raise RuntimeError( - "All registry organization vacancy fetches failed; " - f"first error: {errors[0]}" - ) - if errors: - logger.warning( - "Vacancy registry organization sync completed with %d failed " - "organizations", - len(errors), - ) - - logger.info( - "Fetched %d vacancy records for %d active registry organizations", - len(records), - len(targets), - ) - return records - - # ============================================================================= # FNS Tasks (File Watch & Processing) # ============================================================================= diff --git a/src/organizations/views.py b/src/organizations/views.py index ecb0085..d98d790 100644 --- a/src/organizations/views.py +++ b/src/organizations/views.py @@ -95,6 +95,7 @@ SOURCE_GROUP_VALUES = [choice.value for choice in SourceGroup] SOURCE_RECORD_ORDERING_FIELDS = ( "record_date", "extension__organization__name", + "extension__organization__full_name", "created_at", "updated_at", "title", @@ -103,6 +104,7 @@ SOURCE_RECORD_ORDERING_FIELDS = ( "extension__organization__ogrn", "extension__organization__okpo", "status", + "payload__attestation_status", ) SOURCE_RECORD_ORDERING_VALUES = [ value @@ -222,8 +224,10 @@ SOURCE_RECORD_LIST_PARAMS = [ "ordering", description=( "Сортировка по полям: record_date, extension__organization__name, " - "created_at, updated_at, title, uid, extension__organization__inn, " - "extension__organization__ogrn. Для обратной сортировки используйте " + "extension__organization__full_name, created_at, updated_at, title, " + "uid, extension__organization__inn, extension__organization__ogrn, " + "extension__organization__okpo, status, " + "payload__attestation_status. Для обратной сортировки используйте " "префикс -. Значения record_date с null сортируются последними." ), enum=SOURCE_RECORD_ORDERING_VALUES, diff --git a/src/settings/base.py b/src/settings/base.py index a59e6d0..2bcc505 100644 --- a/src/settings/base.py +++ b/src/settings/base.py @@ -212,8 +212,6 @@ DATA_UPLOAD_MAX_MEMORY_SIZE = int( # ============================================================================= ZAKUPKI_TOKEN = os.getenv("ZAKUPKI_TOKEN", "") -SUPERJOB_APP_ID = os.getenv("SUPERJOB_APP_ID", "").strip() -HH_USER_AGENT = os.getenv("HH_USER_AGENT", "").strip() FNS_LOCK_TTL_SECONDS = 3600 PARSER_PROXIES = [ item.strip() for item in os.getenv("PARSER_PROXIES", "").split(",") if item.strip() diff --git a/tests/apps/organizations/test_api_v2_source_extensions.py b/tests/apps/organizations/test_api_v2_source_extensions.py index 5e5bf86..a72543b 100644 --- a/tests/apps/organizations/test_api_v2_source_extensions.py +++ b/tests/apps/organizations/test_api_v2_source_extensions.py @@ -6,6 +6,7 @@ from django.core.cache import cache from django.urls import reverse from organizations.models import ( ArbitrationExtension, + ElectronicDocumentExchangeExtension, Organization, OrganizationSourceRecord, PlannedInspectionExtension, @@ -335,6 +336,58 @@ class OrganizationSourceExtensionsApiV2Test(APITestCase): self.assertIsNone(ascending.data["data"][-1]["record_date"]) self.assertIsNone(descending.data["data"][-1]["record_date"]) + def test_flat_gosedo_records_support_frontend_sorting_fields(self): + records = [] + for external_id, name, full_name, attestation_status in ( + ("GOSEDO-BETA", "Организация Бета", "Полное имя Бета", "not_attested"), + ("GOSEDO-ALPHA", "Организация Альфа", "Полное имя Альфа", "attested"), + ): + organization = create_frontend_organization( + name=name, + full_name=full_name, + inn="7707083815" if external_id == "GOSEDO-BETA" else "7707083816", + ogrn=( + "1027700132015" if external_id == "GOSEDO-BETA" else "1027700132016" + ), + ) + extension = ElectronicDocumentExchangeExtension.objects.create( + organization=organization, + title="Электронный документооборот", + ) + records.append( + OrganizationSourceRecord.objects.create( + extension=extension, + record_type="participant", + source="gosedo_address_directory", + external_id=external_id, + payload={"attestation_status": attestation_status}, + ) + ) + + endpoint = reverse("api_v2:organizations:organization-source-records-list") + params = { + "source_group": "electronic_document_exchange", + "source": "gosedo_address_directory", + "record_type": "participant", + } + expected_ascending = [records[1].external_id, records[0].external_id] + expected_descending = list(reversed(expected_ascending)) + + for ordering, expected in ( + ("extension__organization__full_name", expected_ascending), + ("-extension__organization__full_name", expected_descending), + ("payload__attestation_status", expected_ascending), + ("-payload__attestation_status", expected_descending), + ): + with self.subTest(ordering=ordering): + response = self.client.get(endpoint, {**params, "ordering": ordering}) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + self.assertEqual( + [item["external_id"] for item in response.data["data"]], + expected, + ) + def test_flat_source_records_rejects_invalid_date_range_and_ordering(self): response = self.client.get( reverse("api_v2:organizations:organization-source-records-list"), diff --git a/tests/apps/parsers/test_gosedo_media_news.py b/tests/apps/parsers/test_gosedo_media_news.py index ea30bf6..dcadc50 100644 --- a/tests/apps/parsers/test_gosedo_media_news.py +++ b/tests/apps/parsers/test_gosedo_media_news.py @@ -2,9 +2,11 @@ from __future__ import annotations from datetime import datetime, timedelta from io import BytesIO +from tempfile import TemporaryDirectory from types import SimpleNamespace from unittest.mock import patch +from apps.core.models import BackgroundJob, JobStatus from apps.parsers.gosedo import ( GosedoNotModified, GosedoParseResult, @@ -22,6 +24,9 @@ from apps.parsers.media_news import ( ) from apps.parsers.models import ParserLoadLog, ParserSourceArtifact, ParserStagedRecord from apps.parsers.source_artifacts import cleanup_parser_source_artifacts +from apps.parsers.tasks import parse_media_news +from django.core.files.base import ContentFile +from django.core.files.storage import default_storage from django.core.files.uploadedfile import SimpleUploadedFile from django.test import TestCase from django.urls import reverse @@ -519,6 +524,55 @@ class MediaNewsImportTest(TestCase): ) +class MediaNewsTaskTest(TestCase): + def test_uploaded_workbook_task_completes_background_job(self): + Organization.objects.create( + name="АО Задача СМИ", + inn="0012345678", + ogrn="1027700132195", + okpo="00123456", + directory_imported_at=timezone.now(), + ) + workbook = _media_workbook( + [ + [ + "00123456", + "0012345678", + "2026-07-01", + "СМИ", + "https://example.test/news", + "Источник\nЗаголовок\nЛид\nСтрока 4\nПолный текст", + "Положительная", + ] + ] + ) + + with TemporaryDirectory() as media_root, self.settings(MEDIA_ROOT=media_root): + file_path = default_storage.save( + "parser_uploads/media.xlsx", + ContentFile(workbook.getvalue()), + ) + task_result = parse_media_news.apply( + kwargs={ + "file_path": file_path, + "original_name": "media.xlsx", + }, + task_id="media-news-task-test", + ) + + self.assertTrue(task_result.successful()) + self.assertFalse(default_storage.exists(file_path)) + + job = BackgroundJob.objects.get(task_id="media-news-task-test") + self.assertEqual(job.status, JobStatus.SUCCESS) + self.assertEqual(job.progress, 100) + self.assertEqual(task_result.result["published"], 1) + self.assertEqual( + ParserLoadLog.objects.get(source=ParserLoadLog.Source.MEDIA_NEWS).status, + ParserLoadLog.Status.SUCCESS, + ) + + class MediaNewsPermissionsTest(APITestCase): def setUp(self): self.user = UserFactory.create_user() @@ -536,8 +590,11 @@ class MediaNewsPermissionsTest(APITestCase): self.client.force_authenticate(self.admin) with patch( + "apps.parsers.views._save_uploaded_parser_file", + return_value="parser_uploads/media.xlsx", + ), patch( "apps.parsers.tasks.parse_media_news.apply_async", - return_value=SimpleNamespace(id="media-task-1"), + side_effect=lambda **kwargs: SimpleNamespace(id=kwargs["task_id"]), ): response = self.client.post( self.url, @@ -546,8 +603,21 @@ class MediaNewsPermissionsTest(APITestCase): ) self.assertEqual(response.status_code, status.HTTP_202_ACCEPTED) - self.assertEqual(response.data["data"]["task_id"], "media-task-1") - self.assertEqual(response.data["data"]["task_ids"], ["media-task-1"]) + task_id = response.data["data"]["task_id"] + self.assertEqual(response.data["data"]["task_ids"], [task_id]) + + card_response = self.client.get( + reverse( + "api_v1:sources:source-cards-detail", + kwargs={"slug": "media-mentions"}, + ) + ) + self.assertEqual(card_response.status_code, status.HTTP_200_OK) + self.assertEqual(card_response.data["data"]["status"], "in_progress") + self.assertEqual( + card_response.data["data"]["active_tasks"][0]["task_id"], + task_id, + ) def test_gosedo_manual_run_is_admin_only_and_returns_task_ids(self): url = reverse( diff --git a/tests/apps/parsers/test_source_cards_service.py b/tests/apps/parsers/test_source_cards_service.py index 4533b41..3bfdf36 100644 --- a/tests/apps/parsers/test_source_cards_service.py +++ b/tests/apps/parsers/test_source_cards_service.py @@ -441,6 +441,31 @@ class SourceCardServiceDatabaseTest(TestCase): def setUp(self): SourceCardService.clear_cache() + def test_enqueue_vacancy_refresh_reuses_fresh_active_job(self): + existing = BackgroundJob.objects.create( + task_id="active-vacancies", + task_name="apps.parsers.tasks.parse_trudvsem_vacancies", + status=JobStatus.STARTED, + ) + task = MagicMock() + + result = SourceCardService._enqueue_task( + task=task, + task_name="apps.parsers.tasks.parse_trudvsem_vacancies", + requested_by_id=5, + meta={"source_card": "labor-vacancies"}, + kwargs={"requested_by_id": 5}, + ) + + self.assertEqual( + result, + { + "task_id": existing.task_id, + "task_name": existing.task_name, + }, + ) + task.apply_async.assert_not_called() + def test_defense_unreliable_suppliers_counts_unique_generic_organizations(self): _save_source_record( source=ParserLoadLog.Source.UNFAIR_SUPPLIERS, diff --git a/tests/apps/parsers/test_source_registry.py b/tests/apps/parsers/test_source_registry.py index 355875b..b038a85 100644 --- a/tests/apps/parsers/test_source_registry.py +++ b/tests/apps/parsers/test_source_registry.py @@ -2,6 +2,7 @@ from dataclasses import asdict from apps.parsers import tasks from apps.parsers.clients.common.structured import MAX_FILE_SIZE_BYTES +from apps.parsers.serializers import ParserRunRequestSerializer from apps.parsers.source_registry import PARSER_SOURCES from apps.parsers.views import TASKS_BY_NAME from django.conf import settings @@ -53,3 +54,19 @@ class ParserSourceRegistryPresentationTest(SimpleTestCase): ) self.assertNotIn("checko", public_metadata.lower()) + + def test_vacancy_source_metadata_names_only_trudvsem(self): + source = PARSER_SOURCES["trudvsem"] + + self.assertEqual(source.agency, "Работа России") + self.assertEqual(source.parser_strategy, "incremental_trudvsem_api") + self.assertNotIn("HH", source.data_scope) + self.assertNotIn("SuperJob", source.data_scope) + + def test_vacancy_run_request_rejects_disabled_sources(self): + serializer = ParserRunRequestSerializer( + data={"vacancy_sources": ["hh", "superjob"]} + ) + + self.assertFalse(serializer.is_valid()) + self.assertIn("vacancy_sources", serializer.errors) diff --git a/tests/apps/parsers/test_tasks.py b/tests/apps/parsers/test_tasks.py index 0eb70c2..817977a 100644 --- a/tests/apps/parsers/test_tasks.py +++ b/tests/apps/parsers/test_tasks.py @@ -16,6 +16,7 @@ from types import SimpleNamespace from unittest.mock import patch from urllib.parse import urlparse +from apps.core.models import BackgroundJob, JobStatus from apps.core.services import BackgroundJobService from apps.parsers import tasks as parser_tasks from apps.parsers.clients.base import HTTPError @@ -37,7 +38,11 @@ from apps.parsers.models import ( FinancialReport, ParserLoadLog, ) -from apps.parsers.services import FNSReportService, ParserLoadLogService +from apps.parsers.services import ( + FNSReportService, + GenericParserRecordService, + ParserLoadLogService, +) from apps.parsers.tasks import ( FNS_API_DEFAULT_REGISTRY_ORGANIZATION_LIMIT, FNS_API_MAX_REGISTRY_ORGANIZATION_LIMIT, @@ -2869,6 +2874,19 @@ class TaskHelpersTestCase(TestCase): class ParseVacanciesTaskTestCase(TestCase): + def _make_registry_organization(self, *, inn: int, name: str) -> Organization: + organization = OrganizationFactory( + pn_name=name, + mn_inn=inn, + mn_ogrn=inn + 1000000000000, + ) + RegistryMembershipPeriodFactory(organization=organization, ended_at=None) + return organization + + def test_vacancy_task_requeues_after_worker_loss(self): + self.assertTrue(parse_trudvsem_vacancies.acks_late) + self.assertTrue(parse_trudvsem_vacancies.reject_on_worker_lost) + def test_parse_trudvsem_vacancies_without_filters_fetches_active_registry_orgs( self, ): @@ -3052,170 +3070,28 @@ class ParseVacanciesTaskTestCase(TestCase): {"trudvsem:7701000102"}, ) - @override_settings(SUPERJOB_APP_ID="test-superjob-app-id") - def test_parse_trudvsem_vacancies_matches_job_boards_by_employer_name(self): - organization = OrganizationFactory( - pn_name='Общество с ограниченной ответственностью "Ромашка"', - mn_inn=7701000301, - mn_ogrn=1027700000301, - ) - RegistryMembershipPeriodFactory(organization=organization, ended_at=None) - captured_client_kwargs = {} - captured_text_queries = {} + def test_parse_trudvsem_vacancies_rejects_disabled_sources(self): + _ensure_directory_organization(inn="7701000401") - class _Provider: - def __init__(self, source_name, *, supports_company_inn): - self.source_name = source_name - self.supports_company_inn = supports_company_inn + with self.assertRaisesMessage( + ValueError, + "Поддерживается только источник trudvsem", + ): + parse_trudvsem_vacancies( + limit=25, + region_code="1", + text="инженер", + proxies=[], + vacancy_sources=["hh", "superjob"], + ) - def fetch_vacancies(self, **kwargs): - if self.source_name == "trudvsem": - return [ - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id="trudvsem:romashka", - inn=kwargs["company_inn"], - title="Работа России", - payload={"vacancy_source": "trudvsem"}, - ) - ] + self.assertFalse(OrganizationSourceRecord.objects.exists()) - captured_text_queries[self.source_name] = kwargs["text"] - return [ - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id=f"{self.source_name}:romashka", - organisation_name='ООО "Ромашка"', - title=f"{self.source_name} matching vacancy", - payload={"vacancy_source": self.source_name}, - ), - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id=f"{self.source_name}:other", - organisation_name='ООО "Лютик"', - title=f"{self.source_name} unrelated vacancy", - payload={"vacancy_source": self.source_name}, - ), - ] - - class _VacanciesClient: - def __init__(self, **kwargs): - captured_client_kwargs.update(kwargs) - - def __enter__(self): - return self - - def __exit__(self, exc_type, exc_val, exc_tb): - return None - - def fetch_vacancies(self, **kwargs): - return [ - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id="trudvsem:romashka", - inn=kwargs["company_inn"], - title="Работа России", - payload={"vacancy_source": "trudvsem"}, - ) - ] - - def iter_source_clients(self): - return [ - ("trudvsem", _Provider("trudvsem", supports_company_inn=True)), - ("hh", _Provider("hh", supports_company_inn=False)), - ("superjob", _Provider("superjob", supports_company_inn=False)), - ] - - original_client = parser_tasks.VacanciesClient - parser_tasks.VacanciesClient = _VacanciesClient - try: - result = parse_trudvsem_vacancies(limit=50, proxies=[]) - finally: - parser_tasks.VacanciesClient = original_client - - self.assertEqual(result["status"], "success") - self.assertIsNone(captured_client_kwargs["sources"]) - self.assertEqual( - captured_text_queries, - { - "hh": "ромашка", - "superjob": "ромашка", - }, - ) - self.assertEqual(result["saved"], 3) - self.assertEqual( - { - ( - record.external_id, - record.source, - record.registry_organization_id, - ) - for record in OrganizationSourceRecord.objects.all() - }, - { - ("trudvsem:romashka", "trudvsem", organization.id), - ("hh:romashka", "hh", organization.id), - ("superjob:romashka", "superjob", organization.id), - }, - ) - - def test_registry_job_board_matching_fetches_only_first_text_search_page(self): - organization = OrganizationFactory( - pn_name='Общество с ограниченной ответственностью "Ромашка"', - mn_inn=7701000302, - mn_ogrn=1027700000302, - ) - RegistryMembershipPeriodFactory(organization=organization, ended_at=None) - captured_offsets = [] - - class _Provider: - supports_company_inn = False - - def fetch_vacancies(self, **kwargs): - captured_offsets.append(kwargs["offset"]) - return [ - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id=f"hh:romashka:{kwargs['offset']}", - organisation_name='ООО "Ромашка"', - title="HeadHunter", - payload={"vacancy_source": "hh"}, - ) - ] - - class _VacanciesClient: - def __init__(self, **kwargs): - pass - - def __enter__(self): - return self - - def __exit__(self, exc_type, exc_val, exc_tb): - return None - - def iter_source_clients(self): - return [("hh", _Provider())] - - original_client = parser_tasks.VacanciesClient - parser_tasks.VacanciesClient = _VacanciesClient - try: - result = parse_trudvsem_vacancies(limit=1, proxies=[]) - finally: - parser_tasks.VacanciesClient = original_client - - self.assertEqual(result["status"], "success") - self.assertEqual(result["saved"], 1) - self.assertEqual(captured_offsets, [0]) - - @override_settings( - SUPERJOB_APP_ID="test-superjob-app-id", - HH_USER_AGENT="Mostovik/1.0 (ops@example.test)", - ) - def test_parse_trudvsem_vacancies_uses_combined_vacancies_client(self): - for inn in ("7701000401", "7701000402", "7701000403"): - _ensure_directory_organization(inn=inn) + def test_registry_run_persists_a_chunk_and_queues_continuation(self): + self._make_registry_organization(inn=7701000501, name="АО Первый пакет") + self._make_registry_organization(inn=7701000502, name="АО Второй пакет") captured_kwargs = {} - captured_fetch_kwargs = {} + captured_fetches = [] class _VacanciesClient: def __init__(self, **kwargs): @@ -3228,70 +3104,175 @@ class ParseVacanciesTaskTestCase(TestCase): return None def fetch_vacancies(self, **kwargs): - captured_fetch_kwargs.update(kwargs) + captured_fetches.append(kwargs["company_inn"]) return [ GenericParserItem( source=ParserLoadLog.Source.TRUDVSEM, - external_id="trudvsem:1", - inn="7701000401", + external_id=f"trudvsem:{kwargs['company_inn']}", + inn=kwargs["company_inn"], title="Работа России", payload={"vacancy_source": "trudvsem"}, - ), - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id="hh:1", - inn="7701000402", - title="HeadHunter", - payload={"vacancy_source": "hh"}, - ), - GenericParserItem( - source=ParserLoadLog.Source.TRUDVSEM, - external_id="superjob:1", - inn="7701000403", - title="SuperJob", - payload={"vacancy_source": "superjob"}, - ), + ) ] original_client = parser_tasks.VacanciesClient parser_tasks.VacanciesClient = _VacanciesClient try: - result = parse_trudvsem_vacancies( - limit=25, - offset=5, - region_code="1", - text="инженер", - proxies=[], - vacancy_sources=["trudvsem", "hh", "superjob"], - ) + with patch.object( + parser_tasks, + "VACANCY_REGISTRY_ORGANIZATIONS_PER_TASK", + 1, + create=True, + ), patch.object( + parser_tasks.parse_trudvsem_vacancies, + "apply_async", + return_value=SimpleNamespace(id="vacancy-continuation"), + ) as apply_async_mock: + first_result = parse_trudvsem_vacancies(limit=25, proxies=[]) + continuation_kwargs = apply_async_mock.call_args.kwargs["kwargs"] + + self.assertEqual(first_result["status"], "in_progress") + self.assertEqual(captured_fetches, ["7701000501"]) + self.assertEqual(OrganizationSourceRecord.objects.count(), 1) + + job = BackgroundJob.objects.get( + task_name="apps.parsers.tasks.parse_trudvsem_vacancies" + ) + load = ParserLoadLog.objects.get( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=job.meta["batch_id"], + ) + self.assertEqual(job.status, JobStatus.STARTED) + self.assertEqual(job.progress, 50) + self.assertEqual(job.meta["next_offset"], 1) + self.assertEqual(load.status, ParserLoadLog.Status.IN_PROGRESS) + self.assertEqual(load.records_count, 1) + + final_result = parse_trudvsem_vacancies(**continuation_kwargs) finally: parser_tasks.VacanciesClient = original_client - self.assertEqual(result["status"], "success") - self.assertEqual(result["saved"], 3) - self.assertEqual(captured_kwargs["superjob_app_id"], "test-superjob-app-id") - self.assertEqual( - captured_kwargs["hh_user_agent"], - "Mostovik/1.0 (ops@example.test)", + self.assertEqual(final_result["status"], "success") + self.assertEqual(final_result["saved"], 2) + self.assertEqual(captured_fetches, ["7701000501", "7701000502"]) + self.assertEqual(captured_kwargs["sources"], ["trudvsem"]) + self.assertEqual(OrganizationSourceRecord.objects.count(), 2) + job.refresh_from_db() + load.refresh_from_db() + self.assertEqual(job.status, JobStatus.SUCCESS) + self.assertEqual(job.progress, 100) + self.assertEqual(load.status, ParserLoadLog.Status.SUCCESS) + self.assertEqual(load.records_count, 2) + + def test_registry_run_preserves_and_resumes_records_after_revoke(self): + self._make_registry_organization(inn=7701000601, name="АО До остановки") + self._make_registry_organization(inn=7701000602, name="АО После остановки") + captured_fetches = [] + + class _VacanciesClient: + def __init__(self, **kwargs): + pass + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + return None + + def fetch_vacancies(self, **kwargs): + captured_fetches.append(kwargs["company_inn"]) + return [ + GenericParserItem( + source=ParserLoadLog.Source.TRUDVSEM, + external_id=f"trudvsem:{kwargs['company_inn']}", + inn=kwargs["company_inn"], + title="Работа России", + payload={"vacancy_source": "trudvsem"}, + ) + ] + + original_client = parser_tasks.VacanciesClient + original_save_records = GenericParserRecordService.save_records + + def save_then_revoke(records, batch_id, *, source, chunk_size=500): + saved = original_save_records( + records, + batch_id, + source=source, + chunk_size=chunk_size, + ) + job = BackgroundJob.objects.get( + task_name="apps.parsers.tasks.parse_trudvsem_vacancies" + ) + job.revoke() + return saved + + parser_tasks.VacanciesClient = _VacanciesClient + try: + with patch.object( + GenericParserRecordService, + "save_records", + side_effect=save_then_revoke, + ): + stopped_result = parse_trudvsem_vacancies(limit=25, proxies=[]) + + stopped_job = BackgroundJob.objects.get( + task_name="apps.parsers.tasks.parse_trudvsem_vacancies" + ) + stopped_load = ParserLoadLog.objects.get( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=stopped_job.meta["batch_id"], + ) + self.assertEqual(stopped_result["status"], "revoked") + self.assertEqual(stopped_job.meta["next_offset"], 1) + self.assertEqual(stopped_load.status, ParserLoadLog.Status.SKIPPED) + self.assertEqual(stopped_load.records_count, 1) + self.assertEqual(OrganizationSourceRecord.objects.count(), 1) + + resumed_result = parse_trudvsem_vacancies(limit=25, proxies=[]) + finally: + parser_tasks.VacanciesClient = original_client + + self.assertEqual(resumed_result["status"], "success") + self.assertEqual(resumed_result["batch_id"], stopped_result["batch_id"]) + self.assertTrue(resumed_result["resumed"]) + self.assertEqual(captured_fetches, ["7701000601", "7701000602"]) + self.assertEqual(OrganizationSourceRecord.objects.count(), 2) + + def test_registry_run_fails_when_trudvsem_fails_for_every_organization(self): + self._make_registry_organization(inn=7701000701, name="АО Ошибка источника") + + class _VacanciesClient: + def __init__(self, **kwargs): + pass + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + return None + + def fetch_vacancies(self, **kwargs): + raise RuntimeError("external 503") + + original_client = parser_tasks.VacanciesClient + parser_tasks.VacanciesClient = _VacanciesClient + try: + result = parse_trudvsem_vacancies(limit=25, proxies=[]) + finally: + parser_tasks.VacanciesClient = original_client + + job = BackgroundJob.objects.get( + task_name="apps.parsers.tasks.parse_trudvsem_vacancies" ) - self.assertEqual(captured_kwargs["sources"], ["trudvsem", "hh", "superjob"]) - self.assertEqual(captured_fetch_kwargs["limit"], 25) - self.assertEqual(captured_fetch_kwargs["offset"], 5) - self.assertEqual(captured_fetch_kwargs["region_code"], "1") - self.assertEqual(captured_fetch_kwargs["text"], "инженер") - self.assertEqual( - set( - OrganizationSourceRecord.objects.values_list( - "external_id", - "source", - ) - ), - { - ("trudvsem:1", "trudvsem"), - ("hh:1", "hh"), - ("superjob:1", "superjob"), - }, + load = ParserLoadLog.objects.get( + source=ParserLoadLog.Source.TRUDVSEM, + batch_id=job.meta["batch_id"], ) + self.assertEqual(result["status"], "failure") + self.assertEqual(job.status, JobStatus.FAILURE) + self.assertEqual(load.status, ParserLoadLog.Status.FAILED) + self.assertFalse(OrganizationSourceRecord.objects.exists()) class ParserLoadLogServiceTestCase(TestCase): diff --git a/tests/apps/parsers/test_vacancy_clients.py b/tests/apps/parsers/test_vacancy_clients.py index 2b4672f..938ef45 100644 --- a/tests/apps/parsers/test_vacancy_clients.py +++ b/tests/apps/parsers/test_vacancy_clients.py @@ -196,7 +196,7 @@ class SuperJobVacanciesClientTest(SimpleTestCase): class VacanciesClientTest(SimpleTestCase): - def test_fetch_vacancies_combines_enabled_sources(self): + def test_fetch_vacancies_uses_only_trudvsem_by_default(self): class _Provider: def __init__(self, source_name: str): self.source_name = source_name @@ -222,7 +222,7 @@ class VacanciesClientTest(SimpleTestCase): self.assertEqual( [record.payload["vacancy_source"] for record in records], - ["trudvsem", "hh", "superjob"], + ["trudvsem"], ) def test_fetch_vacancies_skips_sources_without_company_inn_support(self): @@ -278,8 +278,11 @@ class VacanciesClientTest(SimpleTestCase): self.assertEqual([record.external_id for record in records], ["trudvsem:1"]) - def test_fetch_vacancies_requires_superjob_app_id_when_explicitly_selected(self): - client = VacanciesClient(sources=["superjob"], superjob_app_id="") + def test_fetch_vacancies_rejects_disabled_sources(self): + client = VacanciesClient(sources=["hh", "superjob"]) - with self.assertRaises(VacanciesClientError): + with self.assertRaisesMessage( + VacanciesClientError, + "Disabled vacancy sources: hh, superjob", + ): client.fetch_vacancies()