From 78e27806b9570036fd0e12e9681a783281698ad6 Mon Sep 17 00:00:00 2001 From: jze9 Date: Thu, 27 Aug 2026 17:40:23 +0500 Subject: [PATCH] =?UTF-8?q?feat(admin):=20=D1=88=D0=BA=D0=B0=D0=BB=D0=B0?= =?UTF-8?q?=20=D0=B7=D0=B0=D0=B3=D1=80=D1=83=D0=B7=D0=BA=D0=B8=20=D0=B8?= =?UTF-8?q?=D1=81=D1=82=D0=BE=D1=87=D0=BD=D0=B8=D0=BA=D0=BE=D0=B2,=20?= =?UTF-8?q?=D0=BE=D1=82=D0=BB=D0=B0=D0=B4=D0=BA=D0=B0=20=D0=B8=20=D0=B7?= =?UTF-8?q?=D0=B0=D0=B3=D1=80=D1=83=D0=B7=D0=BA=D0=B0=20=D1=80=D0=B0=D0=B1?= =?UTF-8?q?=D0=BE=D1=82=20=D0=B2=20=D0=BA=D0=BE=D1=80=D0=BF=D1=83=D1=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Заливка корпуса была чёрным ящиком: у источника только last_status (idle/running/done/error), без «сколько из скольки», без причины падения и без способа остановить начатое. Теперь каждый запуск создаёт строку parse_runs, куда воркер раз в ~2с пишет стадию, счётчики и журнал событий. Админка: - шкала загрузки у каждого источника (0→50% выборка, 50→100% индексация), раскрытая строка — журнал прогона по шагам с таймингами; - «Запустить всё» / «Остановить всё» и остановка по одному источнику (кооперативная отмена: воркер останавливается сам, не рвя запись в базу); - пакетное добавление источников (тип + список тем), тип pmc в форме; - загрузка PDF/DOCX/TXT прямо в базу сравнения (index.ingest_upload); - страница «Отладка»: воркеры Celery и их текущие таски, очереди RabbitMQ, покрытие корпуса эмбеддингами, зависшие и упавшие прогоны, конфиг бэкендов. Защита от краш-лупа по consumer_timeout RabbitMQ (docs/DR-HA.md §6), без неё массовый запуск 170+ источников гарантированно ронял воркер: - PARSER_TIME_BUDGET_S (1500с) — прогон закругляется сам и помечается partial; - worker_prefetch_multiplier=1 — таймаут считается от ДОСТАВКИ сообщения, и с дефолтным префетчем очередь долгих run_parser убивала канал на задачах, которые ещё не начинались. Co-Authored-By: Claude Opus 5 --- README.md | 10 +- docs/ARCHITECTURE.md | 29 +- docs/DR-HA.md | 24 +- docs/INGESTION.md | 31 +- scripts/parsers/arxiv.py | 8 +- scripts/parsers/base.py | 13 +- scripts/parsers/cyberleninka.py | 10 +- scripts/parsers/openalex.py | 7 +- scripts/parsers/pmc.py | 7 +- scripts/parsers/tests/test_progress_cb.py | 68 +++ scripts/run_lint.sh | 6 +- .../api/alembic/versions/005_parse_runs.py | 68 +++ services/api/app/api/admin.py | 541 ++++++++++++++++-- services/api/app/core/progress.py | 26 + services/api/app/models/admin.py | 49 +- services/api/app/schemas/admin.py | 54 +- services/api/tests/test_progress.py | 41 ++ services/frontend/src/api/client.ts | 19 + .../frontend/src/components/ProgressBar.tsx | 47 ++ services/frontend/src/main.tsx | 2 + .../frontend/src/pages/admin/AdminLayout.tsx | 2 + services/frontend/src/pages/admin/Debug.tsx | 255 +++++++++ services/frontend/src/pages/admin/Sources.tsx | 416 ++++++++++++-- services/worker-indexer/app/celery_app.py | 7 + services/worker-indexer/app/config.py | 7 + .../worker-indexer/app/models/__init__.py | 38 +- services/worker-indexer/app/progress.py | 124 ++++ services/worker-indexer/app/tasks/index.py | 273 +++++++-- .../worker-indexer/tests/test_progress.py | 104 ++++ 29 files changed, 2118 insertions(+), 168 deletions(-) create mode 100644 scripts/parsers/tests/test_progress_cb.py create mode 100644 services/api/alembic/versions/005_parse_runs.py create mode 100644 services/api/app/core/progress.py create mode 100644 services/api/tests/test_progress.py create mode 100644 services/frontend/src/components/ProgressBar.tsx create mode 100644 services/frontend/src/pages/admin/Debug.tsx create mode 100644 services/worker-indexer/app/progress.py create mode 100644 services/worker-indexer/tests/test_progress.py diff --git a/README.md b/README.md index 9bc4953..47980e0 100644 --- a/README.md +++ b/README.md @@ -232,9 +232,9 @@ docker compose -f docker-compose.prod.yml --profile observability up -d promethe 1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика). Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini). -2. **Юнит-тесты** — `pytest` по сервисам: 113 тестов на ядро детекции, скоринга, - парсеров, форматирования и OAuth, без внешней инфры (БД/Redis/GPU/Ollama - замоканы либо не нужны). +2. **Юнит-тесты** — `pytest` по сервисам: 130 тестов на ядро детекции, скоринга, + парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры + (БД/Redis/GPU/Ollama замоканы либо не нужны). ```bash make lint # ruff + mypy в изолированном контейнере @@ -262,8 +262,10 @@ make test-one SVC=worker-gost # тесты одного сервиса | Итоговый % плагиата + цитаты | `worker-gpu/app/scoring.py` | 18 | | ГОСТ 7.1 / 7.0.5 | `worker-gost/app/formatters/` | 17 | | Список литературы | `worker-gost/app/bibliography.py` | 7 | -| Парсеры источников (CyberLeninka, PMC) | `scripts/parsers/` | 14 | +| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 17 | | OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 | +| Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 | +| Шкала загрузки источников | `api/app/core/progress.py` | 7 | ## Лицензия diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index f75a959..621a376 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -77,7 +77,11 @@ `faiss_id`. - **fingerprints** — `doc_id`, `hash_value` (BIGINT, Winnowing), `position` — для L1. - **usage_logs** — `user_id`, `action` — учёт лимитов по тарифу. -- **parse_sources** — задания парсеров (админка): тип, query, годы, лимит, статус. +- **parse_sources** — задания парсеров (админка): тип, query, годы, лимит, статус, + `last_run_id` — ссылка на последний прогон. +- **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/ + added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал + событий (JSON). Из них админка рисует шкалу загрузки, см. §12. - **staged_works** — пользовательские загрузки на модерацию перед добавлением в корпус. - **admin_sessions** — одноразовые коды входа в админку. @@ -87,7 +91,7 @@ | Очередь | Задачи | Воркер | |---------|--------|--------| -| `queue.index` | `index.extract_and_check`, `index.add_document`, `index.run_parser`, `index.enrich_full_text` | worker-indexer | +| `queue.index` | `index.extract_and_check`, `index.add_document`, `index.run_parser`, `index.enrich_full_text`, `index.ingest_upload` | worker-indexer | | `queue.gpu` | `gpu.check_plagiarism`, `gpu.embed_documents`, `gpu.search_semantic` | worker-gpu | | `queue.gost` | `gost.format_bibliography` | worker-gost | | `queue.notify` | `notify.send_task_done`, `notify.send_verification` | worker-notifier | @@ -115,7 +119,11 @@ **Наполнение корпуса.** админка/CLI → `index.run_parser` (OpenAlex/arXiv/PMC/ КиберЛенинка, фильтр `is_oa`) → `index.add_document` (дедуп по `ext_id`, fingerprints, MinHash, эмбеддинг) → `index.enrich_full_text` (скачать OA-PDF → -MinIO → переиндексация). +MinIO → переиндексация). Ход заливки пишется в `parse_runs` (§12). +Второй путь наполнения — ручная загрузка файлов админом: api сохраняет их в +MinIO (`corpus-upload/`) → `index.ingest_upload` (извлечь текст → `add_document`, +`source=manual_upload`). Это не проверка на плагиат: файл сразу становится +источником для сравнения, минуя отстойник. ## 7. Детекция плагиата — 4 уровня @@ -203,6 +211,21 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica статуса). Опционально — Prometheus+Grafana (профиль `observability`, метрики API + Flower). - **Бэкапы**: `pg_dump→gzip→MinIO`, cron 03:00, ротация 14. Проверка восстановления — `scripts/ops/pg_restore_verify.sh`. HA/DR — [DR-HA.md](DR-HA.md). +- **Шкала загрузки источников** (админка → «Источники»): каждый запуск создаёт строку + `parse_runs`, воркер пишет туда прогресс раз в ~2с. Шкала считается по формуле + `api/app/core/progress.py`: 0→50% — выборка из источника, 50→100% — индексация в + базу. Раскрытая строка показывает журнал прогона (по шагам, с таймингами). + Кнопки «Запустить всё» / «Остановить всё» — массовый старт и кооперативная отмена + (флаг `cancel_requested`, воркер останавливается сам на ближайшем тике; уже + начатый прогон не рвём посреди записи в базу). +- **Панель отладки** (админка → «Отладка», `GET /api/admin/debug`): один срез — + живые воркеры Celery и что именно они крутят, глубина очередей RabbitMQ + (в т.ч. `unacked` и число потребителей), покрытие корпуса эмбеддингами, + активные и проблемные прогоны, зависшие прогоны (нет heartbeat >10 мин), + упавшие проверки за сутки и текущие бэкенды (`EMBED_BACKEND`/`LLM_BACKEND`/ + `VECTOR_BACKEND`). +- **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) — + защита от краш-лупа по `consumer_timeout` RabbitMQ, см. [DR-HA.md](DR-HA.md) §6. ## 13. Безопасность diff --git a/docs/DR-HA.md b/docs/DR-HA.md index d6cf033..eab7c09 100644 --- a/docs/DR-HA.md +++ b/docs/DR-HA.md @@ -66,8 +66,22 @@ Redis у нас — кэш/rate-limits/LSH-индекс (префикс `antipla (`restart: unless-stopped`) поднимает его заново, недоставленное сообщение передоставляется — и цикл повторяется бесконечно, монополизируя весь пул воркера (было и на `worker-indexer` из-за backoff `scripts/parsers/openalex.py`, и на -`worker-gpu` из-за зависшей Ollama на embedding-gpu). Текущая митигация — держать -retry-бэкоффы заведомо короче 1800с (пример: `MAX_RATE_LIMIT_RETRIES` в -`openalex.py`). Более глубокий фикс — поднять `consumer_timeout` на самом RabbitMQ -(доступа для этого пока не заводили) — актуально при добавлении любых новых -долгих retry-циклов в Celery-тасках. +`worker-gpu` из-за зависшей Ollama на embedding-gpu). + +Митигации (2026-08-27, сделано для `worker-indexer`): + +1. **Бюджет времени прогона** — `PARSER_TIME_BUDGET_S` (1500с) в + `worker-indexer/app/config.py`: `index.run_parser` сам останавливается раньше + дедлайна и помечает прогон `partial` вместо того, чтобы довести воркер до + падения. Остаток дозаливается повторным запуском источника. +2. **`worker_prefetch_multiplier=1`** (`worker-indexer/app/celery_app.py`) — + ключевое при массовой заливке. Таймаут отсчитывается от **доставки** + сообщения, а не от начала выполнения: с дефолтным префетчем (4×concurrency) + сотни поставленных в очередь долгих `run_parser` висят unacked и убивают + канал на задачах, которые ещё даже не начинались. +3. Retry-бэкоффы держим заведомо короче 1800с (`MAX_RATE_LIMIT_RETRIES` в + `openalex.py`). + +`worker-gpu` этих защит пока не имеет (там нет своих долгих retry-циклов, но +есть зависание внешней Ollama — см. выше). Более глубокий общий фикс — поднять +`consumer_timeout` на самом RabbitMQ (доступа для этого пока не заводили). diff --git a/docs/INGESTION.md b/docs/INGESTION.md index 44c85dc..baf0fa3 100644 --- a/docs/INGESTION.md +++ b/docs/INGESTION.md @@ -14,11 +14,9 @@ (русскоязычные, CyberLeninka/OpenAlex `lang=ru`) и **`scripts/seed_broad_corpus.py`** (англоязычные — OpenAlex/arXiv/PMC по широкому списку дисциплин, добавлен позже первой волны). -- Известный операционный риск при массовой докачке: retry-бэкофф на 429 от OpenAlex - не должен по сумме ожиданий превышать `consumer_timeout` RabbitMQ (по умолчанию - 1800с) — иначе брокер рвёт канал до того, как таск успевает сдаться, и - `worker-indexer` уходит в бесконечный краш-луп на редоставленном сообщении - (см. `scripts/parsers/openalex.py`, `MAX_RATE_LIMIT_RETRIES`). +- Операционный риск массовой докачки (`consumer_timeout` RabbitMQ vs долгие таски) + закрыт бюджетом времени прогона и `worker_prefetch_multiplier=1` — подробности + и почему префетч тут главный, см. [DR-HA.md](DR-HA.md) §6. ## Что подготовлено @@ -42,10 +40,31 @@ python scripts/seed_ru_sources.py python scripts/seed_ru_sources.py --apply --limit 500 # 2. Проверить пару источников на темпе/качестве, затем запустить заливку: -# • Админ-панель → «Источники» → «Запустить», ЛИБО +# • Админ-панель → «Источники» → «Запустить» (или «Запустить всё»), ЛИБО # • Celery: index.run_parser.delay(source_id) по каждому id ``` +## Управление заливкой из админки + +Всё, что ниже, доступно на странице «Источники» — CLI для этого больше не нужен: + +- **Шкала загрузки** у каждого источника: стадия (выборка → индексация), сколько + получено из скольки, сколько добавлено/дублей/ошибок. Раскрытая строка — + журнал прогона по шагам с таймингами (таблица `parse_runs`). +- **«Запустить всё»** — прогон по всем включённым источникам; уже идущие + пропускаются. **«Остановить всё»** и остановка по одному — кооперативная + отмена: воркер останавливается сам на ближайшем тике, не обрывая запись в базу. +- **Пакетное добавление** — один тип источника + список тем (по строке), + опционально с немедленным запуском. Заменяет `seed_*.py` для разовых расширений. +- **Загрузка работ в базу** — PDF/DOCX/TXT прямо в корпус сравнения + (`index.ingest_upload`, `source=manual_upload`), минуя проверку и отстойник. +- **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и + упавшие прогоны (см. ARCHITECTURE.md §12). + +Прогон со статусом `partial` — это не ошибка: сработал бюджет времени +(`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout` +RabbitMQ. Остаток добирается повторным запуском источника. + Заливка сама: fetch (rate-limit 1 req/s) → `add_document` (дедуп по `ext_id`, fingerprints L1, MinHash L2) → батч-эмбеддинги `gpu.embed_documents` (L3). ~30 тем × 500 ≈ 15K русских документов на первый заход. diff --git a/scripts/parsers/arxiv.py b/scripts/parsers/arxiv.py index fee3d18..ff5c642 100644 --- a/scripts/parsers/arxiv.py +++ b/scripts/parsers/arxiv.py @@ -12,7 +12,7 @@ import xml.etree.ElementTree as ET from typing import Any import httpx -from base import BaseParser +from base import BaseParser, ProgressCallback logger = logging.getLogger(__name__) @@ -46,6 +46,7 @@ class ArxivParser(BaseParser): limit: int = 500, categories: list[str] | None = None, year_from: int | None = None, + progress_cb: ProgressCallback | None = None, ) -> list[dict[str, Any]]: """ Получить препринты из arXiv. @@ -55,6 +56,7 @@ class ArxivParser(BaseParser): limit: Максимальное количество документов categories: Список категорий arXiv (cs.AI, math.ST и т.д.) year_from: Год публикации от + progress_cb: см. base.ProgressCallback Returns: Список сырых словарей @@ -94,6 +96,10 @@ class ArxivParser(BaseParser): results.extend(entries) start += len(entries) + if progress_cb and not progress_cb(min(len(results), limit)): + logger.info(f"arXiv: выборка остановлена по запросу (получено {len(results)})") + break + if len(entries) < max_results: break diff --git a/scripts/parsers/base.py b/scripts/parsers/base.py index 8e8d22c..2c95058 100644 --- a/scripts/parsers/base.py +++ b/scripts/parsers/base.py @@ -6,11 +6,19 @@ import json import logging from abc import ABC, abstractmethod +from collections.abc import Callable from pathlib import Path from typing import Any logger = logging.getLogger(__name__) +# Колбэк прогресса выборки: вызывается после каждой страницы результатов с +# накопленным количеством документов. Возврат False — просьба остановиться +# (кооперативная отмена из админки или исчерпанный бюджет времени таска, +# см. worker-indexer/app/progress.py). Парсер обязан вернуть уже собранное, +# а не бросать исключение: частичная выборка — валидный результат. +ProgressCallback = Callable[[int], bool] + # Унифицированный формат документа UNIFIED_SCHEMA = { @@ -45,7 +53,10 @@ class BaseParser(ABC): Получить сырые документы из источника. Args: - **kwargs: Специфичные для источника параметры (query, limit и т.д.) + **kwargs: Специфичные для источника параметры (query, limit и т.д.). + Все парсеры принимают ещё и `progress_cb` (ProgressCallback) — + отчёт о прогрессе после каждой страницы и точка кооперативной + остановки. Returns: Список сырых словарей из API источника diff --git a/scripts/parsers/cyberleninka.py b/scripts/parsers/cyberleninka.py index 822ec0d..86181aa 100644 --- a/scripts/parsers/cyberleninka.py +++ b/scripts/parsers/cyberleninka.py @@ -16,7 +16,7 @@ import time from typing import Any import httpx -from base import BaseParser +from base import BaseParser, ProgressCallback logger = logging.getLogger(__name__) @@ -47,6 +47,7 @@ class CyberLeninkaParser(BaseParser): query: str = "", limit: int = 500, subject: str | None = None, + progress_cb: ProgressCallback | None = None, ) -> list[dict[str, Any]]: """ Получить статьи из КиберЛенинки. @@ -55,6 +56,7 @@ class CyberLeninkaParser(BaseParser): query: Поисковый запрос limit: Максимальное количество статей subject: Предметная область (опционально) + progress_cb: см. base.ProgressCallback Returns: Список сырых словарей статей @@ -85,6 +87,12 @@ class CyberLeninkaParser(BaseParser): results.extend(items) page += 1 + if progress_cb and not progress_cb(min(len(results), limit)): + logger.info( + f"КиберЛенинка: выборка остановлена по запросу (получено {len(results)})" + ) + break + if len(items) < 10: break diff --git a/scripts/parsers/openalex.py b/scripts/parsers/openalex.py index fb3169f..7c1edc4 100644 --- a/scripts/parsers/openalex.py +++ b/scripts/parsers/openalex.py @@ -15,7 +15,7 @@ from collections.abc import Generator from typing import Any import httpx -from base import BaseParser +from base import BaseParser, ProgressCallback logger = logging.getLogger(__name__) @@ -45,6 +45,7 @@ class OpenAlexParser(BaseParser): year_to: int | None = None, type_filter: str = "article", open_access_only: bool = False, + progress_cb: ProgressCallback | None = None, ) -> list[dict[str, Any]]: """ Получить документы из OpenAlex. @@ -58,6 +59,7 @@ class OpenAlexParser(BaseParser): type_filter: Тип документа (journal-article, book, и т.д.) open_access_only: только open-access работы (is_oa:true) — резко повышает долю статей с доступным PDF (для наполнения корпуса) + progress_cb: см. base.ProgressCallback Returns: Список сырых словарей из OpenAlex API @@ -74,6 +76,9 @@ class OpenAlexParser(BaseParser): open_access_only=open_access_only, ): results.extend(page) + if progress_cb and not progress_cb(min(len(results), limit)): + logger.info("OpenAlex: выборка остановлена по запросу (получено %d)", len(results)) + break if len(results) >= limit: break diff --git a/scripts/parsers/pmc.py b/scripts/parsers/pmc.py index acd0af8..6faccee 100644 --- a/scripts/parsers/pmc.py +++ b/scripts/parsers/pmc.py @@ -17,7 +17,7 @@ import xml.etree.ElementTree as ET from typing import Any import httpx -from base import BaseParser +from base import BaseParser, ProgressCallback logger = logging.getLogger(__name__) @@ -52,6 +52,7 @@ class PMCParser(BaseParser): limit: int = 500, year_from: int | None = None, year_to: int | None = None, + progress_cb: ProgressCallback | None = None, ) -> list[dict[str, Any]]: """ Получить статьи из PMC Open Access Subset. @@ -61,6 +62,7 @@ class PMCParser(BaseParser): limit: Максимальное количество документов year_from: Год публикации от year_to: Год публикации до + progress_cb: см. base.ProgressCallback Returns: Список сырых словарей (уже с извлечёнными метаданными + full_text) @@ -85,6 +87,9 @@ class PMCParser(BaseParser): ) response.raise_for_status() results.extend(self._parse_articles(response.text)) + if progress_cb and not progress_cb(min(len(results), limit)): + logger.info(f"PMC: выборка остановлена по запросу (получено {len(results)})") + break time.sleep(RATE_LIMIT_DELAY) except httpx.HTTPStatusError as e: logger.error(f"PMC efetch HTTP ошибка: {e.response.status_code}") diff --git a/scripts/parsers/tests/test_progress_cb.py b/scripts/parsers/tests/test_progress_cb.py new file mode 100644 index 0000000..01a5052 --- /dev/null +++ b/scripts/parsers/tests/test_progress_cb.py @@ -0,0 +1,68 @@ +"""Контракт progress_cb у парсеров: отчёт по страницам и остановка по запросу. + +Сеть не трогаем — подменяем HTTP-клиент парсера заглушкой. Проверяем то, на что +опирается заливка: воркер видит рост выборки и может остановить парсер, не +дожидаясь конца (отмена из админки, исчерпанный бюджет времени таска). +""" + +from typing import Any + +from cyberleninka import CyberLeninkaParser + + +class FakeResponse: + def __init__(self, payload: dict[str, Any]) -> None: + self._payload = payload + + def raise_for_status(self) -> None: + pass + + def json(self) -> dict[str, Any]: + return self._payload + + +class FakeClient: + """Отдаёт бесконечные полные страницы по 10 статей — как большая выдача.""" + + def __init__(self) -> None: + self.calls = 0 + + def post(self, url: str, json: dict[str, Any]) -> FakeResponse: + self.calls += 1 + start = json["from"] + return FakeResponse({ + "articles": [{"id": start + i, "name": f"статья {start + i}"} for i in range(10)] + }) + + +def _parser(monkeypatch) -> CyberLeninkaParser: + p = CyberLeninkaParser() + p.client = FakeClient() # type: ignore[assignment] + monkeypatch.setattr("cyberleninka.RATE_LIMIT_DELAY", 0) + return p + + +def test_progress_cb_reports_growing_count(monkeypatch): + p = _parser(monkeypatch) + seen: list[int] = [] + + docs = p.fetch(query="x", limit=30, progress_cb=lambda n: seen.append(n) or True) + + assert len(docs) == 30 + assert seen == [10, 20, 30] + + +def test_progress_cb_false_stops_fetch_early(monkeypatch): + p = _parser(monkeypatch) + + # Останавливаем после второй страницы — как отмена прогона из админки + docs = p.fetch(query="x", limit=1000, progress_cb=lambda n: n < 20) + + assert len(docs) == 20 + assert p.client.calls == 2 # type: ignore[attr-defined] + + +def test_fetch_works_without_callback(monkeypatch): + """Скрипты заливки зовут парсеры без progress_cb — поведение прежнее.""" + p = _parser(monkeypatch) + assert len(p.fetch(query="x", limit=20)) == 20 diff --git a/scripts/run_lint.sh b/scripts/run_lint.sh index 8d64b28..9509467 100755 --- a/scripts/run_lint.sh +++ b/scripts/run_lint.sh @@ -24,10 +24,10 @@ docker run --rm \ echo '▶ ruff (services/ scripts/)' ruff check services/ scripts/ - echo '▶ mypy (чистая логика L1/L2 + ГОСТ + скоринг + OAuth)' - ( cd services/worker-indexer && mypy --config-file /repo/mypy.ini app/algorithms/ app/fragments.py app/staging.py ) + echo '▶ mypy (чистая логика L1/L2 + ГОСТ + скоринг + OAuth + прогресс заливки)' + ( cd services/worker-indexer && mypy --config-file /repo/mypy.ini app/algorithms/ app/fragments.py app/staging.py app/progress.py ) ( cd services/worker-gost && mypy --config-file /repo/mypy.ini app/formatters/ app/bibliography.py ) ( cd services/worker-gpu && mypy --config-file /repo/mypy.ini app/scoring.py ) - ( cd services/api && mypy --config-file /repo/mypy.ini app/core/oauth.py ) + ( cd services/api && mypy --config-file /repo/mypy.ini app/core/oauth.py app/core/progress.py ) " echo "✅ Линт (ruff + mypy) пройден" diff --git a/services/api/alembic/versions/005_parse_runs.py b/services/api/alembic/versions/005_parse_runs.py new file mode 100644 index 0000000..b23ebbe --- /dev/null +++ b/services/api/alembic/versions/005_parse_runs.py @@ -0,0 +1,68 @@ +"""Журнал прогонов парсинга (parse_runs) + ссылка на последний прогон источника. + +Revision ID: 005 +Revises: 004 +Create Date: 2026-08-27 + +Прогресс заливки раньше был не виден: у источника был только last_status +(idle/running/done/error) без «сколько из скольки». Таблица parse_runs хранит +счётчики и структурный журнал каждого прогона — по ним админка рисует шкалу +загрузки и показывает, на чём именно споткнулся источник. +""" + +import sqlalchemy as sa + +from alembic import op + +revision = "005" +down_revision = "004" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "parse_runs", + sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True), + sa.Column( + "source_id", + sa.Integer(), + sa.ForeignKey("parse_sources.id", ondelete="CASCADE"), + nullable=False, + ), + sa.Column("celery_task_id", sa.String(64), nullable=True), + # queued / running / done / partial / error / cancelled + sa.Column("status", sa.String(20), nullable=False, server_default="queued"), + # queued / fetch / index / finished + sa.Column("stage", sa.String(20), nullable=False, server_default="queued"), + sa.Column("target", sa.Integer(), nullable=False, server_default="0"), + sa.Column("fetched", sa.Integer(), nullable=False, server_default="0"), + sa.Column("processed", sa.Integer(), nullable=False, server_default="0"), + sa.Column("added", sa.Integer(), nullable=False, server_default="0"), + sa.Column("duplicates", sa.Integer(), nullable=False, server_default="0"), + # Мусор от источника (нет title/ext_id) — считаем отдельно от ошибок + sa.Column("skipped", sa.Integer(), nullable=False, server_default="0"), + sa.Column("failed", sa.Integer(), nullable=False, server_default="0"), + # Флаг кооперативной отмены: воркер проверяет его на каждом тике прогресса + sa.Column("cancel_requested", sa.Boolean(), nullable=False, server_default=sa.false()), + sa.Column("error", sa.Text(), nullable=True), + sa.Column("log", sa.JSON(), nullable=True), + sa.Column("started_at", sa.DateTime(), server_default=sa.func.now(), nullable=False), + sa.Column("heartbeat_at", sa.DateTime(), nullable=True), + sa.Column("finished_at", sa.DateTime(), nullable=True), + ) + op.create_index("ix_parse_runs_source_id", "parse_runs", ["source_id"]) + op.create_index("ix_parse_runs_status", "parse_runs", ["status"]) + op.create_index("ix_parse_runs_started_at", "parse_runs", ["started_at"]) + + # Последний (он же текущий, пока идёт) прогон источника — чтобы список + # источников отдавался одним джойном, без подзапроса «максимальный id». + op.add_column("parse_sources", sa.Column("last_run_id", sa.Integer(), nullable=True)) + + +def downgrade() -> None: + op.drop_column("parse_sources", "last_run_id") + op.drop_index("ix_parse_runs_started_at", table_name="parse_runs") + op.drop_index("ix_parse_runs_status", table_name="parse_runs") + op.drop_index("ix_parse_runs_source_id", table_name="parse_runs") + op.drop_table("parse_runs") diff --git a/services/api/app/api/admin.py b/services/api/app/api/admin.py index 7f3a11c..0ac53eb 100644 --- a/services/api/app/api/admin.py +++ b/services/api/app/api/admin.py @@ -5,12 +5,17 @@ """ import contextlib +import io import logging -from datetime import datetime +import os +import uuid +from datetime import datetime, timedelta +from pathlib import Path import httpx -from fastapi import APIRouter, Depends, HTTPException, Query, status -from sqlalchemy import delete, func, select +from fastapi import APIRouter, Depends, HTTPException, Query, UploadFile, status +from fastapi.concurrency import run_in_threadpool +from sqlalchemy import delete, func, select, text from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings @@ -19,7 +24,7 @@ from app.core.celery_app import celery_app from app.core.minio_client import get_minio from app.core.redis_client import get_redis from app.database import get_db -from app.models.admin import ParseSource, StagedWork +from app.models.admin import ParseRun, ParseSource, StagedWork from app.models.document import Document, Fingerprint from app.models.task import Task from app.models.user import User @@ -30,6 +35,9 @@ from app.schemas.admin import ( AdminTaskResponse, AdminUserResponse, AdminUserUpdate, + ParseRunDetail, + ParseRunResponse, + ParseSourceBulkCreate, ParseSourceCreate, ParseSourceResponse, ParseSourceUpdate, @@ -42,6 +50,11 @@ logger = logging.getLogger(__name__) router = APIRouter(prefix="/admin", tags=["admin"], dependencies=[Depends(get_admin_user)]) +# Типы источников, для которых в worker-indexer есть парсер (index.run_parser) +SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc") +# Прогон в этих статусах ещё может двигаться — повторно запускать источник нельзя +RUN_ACTIVE_STATUSES = ("queued", "running") + # ═══════════════════════════════════════════════════════════════════════════════ # СЕССИЯ ДОСТУПА @@ -86,6 +99,28 @@ async def verify_session( # ═══════════════════════════════════════════════════════════════════════════════ # ДАШБОРД / СЕРВИСЫ # ═══════════════════════════════════════════════════════════════════════════════ +async def _rabbitmq_queues() -> dict: + """Очереди queue.* из management API RabbitMQ: {имя: сырой словарь очереди}. + + При недоступности брокера/плагина management возвращает {"error": ...} — + это не повод валить весь дашборд. + """ + try: + # amqp://user:pass@host:port/ + creds, _, hostpart = settings.RABBITMQ_URL.split("//", 1)[1].rpartition("@") + rmq_user, _, rmq_pass = creds.partition(":") + rmq_host = hostpart.split(":")[0].split("/")[0] + async with httpx.AsyncClient(timeout=5) as c: + resp = await c.get(f"http://{rmq_host}:15672/api/queues", auth=(rmq_user, rmq_pass)) + if resp.status_code != 200: + return {"error": f"management API вернул {resp.status_code}"} + return { + q["name"]: q for q in resp.json() if str(q.get("name", "")).startswith("queue.") + } + except Exception as e: + return {"error": str(e)[:120]} + + @router.get("/health", response_model=list[ServiceHealth]) async def health(db: AsyncSession = Depends(get_db)) -> list[ServiceHealth]: """Статус всех сервисов инфраструктуры.""" @@ -183,23 +218,12 @@ async def stats(db: AsyncSession = Depends(get_db)) -> AdminStats: storage = {"error": str(e)[:200]} # Длины очередей через RabbitMQ management API (если доступно) - queues: dict = {} - try: - rmq_user = settings.RABBITMQ_URL.split("//")[1].split(":")[0] - rmq_pass = settings.RABBITMQ_URL.split(":")[2].split("@")[0] - rmq_host = settings.RABBITMQ_URL.split("@")[1].split(":")[0] - async with httpx.AsyncClient(timeout=5) as c: - resp = await c.get( - f"http://{rmq_host}:15672/api/queues", - auth=(rmq_user, rmq_pass), - ) - if resp.status_code == 200: - for q in resp.json(): - name = q.get("name", "") - if name.startswith("queue."): - queues[name] = q.get("messages", 0) - except Exception as e: - queues = {"error": str(e)[:120]} + rmq = await _rabbitmq_queues() + queues: dict = ( + {"error": rmq["error"]} + if "error" in rmq + else {name: q.get("messages", 0) for name, q in rmq.items()} + ) return AdminStats( users_total=users_total, @@ -423,15 +447,63 @@ async def storage_info() -> dict: # ═══════════════════════════════════════════════════════════════════════════════ # ИСТОЧНИКИ ПАРСИНГА # ═══════════════════════════════════════════════════════════════════════════════ +async def _attach_runs( + sources: list[ParseSource], db: AsyncSession +) -> list[ParseSourceResponse]: + """Подтянуть последние прогоны одним запросом (иначе N+1 на 170+ источников).""" + run_ids = [s.last_run_id for s in sources if s.last_run_id] + runs: dict[int, ParseRun] = {} + if run_ids: + rows = (await db.execute(select(ParseRun).where(ParseRun.id.in_(run_ids)))).scalars().all() + runs = {r.id: r for r in rows} + + result = [] + for s in sources: + item = ParseSourceResponse.model_validate(s) + run = runs.get(s.last_run_id) if s.last_run_id else None + item.last_run = ParseRunResponse.model_validate(run) if run else None + result.append(item) + return result + + +async def _start_run(src: ParseSource, db: AsyncSession) -> ParseRun: + """Создать строку прогона и отправить таск парсера. + + Строку создаёт API, а не воркер: при массовом запуске 170+ источников + очередь разбирается минутами, и без неё в админке не было бы видно, что + источник уже поставлен в очередь (шкала висела бы на прошлом прогоне). + """ + run = ParseRun(source_id=src.id, status="queued", stage="queued", target=src.limit, log=[]) + db.add(run) + await db.flush() + + src.last_status = "running" + src.last_error = None + src.last_run_at = datetime.utcnow() + src.last_run_id = run.id + await db.commit() + await db.refresh(run) + + celery_result = celery_app.send_task( + "index.run_parser", args=[src.id, run.id], queue="queue.index" + ) + run.celery_task_id = celery_result.id + await db.commit() + await db.refresh(run) + return run + + @router.get("/sources", response_model=list[ParseSourceResponse]) async def list_sources(db: AsyncSession = Depends(get_db)) -> list[ParseSourceResponse]: - rows = (await db.execute(select(ParseSource).order_by(ParseSource.created_at.desc()))).scalars().all() - return [ParseSourceResponse.model_validate(s) for s in rows] + rows = ( + await db.execute(select(ParseSource).order_by(ParseSource.created_at.desc())) + ).scalars().all() + return await _attach_runs(list(rows), db) @router.post("/sources", response_model=ParseSourceResponse, status_code=201) async def create_source(data: ParseSourceCreate, db: AsyncSession = Depends(get_db)) -> ParseSourceResponse: - if data.source_type not in ("openalex", "cyberleninka", "arxiv"): + if data.source_type not in SOURCE_TYPES: raise HTTPException(status_code=400, detail="Недопустимый тип источника") src = ParseSource(**data.model_dump()) db.add(src) @@ -440,6 +512,44 @@ async def create_source(data: ParseSourceCreate, db: AsyncSession = Depends(get_ return ParseSourceResponse.model_validate(src) +@router.post("/sources/bulk", status_code=201) +async def create_sources_bulk( + data: ParseSourceBulkCreate, db: AsyncSession = Depends(get_db) +) -> dict: + """Добавить пачку источников одного типа — по одному на каждый запрос-тему.""" + if data.source_type not in SOURCE_TYPES: + raise HTTPException(status_code=400, detail="Недопустимый тип источника") + + queries = [q.strip() for q in data.queries if q.strip()] + if not queries: + raise HTTPException(status_code=400, detail="Пустой список запросов") + + prefix = data.name_prefix or data.source_type + created: list[ParseSource] = [] + for q in queries: + src = ParseSource( + source_type=data.source_type, + name=f"{prefix}:{q}"[:255], + query=q, + lang=data.lang, + year_from=data.year_from, + year_to=data.year_to, + limit=data.limit, + ) + db.add(src) + created.append(src) + await db.commit() + + started = 0 + if data.run_now: + for src in created: + await db.refresh(src) + await _start_run(src, db) + started += 1 + + return {"created": len(created), "started": started} + + @router.patch("/sources/{source_id}", response_model=ParseSourceResponse) async def update_source( source_id: int, data: ParseSourceUpdate, db: AsyncSession = Depends(get_db) @@ -463,19 +573,386 @@ async def delete_source(source_id: int, db: AsyncSession = Depends(get_db)) -> N await db.commit() -@router.post("/sources/{source_id}/run") -async def run_source(source_id: int, db: AsyncSession = Depends(get_db)) -> dict: +@router.post("/sources/run-all") +async def run_all_sources( + source_type: str | None = None, + db: AsyncSession = Depends(get_db), +) -> dict: + """Запустить заливку по всем включённым источникам (кроме уже идущих).""" + stmt = select(ParseSource).where(ParseSource.enabled.is_(True)) + if source_type: + stmt = stmt.where(ParseSource.source_type == source_type) + sources = (await db.execute(stmt.order_by(ParseSource.id))).scalars().all() + + active_ids = set( + ( + await db.execute( + select(ParseRun.source_id).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES)) + ) + ).scalars().all() + ) + + started, skipped = 0, 0 + for src in sources: + if src.id in active_ids: + skipped += 1 + continue + await _start_run(src, db) + started += 1 + + logger.info("Массовый запуск источников: старт %d, пропущено %d", started, skipped) + return {"started": started, "skipped_active": skipped, "total": len(sources)} + + +@router.post("/sources/stop-all") +async def stop_all_sources(db: AsyncSession = Depends(get_db)) -> dict: + """Отменить все активные прогоны (кооперативно — воркер увидит на своём тике).""" + runs = ( + await db.execute(select(ParseRun).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES))) + ).scalars().all() + for run in runs: + await _cancel_run(run, db) + await db.commit() + return {"cancelled": len(runs)} + + +@router.post("/sources/{source_id}/run", response_model=ParseRunResponse) +async def run_source(source_id: int, db: AsyncSession = Depends(get_db)) -> ParseRunResponse: src = (await db.execute(select(ParseSource).where(ParseSource.id == source_id))).scalar_one_or_none() if src is None: raise HTTPException(status_code=404, detail="Источник не найден") - if src.last_status == "running": + + active = ( + await db.execute( + select(ParseRun).where( + ParseRun.source_id == source_id, ParseRun.status.in_(RUN_ACTIVE_STATUSES) + ) + ) + ).scalars().first() + if active is not None: raise HTTPException(status_code=409, detail="Источник уже парсится") - src.last_status = "running" - src.last_error = None - src.last_run_at = datetime.utcnow() + + run = await _start_run(src, db) + return ParseRunResponse.model_validate(run) + + +async def _cancel_run(run: ParseRun, db: AsyncSession) -> None: + """Отменить прогон: флаг для воркера + revoke, если таск ещё не начали. + + Пока таск в очереди — revoke снимает его, и никто уже не переведёт прогон + из queued, поэтому закрываем строку прямо здесь. Начатый прогон трогать + жёстко нельзя (оборвётся посреди записи в базу) — ставим флаг, воркер + остановится сам на ближайшем тике прогресса. + """ + run.cancel_requested = True + if run.status == "queued": + if run.celery_task_id: + with contextlib.suppress(Exception): + celery_app.control.revoke(run.celery_task_id) + run.status = "cancelled" + run.stage = "finished" + run.finished_at = datetime.utcnow() + src = await db.get(ParseSource, run.source_id) + if src and src.last_run_id == run.id: + src.last_status = "cancelled" + + +@router.post("/sources/{source_id}/cancel") +async def cancel_source_run(source_id: int, db: AsyncSession = Depends(get_db)) -> dict: + runs = ( + await db.execute( + select(ParseRun).where( + ParseRun.source_id == source_id, ParseRun.status.in_(RUN_ACTIVE_STATUSES) + ) + ) + ).scalars().all() + if not runs: + raise HTTPException(status_code=404, detail="Активных прогонов нет") + for run in runs: + await _cancel_run(run, db) await db.commit() - celery_app.send_task("index.run_parser", args=[source_id], queue="queue.index") - return {"status": "started", "source_id": source_id} + return {"cancelled": len(runs), "source_id": source_id} + + +@router.get("/sources/{source_id}/runs", response_model=list[ParseRunResponse]) +async def list_source_runs( + source_id: int, + limit: int = Query(20, le=100), + db: AsyncSession = Depends(get_db), +) -> list[ParseRunResponse]: + rows = ( + await db.execute( + select(ParseRun) + .where(ParseRun.source_id == source_id) + .order_by(ParseRun.started_at.desc()) + .limit(limit) + ) + ).scalars().all() + return [ParseRunResponse.model_validate(r) for r in rows] + + +@router.get("/runs/active", response_model=list[ParseRunResponse]) +async def list_active_runs(db: AsyncSession = Depends(get_db)) -> list[ParseRunResponse]: + """Все прогоны в работе — для панели отладки и общей шкалы заливки.""" + rows = ( + await db.execute( + select(ParseRun) + .where(ParseRun.status.in_(RUN_ACTIVE_STATUSES)) + .order_by(ParseRun.started_at) + ) + ).scalars().all() + return [ParseRunResponse.model_validate(r) for r in rows] + + +@router.get("/runs/{run_id}", response_model=ParseRunDetail) +async def get_run(run_id: int, db: AsyncSession = Depends(get_db)) -> ParseRunDetail: + run = await db.get(ParseRun, run_id) + if run is None: + raise HTTPException(status_code=404, detail="Прогон не найден") + detail = ParseRunDetail.model_validate(run) + detail.log = run.log or [] + return detail + + +# ═══════════════════════════════════════════════════════════════════════════════ +# ЗАГРУЗКА РАБОТ В КОРПУС +# ═══════════════════════════════════════════════════════════════════════════════ +UPLOAD_EXTENSIONS = {".pdf", ".docx", ".txt"} +UPLOAD_MAX_BYTES = 100 * 1024 * 1024 +UPLOAD_CONTENT_TYPES = { + ".pdf": "application/pdf", + ".docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document", + ".txt": "text/plain", +} + + +@router.post("/documents/upload", status_code=202) +async def upload_documents(files: list[UploadFile]) -> dict: + """Залить готовые работы прямо в базу сравнения (минуя проверку и отстойник). + + Отличие от пользовательской загрузки (/documents/check): там файл проверяют + на плагиат, здесь он сам становится источником, с которым будут сравнивать. + """ + if not files: + raise HTTPException(status_code=400, detail="Файлы не переданы") + + accepted: list[dict] = [] + rejected: list[dict] = [] + minio = get_minio() + + for file in files: + name = file.filename or "без имени" + ext = Path(name).suffix.lower() + if ext not in UPLOAD_EXTENSIONS: + rejected.append({"filename": name, "reason": f"формат {ext or '—'} не поддержан"}) + continue + + data = await file.read() + if not data: + rejected.append({"filename": name, "reason": "пустой файл"}) + continue + if len(data) > UPLOAD_MAX_BYTES: + rejected.append({"filename": name, "reason": "больше 100 МБ"}) + continue + + minio_key = f"corpus-upload/{uuid.uuid4()}{ext}" + try: + minio.put_object( + bucket_name=settings.MINIO_BUCKET_DOCS, + object_name=minio_key, + data=io.BytesIO(data), + length=len(data), + content_type=UPLOAD_CONTENT_TYPES.get(ext, "application/octet-stream"), + ) + except Exception as e: + logger.error("Загрузка %s в MinIO не удалась: %s", name, e) + rejected.append({"filename": name, "reason": "ошибка сохранения в хранилище"}) + continue + + celery_app.send_task( + "index.ingest_upload", + args=[minio_key, name], + kwargs={"meta": {"title": Path(name).stem}}, + queue="queue.index", + ) + accepted.append({"filename": name, "minio_key": minio_key, "bytes": len(data)}) + + logger.info("Админ залил в корпус: принято %d, отклонено %d", len(accepted), len(rejected)) + return {"accepted": accepted, "rejected": rejected} + + +# ═══════════════════════════════════════════════════════════════════════════════ +# ОТЛАДКА +# ═══════════════════════════════════════════════════════════════════════════════ +def _celery_snapshot() -> dict: + """Живое состояние воркеров через Celery inspect (блокирующий вызов).""" + try: + insp = celery_app.control.inspect(timeout=4) + ping = insp.ping() or {} + active = insp.active() or {} + reserved = insp.reserved() or {} + stats = insp.stats() or {} + except Exception as e: + return {"error": str(e)[:200], "workers": []} + + workers = [] + for name in sorted(set(ping) | set(stats)): + wstats = stats.get(name, {}) + workers.append({ + "name": name, + "concurrency": (wstats.get("pool") or {}).get("max-concurrency"), + "reserved": len(reserved.get(name, [])), + "active": [ + { + "name": t.get("name"), + "id": t.get("id"), + # args целиком не отдаём: в них бывает текст работы целиком + "args": str(t.get("args"))[:120], + "started_ago_s": ( + round(datetime.utcnow().timestamp() - t["time_start"]) + if t.get("time_start") else None + ), + } + for t in active.get(name, []) + ], + }) + return {"workers": workers} + + +@router.get("/debug") +async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict: + """Полный срез состояния системы для отладки заливки и проверок. + + Один запрос вместо похода по пяти вкладкам: кто из воркеров жив и что + именно сейчас крутит, что копится в очередях, как наполняется корпус, + какие прогоны идут и на чём падали последние. + """ + # Воркеры — блокирующий Celery inspect, в пул потоков, чтобы не вешать луп + celery_state = await run_in_threadpool(_celery_snapshot) + + rmq = await _rabbitmq_queues() + queues = ( + {"error": rmq["error"]} + if "error" in rmq + else { + name: { + "messages": q.get("messages", 0), + "unacked": q.get("messages_unacknowledged", 0), + "consumers": q.get("consumers", 0), + } + for name, q in sorted(rmq.items()) + } + ) + + # ── Корпус ──────────────────────────────────────────────────────────────── + docs_total = (await db.execute(select(func.count()).select_from(Document))).scalar_one() + docs_embedded = ( + await db.execute( + select(func.count()).select_from(Document).where(Document.faiss_id.isnot(None)) + ) + ).scalar_one() + # count(*) по fingerprints — это десятки миллионов строк и секунды ожидания; + # для панели отладки достаточно оценки планировщика + fingerprints_est = ( + await db.execute(text("SELECT reltuples::bigint FROM pg_class WHERE relname='fingerprints'")) + ).scalar() or 0 + + es_docs = None + try: + async with httpx.AsyncClient(timeout=5) as c: + resp = await c.get(f"{settings.ELASTICSEARCH_URL}/documents/_count") + if resp.status_code == 200: + es_docs = resp.json().get("count") + except Exception: + es_docs = None + + # ── Источники и прогоны ─────────────────────────────────────────────────── + src_rows = ( + await db.execute(select(ParseSource.last_status, func.count()).group_by(ParseSource.last_status)) + ).all() + active_runs = ( + await db.execute( + select(ParseRun).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES)).order_by(ParseRun.id) + ) + ).scalars().all() + recent_failed_runs = ( + await db.execute( + select(ParseRun) + .where(ParseRun.status.in_(("error", "partial"))) + .order_by(ParseRun.started_at.desc()) + .limit(10) + ) + ).scalars().all() + + # Зависший прогон: числится в работе, но воркер давно не отчитывался + stale_cutoff = datetime.utcnow() - timedelta(minutes=10) + stale_runs = [ + r.id for r in active_runs + if r.status == "running" and (r.heartbeat_at is None or r.heartbeat_at < stale_cutoff) + ] + + # ── Задачи пользователей ────────────────────────────────────────────────── + day_ago = datetime.utcnow() - timedelta(hours=24) + tasks_24h = dict( + ( + await db.execute( + select(Task.status, func.count()).where(Task.created_at >= day_ago).group_by(Task.status) + ) + ).all() + ) + failed_tasks = ( + await db.execute( + select(Task) + .where(Task.status == "failed") + .order_by(Task.created_at.desc()) + .limit(10) + ) + ).scalars().all() + + return { + "generated_at": datetime.utcnow().isoformat(timespec="seconds"), + "celery": celery_state, + "queues": queues, + "corpus": { + "documents": docs_total, + "documents_embedded": docs_embedded, + "documents_without_embedding": max(docs_total - docs_embedded, 0), + "fingerprints_estimate": int(fingerprints_est), + "elasticsearch_documents": es_docs, + }, + "sources": { + "by_status": dict(src_rows), + "active_runs": [ParseRunResponse.model_validate(r).model_dump() for r in active_runs], + "stale_run_ids": stale_runs, + "recent_problem_runs": [ + ParseRunResponse.model_validate(r).model_dump() for r in recent_failed_runs + ], + }, + "tasks_24h": tasks_24h, + "recent_failed_tasks": [ + { + "public_id": t.public_id, + "type": t.type, + "error": (t.error or "")[:200], + "created_at": t.created_at, + } + for t in failed_tasks + ], + # Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на + # поведение воркеров — при разборе «почему не так считает» нужны первыми + "config": { + "environment": settings.ENVIRONMENT, + "embed_backend": os.environ.get("EMBED_BACKEND", "ollama"), + "llm_backend": os.environ.get("LLM_BACKEND", "ollama"), + "vector_backend": os.environ.get("VECTOR_BACKEND", "faiss"), + "embed_model": os.environ.get("EMBED_MODEL", "—"), + "fetch_full_text": os.environ.get("FETCH_FULL_TEXT", "false"), + "auto_approve_submissions": os.environ.get("AUTO_APPROVE_SUBMISSIONS", "false"), + "parser_time_budget_s": os.environ.get("PARSER_TIME_BUDGET_S", "1500"), + "ollama_url": settings.OLLAMA_URL, + "elasticsearch_url": settings.ELASTICSEARCH_URL, + }, + } # ═══════════════════════════════════════════════════════════════════════════════ diff --git a/services/api/app/core/progress.py b/services/api/app/core/progress.py new file mode 100644 index 0000000..78046b5 --- /dev/null +++ b/services/api/app/core/progress.py @@ -0,0 +1,26 @@ +"""Пересчёт состояния прогона парсинга в шкалу загрузки — чистая логика. + +Стадии прогона несопоставимы по единицам (получено из API источника vs. +записано в базу), а шкала в админке нужна одна. Формула живёт здесь одна на +всех, чтобы не разъезжаться между списком источников, карточкой прогона и +панелью отладки. +""" + +# Прогоны, после которых двигаться уже некуда +TERMINAL_STATUSES = frozenset({"done", "partial", "error", "cancelled"}) + + +def run_percent(status: str, stage: str, target: int, fetched: int, processed: int) -> float: + """Процент готовности: 0→50% — выборка из источника, 50→100% — индексация. + + Завершённый прогон — всегда 100%, каким бы ни был исход: шкала показывает + «работа окончена», а исход — статус рядом с ней. Пока цель неизвестна + (target=0), выборка не даёт прогресса — рисуем 0. + """ + if status in TERMINAL_STATUSES: + return 100.0 + fetch_part = min(fetched / target, 1.0) * 50 if target > 0 else 0.0 + if stage != "index": + return round(fetch_part, 1) + index_part = min(processed / fetched, 1.0) * 50 if fetched > 0 else 0.0 + return round(fetch_part + index_part, 1) diff --git a/services/api/app/models/admin.py b/services/api/app/models/admin.py index 0d0970b..6dabc26 100644 --- a/services/api/app/models/admin.py +++ b/services/api/app/models/admin.py @@ -1,8 +1,8 @@ -"""Модели для админ-панели: источники парсинга, отстойник работ, сессии админа.""" +"""Модели для админ-панели: источники парсинга, прогоны, отстойник, сессии админа.""" from datetime import datetime -from sqlalchemy import JSON, ForeignKey, String, func +from sqlalchemy import JSON, Boolean, ForeignKey, String, Text, func from sqlalchemy.orm import Mapped, mapped_column from app.database import Base @@ -33,12 +33,57 @@ class ParseSource(Base): last_error: Mapped[str | None] = mapped_column(nullable=True) last_run_at: Mapped[datetime | None] = mapped_column(nullable=True) docs_added: Mapped[int] = mapped_column(default=0) + # Последний (пока идёт — текущий) прогон, см. ParseRun + last_run_id: Mapped[int | None] = mapped_column(nullable=True) created_at: Mapped[datetime] = mapped_column(server_default=func.now()) def __repr__(self) -> str: return f"" +class ParseRun(Base): + """Один прогон парсинга источника: счётчики прогресса + журнал событий. + + Пишется воркером (index.run_parser) на каждом тике прогресса, читается + админкой для шкалы загрузки и отладки — по нему видно не только «упало», + но и где именно: на какой стадии, сколько получено/проиндексировано, + сколько дублей и ошибок отдельных документов. + """ + + __tablename__ = "parse_runs" + + id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True) + source_id: Mapped[int] = mapped_column( + ForeignKey("parse_sources.id", ondelete="CASCADE"), nullable=False, index=True + ) + celery_task_id: Mapped[str | None] = mapped_column(String(64), nullable=True) + # queued / running / done / partial / error / cancelled + status: Mapped[str] = mapped_column(String(20), default="queued", index=True) + # queued / fetch / index / finished + stage: Mapped[str] = mapped_column(String(20), default="queued") + # Цель прогона (limit источника) и фактические счётчики + target: Mapped[int] = mapped_column(default=0) + fetched: Mapped[int] = mapped_column(default=0) + processed: Mapped[int] = mapped_column(default=0) + added: Mapped[int] = mapped_column(default=0) + duplicates: Mapped[int] = mapped_column(default=0) + # Мусор от источника (без title/ext_id) — не то же самое, что ошибка + skipped: Mapped[int] = mapped_column(default=0) + failed: Mapped[int] = mapped_column(default=0) + # Кооперативная отмена: воркер сам проверяет флаг между документами + cancel_requested: Mapped[bool] = mapped_column(Boolean, default=False) + error: Mapped[str | None] = mapped_column(Text, nullable=True) + # [{"ts": ISO8601, "level": "info|warning|error", "msg": str}] + log: Mapped[list | None] = mapped_column(JSON, default=list) + started_at: Mapped[datetime] = mapped_column(server_default=func.now(), index=True) + # Последний признак жизни: по нему видно зависший прогон (running, но тишина) + heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True) + finished_at: Mapped[datetime | None] = mapped_column(nullable=True) + + def __repr__(self) -> str: + return f"" + + class StagedWork(Base): """Отстойник: проверенная пользовательская работа, ожидающая решения админа. diff --git a/services/api/app/schemas/admin.py b/services/api/app/schemas/admin.py index 3d5cf92..abbdcc5 100644 --- a/services/api/app/schemas/admin.py +++ b/services/api/app/schemas/admin.py @@ -3,7 +3,9 @@ from datetime import datetime from typing import Any -from pydantic import BaseModel, Field +from pydantic import BaseModel, Field, computed_field + +from app.core.progress import run_percent # ─── Сессия доступа ─────────────────────────────────────────────────────────── @@ -61,9 +63,42 @@ class AdminDocumentResponse(BaseModel): model_config = {"from_attributes": True} +# ─── Прогоны парсинга ───────────────────────────────────────────────────────── +class ParseRunResponse(BaseModel): + """Состояние прогона для шкалы загрузки (без журнала — он в ParseRunDetail).""" + + id: int + source_id: int + status: str + stage: str + target: int + fetched: int + processed: int + added: int + duplicates: int + skipped: int + failed: int + cancel_requested: bool = False + error: str | None = None + started_at: datetime + heartbeat_at: datetime | None = None + finished_at: datetime | None = None + + model_config = {"from_attributes": True} + + @computed_field # type: ignore[prop-decorator] + @property + def percent(self) -> float: + return run_percent(self.status, self.stage, self.target, self.fetched, self.processed) + + +class ParseRunDetail(ParseRunResponse): + log: list[dict[str, Any]] = Field(default_factory=list) + + # ─── Источники парсинга ─────────────────────────────────────────────────────── class ParseSourceCreate(BaseModel): - source_type: str = Field(description="openalex / cyberleninka / arxiv") + source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") name: str = Field(min_length=1, max_length=255) query: str | None = None lang: str | None = None @@ -98,10 +133,25 @@ class ParseSourceResponse(BaseModel): last_run_at: datetime | None = None docs_added: int created_at: datetime + # Последний (или текущий) прогон — источник данных для шкалы в списке + last_run: ParseRunResponse | None = None model_config = {"from_attributes": True} +class ParseSourceBulkCreate(BaseModel): + """Пакетное добавление: один тип/лимит, много запросов (по строке на тему).""" + + source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") + queries: list[str] = Field(min_length=1, max_length=500) + name_prefix: str = "" + lang: str | None = None + year_from: int | None = None + year_to: int | None = None + limit: int = Field(default=1000, ge=1, le=100000) + run_now: bool = False + + # ─── Отстойник ──────────────────────────────────────────────────────────────── class StagedWorkResponse(BaseModel): id: int diff --git a/services/api/tests/test_progress.py b/services/api/tests/test_progress.py new file mode 100644 index 0000000..d14d3ed --- /dev/null +++ b/services/api/tests/test_progress.py @@ -0,0 +1,41 @@ +"""Юнит-тесты шкалы загрузки источников (app.core.progress) — чистая логика.""" + +from app.core.progress import run_percent + + +def test_queued_run_shows_nothing_done(): + assert run_percent("queued", "queued", target=100, fetched=0, processed=0) == 0.0 + + +def test_fetch_stage_fills_first_half(): + assert run_percent("running", "fetch", target=100, fetched=50, processed=0) == 25.0 + assert run_percent("running", "fetch", target=100, fetched=100, processed=0) == 50.0 + + +def test_index_stage_fills_second_half(): + # Выборка закончена (50%), проиндексирована половина полученного → 75% + assert run_percent("running", "index", target=100, fetched=100, processed=50) == 75.0 + assert run_percent("running", "index", target=100, fetched=100, processed=100) == 100.0 + + +def test_partial_fetch_does_not_inflate_index_progress(): + """Источник отдал меньше лимита: выборка даёт свои 20%, индексация — свои.""" + percent = run_percent("running", "index", target=100, fetched=40, processed=20) + assert percent == 45.0 # 20% за выборку + 25% за половину индексации + + +def test_unknown_target_gives_zero_instead_of_division_error(): + assert run_percent("running", "fetch", target=0, fetched=17, processed=0) == 0.0 + assert run_percent("running", "index", target=0, fetched=0, processed=0) == 0.0 + + +def test_overshoot_is_clamped(): + """Источник может вернуть больше лимита — шкала не должна уезжать за 100%.""" + assert run_percent("running", "fetch", target=10, fetched=99, processed=0) == 50.0 + assert run_percent("running", "index", target=10, fetched=10, processed=99) == 100.0 + + +def test_finished_runs_are_always_full(): + """Шкала показывает «работа окончена», исход виден по статусу рядом.""" + for status in ("done", "partial", "error", "cancelled"): + assert run_percent(status, "finished", target=100, fetched=3, processed=1) == 100.0 diff --git a/services/frontend/src/api/client.ts b/services/frontend/src/api/client.ts index a09f9d3..e5eb9e0 100644 --- a/services/frontend/src/api/client.ts +++ b/services/frontend/src/api/client.ts @@ -114,9 +114,28 @@ export const adminApi = { sources: () => api.get('/admin/sources'), createSource: (data: Record) => api.post('/admin/sources', data), + createSourcesBulk: (data: Record) => api.post('/admin/sources/bulk', data), updateSource: (id: number, data: Record) => api.patch(`/admin/sources/${id}`, data), deleteSource: (id: number) => api.delete(`/admin/sources/${id}`), runSource: (id: number) => api.post(`/admin/sources/${id}/run`), + cancelSource: (id: number) => api.post(`/admin/sources/${id}/cancel`), + runAllSources: (sourceType?: string) => + api.post('/admin/sources/run-all', null, { params: sourceType ? { source_type: sourceType } : undefined }), + stopAllSources: () => api.post('/admin/sources/stop-all'), + sourceRuns: (id: number) => api.get(`/admin/sources/${id}/runs`), + activeRuns: () => api.get('/admin/runs/active'), + run: (runId: number) => api.get(`/admin/runs/${runId}`), + + debug: () => api.get('/admin/debug'), + + uploadDocuments: (files: File[]) => { + const form = new FormData(); + files.forEach((f) => form.append('files', f)); + return api.post('/admin/documents/upload', form, { + headers: { 'Content-Type': 'multipart/form-data' }, + timeout: 300000, // пачка файлов грузится дольше одиночной проверки + }); + }, staging: (params?: { status?: string; limit?: number; offset?: number }) => api.get('/admin/staging', { params }), diff --git a/services/frontend/src/components/ProgressBar.tsx b/services/frontend/src/components/ProgressBar.tsx new file mode 100644 index 0000000..7a08fa3 --- /dev/null +++ b/services/frontend/src/components/ProgressBar.tsx @@ -0,0 +1,47 @@ +import { clsx } from 'clsx'; + +/** Цвет шкалы = исход прогона: серый пока ждёт, синий в работе, дальше по итогу. */ +const BAR_COLORS: Record = { + queued: 'bg-gray-300', + running: 'bg-brand-500', + done: 'bg-emerald-500', + partial: 'bg-amber-500', + cancelled: 'bg-gray-400', + error: 'bg-red-500', +}; + +interface ProgressBarProps { + percent: number; + status?: string; + /** Подпись слева под шкалой (что именно сейчас происходит) */ + label?: string; + /** Показывать процент справа */ + showPercent?: boolean; + className?: string; +} + +export function ProgressBar({ + percent, + status = 'running', + label, + showPercent = true, + className, +}: ProgressBarProps) { + const value = Math.max(0, Math.min(100, percent)); + return ( +
+
+
+
+ {(label || showPercent) && ( +
+ {label} + {showPercent && {value.toFixed(0)}%} +
+ )} +
+ ); +} diff --git a/services/frontend/src/main.tsx b/services/frontend/src/main.tsx index 85fa21f..2d5a90a 100644 --- a/services/frontend/src/main.tsx +++ b/services/frontend/src/main.tsx @@ -25,6 +25,7 @@ import { Documents as AdminDocuments } from './pages/admin/Documents'; import { Storage as AdminStorage } from './pages/admin/Storage'; import { Sources as AdminSources } from './pages/admin/Sources'; import { Staging as AdminStaging } from './pages/admin/Staging'; +import { Debug as AdminDebug } from './pages/admin/Debug'; import './index.css'; @@ -77,6 +78,7 @@ ReactDOM.createRoot(document.getElementById('root')!).render( } /> } /> } /> + } /> } /> diff --git a/services/frontend/src/pages/admin/AdminLayout.tsx b/services/frontend/src/pages/admin/AdminLayout.tsx index 546444f..3a13891 100644 --- a/services/frontend/src/pages/admin/AdminLayout.tsx +++ b/services/frontend/src/pages/admin/AdminLayout.tsx @@ -3,6 +3,7 @@ import { NavLink, useParams, useNavigate, Navigate } from 'react-router-dom'; import { useQuery } from '@tanstack/react-query'; import { LayoutDashboard, Users, FileText, Database, HardDrive, Download, Inbox, LogOut, ShieldAlert, + Bug, } from 'lucide-react'; import { clsx } from 'clsx'; import { adminApi } from '../../api/client'; @@ -16,6 +17,7 @@ const NAV = [ { to: 'storage', label: 'Хранилище', icon: HardDrive }, { to: 'sources', label: 'Источники', icon: Download }, { to: 'staging', label: 'Отстойник', icon: Inbox }, + { to: 'debug', label: 'Отладка', icon: Bug }, ]; interface AdminLayoutProps { diff --git a/services/frontend/src/pages/admin/Debug.tsx b/services/frontend/src/pages/admin/Debug.tsx new file mode 100644 index 0000000..20dabf8 --- /dev/null +++ b/services/frontend/src/pages/admin/Debug.tsx @@ -0,0 +1,255 @@ +import React from 'react'; +import { useQuery } from '@tanstack/react-query'; +import { + Activity, AlertTriangle, Cpu, Database, Layers, RefreshCw, Settings2, ListTree, +} from 'lucide-react'; +import { adminApi } from '../../api/client'; +import { ProgressBar } from '../../components/ProgressBar'; + +interface ActiveTask { name: string; id: string; args: string; started_ago_s: number | null } +interface Worker { name: string; concurrency: number | null; reserved: number; active: ActiveTask[] } +interface Run { + id: number; source_id: number; status: string; stage: string; target: number; + fetched: number; processed: number; added: number; duplicates: number; + skipped: number; failed: number; error: string | null; percent: number; + started_at: string; heartbeat_at: string | null; +} + +interface DebugData { + generated_at: string; + celery: { workers: Worker[]; error?: string }; + queues: Record | { error: string }; + corpus: { + documents: number; documents_embedded: number; documents_without_embedding: number; + fingerprints_estimate: number; elasticsearch_documents: number | null; + }; + sources: { + by_status: Record; + active_runs: Run[]; + stale_run_ids: number[]; + recent_problem_runs: Run[]; + }; + tasks_24h: Record; + recent_failed_tasks: { public_id: string; type: string; error: string; created_at: string }[]; + config: Record; +} + +const STAGE_LABELS: Record = { + queued: 'в очереди', fetch: 'выборка', index: 'индексация', finished: 'завершено', +}; + +function fmtNum(n: number | null | undefined): string { + return n == null ? '—' : n.toLocaleString('ru-RU'); +} + +export function Debug() { + const { data, isFetching, error } = useQuery({ + queryKey: ['admin-debug'], + queryFn: () => adminApi.debug().then((r) => r.data as DebugData), + refetchInterval: 5000, + }); + + if (error) { + return
Не удалось получить срез состояния: {String(error)}
; + } + if (!data) { + return
Сбор данных…
; + } + + const queuesErr = 'error' in data.queues ? (data.queues as { error: string }).error : null; + const queues = queuesErr ? {} : (data.queues as Record); + const embedPercent = data.corpus.documents + ? (data.corpus.documents_embedded / data.corpus.documents) * 100 + : 0; + + return ( +
+
+

+ Отладка + {isFetching && } +

+ срез от {new Date(data.generated_at).toLocaleTimeString('ru-RU')} +
+ + {/* Активные прогоны заливки */} + + {data.sources.stale_run_ids.length > 0 && ( +
+ + + Прогоны без признаков жизни больше 10 минут: {data.sources.stale_run_ids.join(', ')}. + Обычно это упавший или перезапущенный worker-indexer — проверьте его логи и очередь queue.index. + +
+ )} + {data.sources.active_runs.length ? ( +
+ {data.sources.active_runs.map((r) => ( +
+
+ #{r.id} · источник {r.source_id} +
+ +
+ ))} +
+ ) : ( +

Сейчас ничего не заливается.

+ )} +
+ {Object.entries(data.sources.by_status).map(([s, c]) => ( + {s}: {c} + ))} +
+
+ + {/* Воркеры */} + + {data.celery.error &&

{data.celery.error}

} + {data.celery.workers.length ? ( +
+ {data.celery.workers.map((w) => ( +
+
+ {w.name} + + параллельно: {w.concurrency ?? '—'} · в работе: {w.active.length} · зарезервировано: {w.reserved} + +
+ {w.active.length ? ( +
+ {w.active.map((t) => ( +
+ + {t.started_ago_s != null ? `${t.started_ago_s}с` : '—'} + + {t.name} + {t.args} +
+ ))} +
+ ) :

простаивает

} +
+ ))} +
+ ) :

Воркеры не отвечают на ping.

} +
+ + {/* Очереди */} + + {queuesErr ? ( +

{queuesErr}

+ ) : ( + + + + + + {Object.entries(queues).map(([name, q]) => ( + + + + + + + ))} + +
ОчередьЖдутВ работеПотребителей
{name}{fmtNum(q.messages)}{fmtNum(q.unacked)}{q.consumers}
+ )} +
+ + {/* Корпус */} + +
+ + + + +
+ 0 ? 'partial' : 'done'} + label={`покрытие эмбеддингами (L3): без вектора ${fmtNum(data.corpus.documents_without_embedding)} документов`} + /> +
+ + {/* Проблемные прогоны */} + {data.sources.recent_problem_runs.length > 0 && ( + +
+ {data.sources.recent_problem_runs.map((r) => ( +
+ #{r.id} + источник {r.source_id} + {r.status} + + получено {r.fetched}/{r.target}, добавлено {r.added}, ошибок {r.failed} + + {r.error && {r.error}} +
+ ))} +
+
+ )} + + {/* Задачи пользователей */} + +
+ {Object.entries(data.tasks_24h).map(([s, c]) => ( + {s}: {c} + ))} + {!Object.keys(data.tasks_24h).length && задач не было} +
+ {data.recent_failed_tasks.length > 0 && ( +
+
Последние упавшие:
+ {data.recent_failed_tasks.map((t) => ( +
+ {new Date(t.created_at).toLocaleString('ru-RU')} + {t.type} + {t.error || '—'} +
+ ))} +
+ )} +
+ + {/* Конфигурация */} + +
+ {Object.entries(data.config).map(([k, v]) => ( +
+ {k} + {v} +
+ ))} +
+
+
+ ); +} + +function Card({ icon: Icon, title, children }: { icon: React.ElementType; title: string; children: React.ReactNode }) { + return ( +
+

+ {title} +

+ {children} +
+ ); +} + +function Metric({ label, value }: { label: string; value: string }) { + return ( +
+
{value}
+
{label}
+
+ ); +} diff --git a/services/frontend/src/pages/admin/Sources.tsx b/services/frontend/src/pages/admin/Sources.tsx index 66db6e2..21acd65 100644 --- a/services/frontend/src/pages/admin/Sources.tsx +++ b/services/frontend/src/pages/admin/Sources.tsx @@ -1,78 +1,290 @@ -import { useState } from 'react'; +import { Fragment, useRef, useState } from 'react'; import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'; -import { Play, Trash2, Plus } from 'lucide-react'; +import { + Play, Trash2, Plus, Square, PlayCircle, StopCircle, ChevronDown, ChevronRight, + Upload, Layers, RefreshCw, +} from 'lucide-react'; import toast from 'react-hot-toast'; import { adminApi } from '../../api/client'; +import { ProgressBar } from '../../components/ProgressBar'; + +interface Run { + id: number; source_id: number; status: string; stage: string; + target: number; fetched: number; processed: number; + added: number; duplicates: number; skipped: number; failed: number; + cancel_requested: boolean; error: string | null; + started_at: string; heartbeat_at: string | null; finished_at: string | null; + percent: number; +} + +interface LogEntry { ts: string; elapsed: number; level: string; msg: string } +interface RunDetail extends Run { log: LogEntry[] } interface Source { id: number; source_type: string; name: string; query: string | null; lang: string | null; year_from: number | null; year_to: number | null; limit: number; enabled: boolean; last_status: string; last_error: string | null; - last_run_at: string | null; docs_added: number; + last_run_at: string | null; docs_added: number; last_run: Run | null; } +const SOURCE_TYPES = [ + { value: 'openalex', label: 'OpenAlex' }, + { value: 'cyberleninka', label: 'КиберЛенинка' }, + { value: 'arxiv', label: 'arXiv' }, + { value: 'pmc', label: 'PubMed Central' }, +]; + const STATUS_COLORS: Record = { idle: 'bg-gray-100 text-gray-600', + queued: 'bg-gray-100 text-gray-600', running: 'bg-blue-100 text-blue-700', done: 'bg-emerald-100 text-emerald-700', + partial: 'bg-amber-100 text-amber-700', + cancelled: 'bg-gray-200 text-gray-600', error: 'bg-red-100 text-red-700', }; +const STAGE_LABELS: Record = { + queued: 'в очереди', + fetch: 'выборка из источника', + index: 'индексация в базу', + finished: 'завершено', +}; + +const ACTIVE = ['queued', 'running']; + +/** Что происходит с источником прямо сейчас — подпись под шкалой. */ +function runLabel(run: Run): string { + if (run.stage === 'fetch') return `${STAGE_LABELS.fetch}: ${run.fetched}/${run.target}`; + if (run.stage === 'index') return `${STAGE_LABELS.index}: ${run.processed}/${run.fetched}`; + return `+${run.added} новых · ${run.duplicates} дублей${run.failed ? ` · ${run.failed} ошибок` : ''}`; +} + export function Sources() { const qc = useQueryClient(); - const [form, setForm] = useState({ source_type: 'openalex', name: '', query: '', lang: '', limit: 100 }); - const { data: sources } = useQuery({ + const [form, setForm] = useState({ source_type: 'openalex', name: '', query: '', lang: '', limit: 1000 }); + const [bulk, setBulk] = useState({ open: false, source_type: 'openalex', queries: '', lang: '', limit: 1000, run_now: false }); + const [expanded, setExpanded] = useState(null); + const fileInput = useRef(null); + + const { data: sources, isFetching } = useQuery({ queryKey: ['admin-sources'], queryFn: () => adminApi.sources().then((r) => r.data as Source[]), - refetchInterval: 5000, + refetchInterval: 3000, }); + const invalidate = () => qc.invalidateQueries({ queryKey: ['admin-sources'] }); + const fail = (e: any) => toast.error(e.response?.data?.detail || 'Ошибка'); + const create = useMutation({ mutationFn: (data: Record) => adminApi.createSource(data), - onSuccess: () => { qc.invalidateQueries({ queryKey: ['admin-sources'] }); toast.success('Источник добавлен'); setForm({ source_type: 'openalex', name: '', query: '', lang: '', limit: 100 }); }, - onError: (e: any) => toast.error(e.response?.data?.detail || 'Ошибка'), + onSuccess: () => { + invalidate(); + toast.success('Источник добавлен'); + setForm({ source_type: form.source_type, name: '', query: '', lang: '', limit: form.limit }); + }, + onError: fail, + }); + const createBulk = useMutation({ + mutationFn: (data: Record) => adminApi.createSourcesBulk(data), + onSuccess: (r) => { + invalidate(); + toast.success(`Добавлено источников: ${r.data.created}${r.data.started ? `, запущено ${r.data.started}` : ''}`); + setBulk({ ...bulk, queries: '' }); + }, + onError: fail, }); const run = useMutation({ mutationFn: (id: number) => adminApi.runSource(id), - onSuccess: () => { qc.invalidateQueries({ queryKey: ['admin-sources'] }); toast.success('Парсинг запущен'); }, - onError: (e: any) => toast.error(e.response?.data?.detail || 'Ошибка'), + onSuccess: () => { invalidate(); toast.success('Парсинг запущен'); }, + onError: fail, + }); + const cancel = useMutation({ + mutationFn: (id: number) => adminApi.cancelSource(id), + onSuccess: () => { invalidate(); toast.success('Остановка запрошена'); }, + onError: fail, + }); + const runAll = useMutation({ + mutationFn: () => adminApi.runAllSources(), + onSuccess: (r) => { + invalidate(); + toast.success(`Запущено ${r.data.started} из ${r.data.total} (уже шли: ${r.data.skipped_active})`); + }, + onError: fail, + }); + const stopAll = useMutation({ + mutationFn: () => adminApi.stopAllSources(), + onSuccess: (r) => { invalidate(); toast.success(`Остановка ${r.data.cancelled} прогонов запрошена`); }, + onError: fail, }); const del = useMutation({ mutationFn: (id: number) => adminApi.deleteSource(id), - onSuccess: () => { qc.invalidateQueries({ queryKey: ['admin-sources'] }); toast.success('Удалён'); }, + onSuccess: () => { invalidate(); toast.success('Удалён'); }, + onError: fail, }); + const upload = useMutation({ + mutationFn: (files: File[]) => adminApi.uploadDocuments(files), + onSuccess: (r) => { + const { accepted, rejected } = r.data; + toast.success(`Принято файлов: ${accepted.length}${rejected.length ? `, отклонено ${rejected.length}` : ''}`); + rejected.forEach((x: { filename: string; reason: string }) => toast.error(`${x.filename}: ${x.reason}`)); + }, + onError: fail, + }); + + const list = sources || []; + const active = list.filter((s) => s.last_run && ACTIVE.includes(s.last_run.status)); + const activeAdded = active.reduce((sum, s) => sum + (s.last_run?.added || 0), 0); + // Общая шкала = средняя готовность идущих прогонов: столько заливки осталось + const overall = active.length + ? active.reduce((sum, s) => sum + (s.last_run?.percent || 0), 0) / active.length + : 0; return (
-

Источники парсинга

+
+

+ Источники парсинга + {isFetching && } +

+
+ + +
+
+ + {/* Общий прогресс заливки */} +
+
+

Заливка корпуса

+ + активных источников: {active.length} из {list.length} + {active.length > 0 && <> · добавлено в этом заходе: {activeAdded}} + +
+ +
{/* Форма добавления */}
-

Добавить источник

-
- - setForm({ ...form, name: e.target.value })} placeholder="Название" - className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> - setForm({ ...form, query: e.target.value })} placeholder="Запрос" - className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> - setForm({ ...form, lang: e.target.value })} placeholder="Язык (ru/en)" - className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> - setForm({ ...form, limit: Number(e.target.value) })} placeholder="Лимит" - className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> +
+

Добавить источник

+
- + + ) : ( + <> +
+ + setBulk({ ...bulk, lang: e.target.value })} placeholder="Язык (ru/en)" + className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> + setBulk({ ...bulk, limit: Number(e.target.value) })} placeholder="Лимит на тему" + className="border border-gray-200 rounded-lg px-3 py-2 text-sm" /> + +
+