Compare commits

..

6 Commits

Author SHA1 Message Date
jze9
2d38dfd1f4 feat(ops): путь к миллионам статей — массовая заливка из бакета PMC (не закончен)
All checks were successful
Deploy / test (push) Successful in 3m12s
Deploy / deploy (push) Successful in 1m56s
Разведка под задачу «нужны миллионы»: постраничные API дают 1-2 статьи в
секунду и упираются в rate limit, а OpenAlex-дамп содержит только метаданные —
миллионы аннотаций бесполезны, что уже доказано на КиберЛенинке (30 отпечатков
на документ, 0 глубоко проиндексированных).

Найден и проверен рабочий источник: AWS Open Data бакет pmc-oa-opendata открыт
без ключа, у каждой статьи лежит готовый извлечённый текст (33-130 КБ) и json с
метаданными. Замер: 12 статей/с в 10 потоков — миллион за сутки.

Замеры узких мест на проде (для планирования масштаба):
- winnowing 95 док/с на ядро — не ограничение;
- поиск L1 31 мс по 500 отпечаткам при 113 млн строк;
- вставка отпечатков построчно 6.7 тыс. строк/с → миллион статей это 2.5 млрд
  строк и четверо суток только на запись, поэтому в скрипте COPY;
- объём: 106 байт на строку → ~400 ГБ на миллион статей с индексами.

Скрипт НЕ закончен: документы вставляются верно, но шаг отпечатков ошибочно
пропускает свежие документы (пробные 200 записей удалены из базы). Ограничение
задокументировано прямо в скрипте, чтобы его не запустили на объёме.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:58:17 +05:00
jze9
d7c004af76 feat(ops): углубление индексации КиберЛенинки + общий store_full_text
Проверка глубины вскрыла главную слабость корпуса: 99 883 документа
КиберЛенинки (56% базы) имеют в среднем 30 отпечатков — заголовок с
аннотацией. Списывание из тела русской статьи L1 не находит, хотя сервис
рассчитан именно на русских студентов. У PMC и arXiv глубина 97% и 94%.

Отдельно: мерить глубину по minio_key оказалось неверно — PMC кладёт тело
статьи в отпечатки при заливке, не сохраняя файл, поэтому «27.5% с полным
текстом» занижало картину для одних источников и скрывало провал у другой.
Честный показатель — число отпечатков на документ, по нему всё и пересчитано.

- backfill_cyberleninka_pdf.py: PDF берётся прямым адресом {url}/pdf (докачка
  по обычному url бесполезна — там HTML, замер 0 из 8). Проверено на 50:
  углублено 44, в среднем 23 тыс. символов, глубина 30 → 1453 отпечатка;
- backfill_pmc_fulltext.py: сохранение тела статьи из API в MinIO (отпечатки
  там уже глубокие — скрипт нужен для отчётов и подсветки, не для детекции);
- store_full_text вынесен из index.enrich_full_text: один путь «текст получен →
  MinIO + пересчёт L1 + обновление L2» для задачи докачки и для бэкфиллов.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:45:48 +05:00
jze9
d8630ca0f9 fix(admin): покрытие L3 мерить по индексу, а не по отметкам в базе
Проверка на живых данных вскрыла, что панель отладки врала в мою же пользу:
показывала 154 046 «документов с эмбеддингом» (87%), тогда как в FAISS реально
93 053 вектора (52.5%). Колонка documents.faiss_id для этого непригодна —
отметка остаётся после пересоздания индекса (смена модели: 768 → 1024) и после
сбоев worker-gpu. Выборочная проверка: у 8 из 20 «отмеченных» вектора нет.

- gpu.index_stats — новая задача, отдаёт реальное содержимое активного
  векторного бэкенда; API спрашивает её для панели отладки;
- панель показывает «векторов в индексе» и отдельно предупреждает о ложных
  отметках (их 60 993), потому что такие документы молча выпадают из L3:
  в индексе их нет, а на пересчёт они не попадут — reembed ищет faiss_id IS NULL;
- scripts/ops/faiss_reconcile.py — сверяет отметки с индексом и обнуляет ложные,
  после чего reembed_missing.py отправляет их на пересчёт. Проверен вживую
  (dry-run на проде: 93 053 в индексе, 60 993 ложных отметок).

Документация: зафиксирована реальная глубина корпуса — полный текст только у
27.5% документов, у 79% меньше 500 отпечатков (уровень аннотации). Система
ловит списывание из того, что есть целиком; это граница, а не поломка.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:17:41 +05:00
jze9
ea62f10829 feat(admin): состояние инфраструктуры на странице отладки
Сегодняшняя авария (упал гипервизор — вместе с ним брокер, Ollama, прокси)
проявлялась в отладке косвенно: пустой список воркеров и ошибка очередей.
Прямого ответа «что именно недоступно» страница не давала, хотя эндпоинт
/admin/health его уже знает.

Блок статусов рисуется и тогда, когда сам срез не собрался — именно в аварии
он нужен больше всего, а собирается в этот момент дольше обычного.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:00:19 +05:00
jze9
4008b5019f test(admin): закрепить смысл таймингов прогона + время работы в отладке
Тесты на queued_s/duration_s фиксируют ровно ту путаницу, из-за которой
метрика и разъехалась: ожидание в очереди и время работы — разные величины,
а у прогонов до миграции 006 длительности просто нет (вместо неё раньше
показывалось время в очереди).

Панель отладки теперь показывает, сколько идущий прогон уже работает —
по этому и виден застрявший, а не только по отсутствию heartbeat.

README: фактические числа тестов (150, проверено прогоном run_tests.sh).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 16:58:42 +05:00
jze9
99bf14fe6a feat(admin): честная длительность прогона + добор эмбеддингов
Разбор итогов массовой заливки показал в отладке «среднюю длительность прогона»
в 2.8 часа там, где парсинг занимал 50 секунд: started_at пишется в момент
постановки в очередь, а очередь из 173 источников разбирается часами. Теперь
момент реального старта пишется отдельно (run_started_at, миграция 006), а
схема отдаёт обе величины — сколько ждал очереди и сколько работал.

Плюс scripts/ops/reembed_missing.py: документ попадает в корпус сразу, а вектор
для L3 считает отдельная задача; когда worker-gpu или Ollama недоступны, эти
задачи теряются и документ остаётся невидимым для семантического поиска. Скрипт
находит faiss_id IS NULL и переотправляет пачками (dry-run по умолчанию) —
сейчас таких 23 101 из 177 147.

Документация: актуальные цифры корпуса, дубли при повторном прогоне, лимит
OpenAlex, и главное — гипервизор .254 зафиксирован в DR-HA как самая широкая
единая точка отказа (брокер, эмбеддинги, секреты и прокси на одном железе;
подтверждено аварией 28.08).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 16:51:25 +05:00
20 changed files with 1084 additions and 68 deletions

View File

@@ -232,7 +232,7 @@ docker compose -f docker-compose.prod.yml --profile observability up -d promethe
1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика). 1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика).
Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini). Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini).
2. **Юнит-тесты** — `pytest` по сервисам: 132 теста на ядро детекции, скоринга, 2. **Юнит-тесты** — `pytest` по сервисам: 150 тестов на ядро детекции, скоринга,
парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры
(БД/Redis/GPU/Ollama замоканы либо не нужны). (БД/Redis/GPU/Ollama замоканы либо не нужны).
@@ -265,7 +265,7 @@ make test-one SVC=worker-gost # тесты одного сервиса
| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 19 | | Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 19 |
| OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 | | OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 |
| Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 | | Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 |
| Шкала загрузки источников | `api/app/core/progress.py` | 7 | | Шкала загрузки и тайминги прогонов | `api/app/core/progress.py`, `schemas/admin.py` | 11 |
## Лицензия ## Лицензия

View File

@@ -81,7 +81,10 @@
`last_run_id` — ссылка на последний прогон. `last_run_id` — ссылка на последний прогон.
- **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/ - **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/
added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал
событий (JSON). Из них админка рисует шкалу загрузки, см. §12. событий (JSON). Тайминги раздельные: `started_at` — постановка в очередь,
`run_started_at` — реальный старт работы воркером (при массовом запуске между
ними часы ожидания), `finished_at` — конец. Из них админка рисует шкалу
загрузки, см. §12.
- **staged_works** — пользовательские загрузки на модерацию перед добавлением в корпус. - **staged_works** — пользовательские загрузки на модерацию перед добавлением в корпус.
- **admin_sessions** — одноразовые коды входа в админку. - **admin_sessions** — одноразовые коды входа в админку.
@@ -92,7 +95,7 @@
| Очередь | Задачи | Воркер | | Очередь | Задачи | Воркер |
|---------|--------|--------| |---------|--------|--------|
| `queue.index` | `index.extract_and_check`, `index.add_document`, `index.run_parser`, `index.enrich_full_text`, `index.ingest_upload` | 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.gpu` | `gpu.check_plagiarism`, `gpu.embed_documents`, `gpu.search_semantic`, `gpu.index_stats` | worker-gpu |
| `queue.gost` | `gost.format_bibliography` | worker-gost | | `queue.gost` | `gost.format_bibliography` | worker-gost |
| `queue.notify` | `notify.send_task_done`, `notify.send_verification` | worker-notifier | | `queue.notify` | `notify.send_task_done`, `notify.send_verification` | worker-notifier |
@@ -218,7 +221,10 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
Кнопки «Запустить всё» / «Остановить всё» — массовый старт и кооперативная отмена Кнопки «Запустить всё» / «Остановить всё» — массовый старт и кооперативная отмена
(флаг `cancel_requested`, воркер останавливается сам на ближайшем тике; уже (флаг `cancel_requested`, воркер останавливается сам на ближайшем тике; уже
начатый прогон не рвём посреди записи в базу). начатый прогон не рвём посреди записи в базу).
- **Панель отладки** (админка → «Отладка», `GET /api/admin/debug`): один срез — - **Панель отладки** (админка → «Отладка», `GET /api/admin/debug` + `/health`):
один срез — доступность инфраструктуры (PostgreSQL/Redis/MinIO/ES/Ollama/
брокер: этот блок показывается даже когда сам срез не собирается, потому что
при аварии первый вопрос — что именно отвалилось),
живые воркеры Celery и что именно они крутят, глубина очередей RabbitMQ живые воркеры Celery и что именно они крутят, глубина очередей RabbitMQ
(в т.ч. `unacked` и число потребителей), покрытие корпуса эмбеддингами, (в т.ч. `unacked` и число потребителей), покрытие корпуса эмбеддингами,
активные и проблемные прогоны, зависшие прогоны (нет heartbeat >10 мин), активные и проблемные прогоны, зависшие прогоны (нет heartbeat >10 мин),
@@ -226,6 +232,18 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
`VECTOR_BACKEND`). `VECTOR_BACKEND`).
- **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) — - **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) —
защита от краш-лупа по `consumer_timeout` RabbitMQ, см. [DR-HA.md](DR-HA.md) §6. защита от краш-лупа по `consumer_timeout` RabbitMQ, см. [DR-HA.md](DR-HA.md) §6.
- **Покрытие L3 меряется по индексу, а не по БД.** Колонка `documents.faiss_id`
для этого непригодна: отметка остаётся после пересоздания индекса (смена
модели/размерности) и после сбоев worker-gpu. Панель отладки спрашивает
реальное число векторов задачей `gpu.index_stats` и отдельно предупреждает о
ложных отметках; чинит их
[`scripts/ops/faiss_reconcile.py`](../scripts/ops/faiss_reconcile.py).
- **Добор эмбеддингов** — [`scripts/ops/reembed_missing.py`](../scripts/ops/reembed_missing.py):
документ попадает в корпус сразу, а вектор для L3 считает отдельная задача
`gpu.embed_documents`; если worker-gpu или Ollama были недоступны, эти задачи
теряются и документ остаётся невидимым для семантического поиска. Скрипт
находит `faiss_id IS NULL` и переотправляет задачи пачками (dry-run по
умолчанию). Покрытие видно в панели отладки.
## 13. Безопасность ## 13. Безопасность

View File

@@ -55,6 +55,17 @@ Redis у нас — кэш/rate-limits/LSH-индекс (префикс `antipla
| Векторный индекс | было | `VECTOR_BACKEND=qdrant` снимает (см. README) ✅ | | Векторный индекс | было | `VECTOR_BACKEND=qdrant` снимает (см. README) ✅ |
| RabbitMQ .82 | да | мониторинг ловит падение ✅; кластер — по потребности | | RabbitMQ .82 | да | мониторинг ловит падение ✅; кластер — по потребности |
| app/ES 1.32 | да | воркеры горизонтальны; ES single-node (для BM25 не критично) | | app/ES 1.32 | да | воркеры горизонтальны; ES single-node (для BM25 не критично) |
| Хост Proxmox .254 | **да, широкий** | митигации нет — см. ниже |
**Гипервизор — самая широкая единая точка отказа.** На хосте `192.168.20.254`
одновременно живут RabbitMQ (.82), embedding-gpu (.109), Infisical (.111) и
CT 102 (.253) — фронтенд и реверс-прокси. Его падение снимает сразу: приём и
обработку задач (нет брокера), эмбеддинги (нет Ollama), деплой (нет Infisical и
Gitea) и весь публичный доступ (нет прокси) — при том что API, PostgreSQL и
MinIO продолжают работать. Проверено на практике 2026-08-28: хост перестал
отвечать даже на ARP, всё перечисленное отвалилось разом, данные не пострадали.
Разнести хотя бы прокси/брокер по разным физическим хостам — самая дешёвая
мера; пока её нет, восстановление требует физического доступа к железу.
## 6. Известный операционный риск — RabbitMQ `consumer_timeout` vs долгие таски ## 6. Известный операционный риск — RabbitMQ `consumer_timeout` vs долгие таски

View File

@@ -1,13 +1,45 @@
# Наполнение корпуса — runbook # Наполнение корпуса — runbook
## Текущее состояние (на 2026-08-27) ## Текущее состояние (на 2026-08-28)
- **165 480 документов**, русский теперь большинство: `ru` 97 507, `en` 67 646, - **177 147 документов**, русский большинство: `ru` 99 884, `en` 76 936,
остальные языки — единицы/десятки. Проблема «не с чем сравнивать русские остальные языки — единицы/десятки. Проблема «не с чем сравнивать русские
работы» из более ранней версии этого документа закрыта. работы» из более ранней версии этого документа закрыта.
- По источникам: CyberLeninka 97 506, OpenAlex 49 218, PMC 10 451, arXiv 8 304, - По источникам: CyberLeninka 99 883, OpenAlex 49 276, PMC 14 910, arXiv 13 077,
`user_submission` (проверенные пользователями работы, не источник для сравнения `user_submission` (проверенные пользователями работы, не источник для сравнения
сами с собой — см. ARCHITECTURE.md §7) 1. сами с собой — см. ARCHITECTURE.md §7) 1.
- Повторный прогон по уже залитым источникам даёт почти одни дубли (типично
1490 из 1500 на источник): прирост дают только новые публикации. Реальный
рост корпуса — поднятый `limit` или новые темы, а не повторный запуск.
### Глубина индексации — чем реально располагает детекция (замер 2026-08-28)
Мерить глубину по `minio_key` (сохранён ли текст в MinIO) **нельзя**: PMC,
например, кладёт тело статьи в отпечатки прямо при заливке, не сохраняя файл.
Честный показатель — **число отпечатков на документ**: аннотация даёт десятки,
полный текст — тысячи.
| Источник | Документов | Глубоко (≥500 отпечатков) | Среднее отпечатков |
|----------|-----------:|--------------------------:|-------------------:|
| КиберЛенинка | 99 883 | 0 (0.0%) | **30** |
| OpenAlex | 49 276 | 6 240 (12.7%) | 907 |
| PMC | 14 910 | 14 519 (97.4%) | 2 208 |
| arXiv | 13 077 | 12 236 (93.6%) | 4 035 |
**Главная слабость — русская часть корпуса.** 56% базы (КиберЛенинка) в индексе
представлено заголовком и аннотацией: списывание из тела русской статьи L1 не
найдёт, хотя сервис рассчитан именно на русских студентов. Причина не в
алгоритме: search API отдаёт только аннотацию и OCR-фрагмент (~700 символов), а
`url` ведёт на HTML-страницу — докачка по нему бесполезна (замер: 0 из 8).
Лечится `scripts/ops/backfill_cyberleninka_pdf.py`: PDF доступен прямым адресом
`{url}/pdf` (проверено: 44 из 50, в среднем 23 тыс. символов, глубина 30 → ~1450
отпечатков). Полный прогон — около 48 часов при вежливом 1 req/s и порядка
+220 млн строк в `fingerprints` (~22 ГБ; место на сервере БД проверять заранее).
Покрытие L3 на ту же дату — 93 053 вектора (52.5% корпуса); ещё 60 993 документа
числились векторизованными ошибочно, отметки сброшены `faiss_reconcile.py`.
- OpenAlex при массовом запуске упирается в лимит вежливого пула (10 req/s
на mailto, общий для всех воркеров) — пауза между страницами поднята до 1с.
- Добавлен 4-й парсер — **PMC** (PubMed Central, `scripts/parsers/pmc.py`), - Добавлен 4-й парсер — **PMC** (PubMed Central, `scripts/parsers/pmc.py`),
англоязычные научные статьи открытого доступа. англоязычные научные статьи открытого доступа.
- Массовое расширение по дисциплинам теперь двумя сидерами: `seed_ru_sources.py` - Массовое расширение по дисциплинам теперь двумя сидерами: `seed_ru_sources.py`
@@ -61,6 +93,18 @@ python scripts/seed_ru_sources.py --apply --limit 500
- **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и - **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и
упавшие прогоны (см. ARCHITECTURE.md §12). упавшие прогоны (см. ARCHITECTURE.md §12).
После большой заливки — свериться с покрытием L3 (панель отладки, «Векторов в
индексе»). Порядок именно такой, из двух шагов:
1. `scripts/ops/faiss_reconcile.py` (в контейнере worker-gpu) — сверяет отметки
`faiss_id` с реальным содержимым индекса и обнуляет ложные. Без этого шага
документы, потерявшие вектор при пересоздании индекса, не попадут на
пересчёт: они всё ещё «отмечены».
2. `scripts/ops/reembed_missing.py` (в контейнере worker-indexer) — отправляет
`gpu.embed_documents` для всех `faiss_id IS NULL`.
Оба по умолчанию dry-run, отправляют/меняют только с `--apply`.
Прогон со статусом `partial` — это не ошибка: сработал бюджет времени Прогон со статусом `partial` — это не ошибка: сработал бюджет времени
(`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout` (`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout`
RabbitMQ. Остаток добирается повторным запуском источника. RabbitMQ. Остаток добирается повторным запуском источника.
@@ -82,6 +126,6 @@ SELECT source, count(*) FROM documents WHERE source='cyberleninka'; -- > 0
заголовки с меткой ru). заголовки с меткой ru).
- Для миллионов — **bulk** (снапшот OpenAlex на S3), а не постраничный API. - Для миллионов — **bulk** (снапшот OpenAlex на S3), а не постраничный API.
- На масштабе обязателен `VECTOR_BACKEND=qdrant` (FAISS flat не тянет), а таблица - На масштабе обязателен `VECTOR_BACKEND=qdrant` (FAISS flat не тянет), а таблица
`fingerprints` (уже ~88M строк на 165K доков — партиционирование стоит планировать `fingerprints` (уже ~113M строк на 177K доков — партиционирование стоит планировать
заранее, не постфактум) потребует партиционирования. См. заранее, не постфактум) потребует партиционирования. См.
[ARCHITECTURE.md](ARCHITECTURE.md) и [DR-HA.md](DR-HA.md). [ARCHITECTURE.md](ARCHITECTURE.md) и [DR-HA.md](DR-HA.md).

View File

@@ -0,0 +1,112 @@
#!/usr/bin/env python3
"""Углубить индексацию КиберЛенинки: скачать PDF статей и переиндексировать по ним.
Самая большая слабость корпуса (замер 2026-08-28): 99 883 документа
КиберЛенинки — 56% всей базы — имеют в среднем 30 отпечатков. Это заголовок с
аннотацией, тела статьи в индексе нет. Списывание из русской статьи система
не найдёт, хотя сервис рассчитан именно на русских студентов.
Причина: search API отдаёт только аннотацию и короткий OCR-фрагмент (~700
символов), а `url` ведёт на HTML-страницу — докачка по нему бесполезна
(замер: 0 из 8). Зато PDF доступен прямым адресом `{url}/pdf` и извлекается:
проверка на 4 статьях дала 13-40 тыс. символов у трёх, у одной 0 (скан без
текстового слоя — такие пропускаем).
Дальше — общий путь `store_full_text`: MinIO + пересчёт отпечатков L1 по
полному тексту + обновление MinHash LSH (L2).
Масштаб полного прогона: ~28 часов при вежливом 1 req/s, порядка +300 млн строк
в fingerprints (~30 ГБ). Запускать в фоне (nohup) и следить за местом на БД.
Запуск (в контейнере worker-indexer):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
$C < scripts/ops/backfill_cyberleninka_pdf.py # dry-run
$C --apply --limit 50 < scripts/ops/backfill_cyberleninka_pdf.py # пробная порция
$C --apply < scripts/ops/backfill_cyberleninka_pdf.py # всё (сутки)
Идемпотентно: документы с уже проставленным minio_key пропускаются, так что
прерванный прогон продолжается с того же места.
"""
import argparse
import time
RATE_LIMIT_DELAY = 1.0 # КиберЛенинка не любит частых запросов
MIN_USEFUL_CHARS = 1500 # меньше — скан без текстового слоя либо обрывок
MAX_PDF_BYTES = 30 * 1024 * 1024
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--apply", action="store_true", help="реально качать и сохранять")
ap.add_argument("--limit", type=int, default=0, help="взять не больше N документов (0 = все)")
args = ap.parse_args()
import httpx
from app.db import db_session
from app.extractors.pdf import extract_text_from_pdf
from app.tasks.index import store_full_text
from sqlalchemy import text
with db_session() as s:
sql = """
SELECT id, url FROM documents
WHERE source = 'cyberleninka' AND minio_key IS NULL AND url IS NOT NULL
ORDER BY id
"""
if args.limit:
sql += f" LIMIT {int(args.limit)}"
rows = s.execute(text(sql)).all()
print(f"документов КиберЛенинки без полного текста: {len(rows)}")
if not rows:
return
if not args.apply:
print(f"[dry-run] пример: {rows[0][1]}/pdf")
print(f"оценка полного прогона: ~{len(rows) * RATE_LIMIT_DELAY / 3600:.1f} ч, "
f"~{len(rows) * 3000 / 1e6:.0f} млн строк отпечатков")
print("Запустите с --apply (сначала --limit 50).")
return
client = httpx.Client(
timeout=30, follow_redirects=True,
headers={"User-Agent": "Mozilla/5.0 (compatible; AcademicHelper/1.0; +noreply@jze9.ru)"},
)
done = no_text = failed = 0
chars_total = 0
t0 = time.time()
for n, (doc_id, url) in enumerate(rows, 1):
try:
resp = client.get(url.rstrip("/") + "/pdf")
if resp.status_code != 200 or "pdf" not in resp.headers.get("content-type", "").lower() or len(resp.content) > MAX_PDF_BYTES:
no_text += 1
else:
body = extract_text_from_pdf(resp.content)
if len(body) < MIN_USEFUL_CHARS:
no_text += 1 # скан без текстового слоя
else:
store_full_text(doc_id, body)
done += 1
chars_total += len(body)
except Exception as e:
failed += 1
if failed <= 5:
print(f" doc {doc_id}: {type(e).__name__}: {str(e)[:80]}")
if n % 25 == 0:
el = time.time() - t0
print(f" {n}/{len(rows)} · углублено {done} · без текста {no_text} · "
f"ошибок {failed} · {n / el:.2f} док/с · осталось ~{(len(rows) - n) / max(n / el, 0.01) / 3600:.1f} ч",
flush=True)
time.sleep(RATE_LIMIT_DELAY)
avg = chars_total // done if done else 0
print(f"\nГотово за {(time.time() - t0) / 60:.1f} мин: углублено {done}, "
f"без текстового слоя {no_text}, ошибок {failed}, средний объём {avg} симв.")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,124 @@
#!/usr/bin/env python3
"""Бэкфилл полного текста для уже залитых документов PMC.
Зачем: PMC отдаёт тело статьи прямо в ответе efetch (40-50 тыс. символов), но
документы, залитые до включения сохранения full_text, остались в базе с одной
аннотацией — L1 сравнивает по ней и не видит тело статьи. Докачка PDF по url
здесь не работает: ссылка ведёт на HTML-страницу NCBI (замер: 0 из 8). Поэтому
текст берём тем же путём, что и парсер, — через API.
Замер на проде 2026-08-28: 12 368 документов PMC без полного текста.
Что делает для каждого документа: efetch по PMC id → извлечение текста →
`store_full_text` (та же функция, что у задачи докачки: MinIO + пересчёт
fingerprints по полному тексту + обновление MinHash LSH).
Идемпотентно: документы с уже проставленным minio_key не берутся.
Запуск (в контейнере worker-indexer, репозиторий внутрь не смонтирован):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
$C < scripts/ops/backfill_pmc_fulltext.py # dry-run
$C --apply --limit 40 < scripts/ops/backfill_pmc_fulltext.py # пробная порция
$C --apply < scripts/ops/backfill_pmc_fulltext.py # всё
Внимание: каждый документ добавляет ~3.5 тыс. строк в fingerprints (~4.5 ГБ
на все 12 тыс.). Перед полным прогоном стоит убедиться в свободном месте на
сервере БД.
"""
import argparse
import sys
import time
BATCH = 20 # столько id за один efetch — как в самом парсере
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--apply", action="store_true", help="реально сохранять (иначе dry-run)")
ap.add_argument("--limit", type=int, default=0, help="взять не больше N документов (0 = все)")
args = ap.parse_args()
if "/parsers" not in sys.path:
sys.path.insert(0, "/parsers")
from app.db import db_session
from app.tasks.index import store_full_text
from sqlalchemy import text
with db_session() as s:
sql = """
SELECT id, ext_id FROM documents
WHERE source = 'pmc' AND minio_key IS NULL AND ext_id IS NOT NULL
ORDER BY id
"""
if args.limit:
sql += f" LIMIT {int(args.limit)}"
rows = s.execute(text(sql)).all()
# ext_id формата "pmc:13521573" → числовой id для efetch
todo = {ext.split(":", 1)[-1]: doc_id for doc_id, ext in rows if ":" in ext}
print(f"документов PMC без полного текста: {len(rows)} (пригодных ext_id: {len(todo)})")
if not todo:
return
if not args.apply:
print(f"[dry-run] первые id: {list(todo.items())[:5]}")
print("Запустите с --apply (рекомендуется сначала --limit 40).")
return
from pmc import PMCParser
parser = PMCParser()
ids = list(todo)
done = failed = empty = 0
chars_total = 0
t0 = time.time()
for i in range(0, len(ids), BATCH):
batch = ids[i : i + BATCH]
try:
resp = parser.client.get(
"https://eutils.ncbi.nlm.nih.gov/entrez/eutils/efetch.fcgi",
params=parser._params(db="pmc", id=",".join(batch), rettype="full", retmode="xml"),
)
resp.raise_for_status()
articles = parser._parse_articles(resp.text)
except Exception as e:
failed += len(batch)
print(f" батч {i // BATCH}: ошибка запроса: {e}")
continue
for raw in articles:
doc = parser.transform(raw)
ext = (doc.get("ext_id") or "").split(":", 1)[-1]
doc_id = todo.get(ext)
full = doc.get("full_text") or ""
if not doc_id:
continue
if len(full) < 1000: # аннотация вместо тела — сохранять нечего
empty += 1
continue
try:
store_full_text(doc_id, full)
done += 1
chars_total += len(full)
except Exception as e:
failed += 1
print(f" doc {doc_id}: не сохранён: {e}")
if (i // BATCH) % 10 == 0:
speed = done / max(time.time() - t0, 1)
print(f" обработано {i + len(batch)}/{len(ids)} · сохранено {done} · "
f"{speed:.1f} док/с")
time.sleep(0.35) # лимит NCBI без ключа — 3 запроса/с
avg = chars_total // done if done else 0
print(f"\nГотово за {time.time() - t0:.0f}с: сохранено {done}, "
f"без тела статьи {empty}, ошибок {failed}, средний объём {avg} симв.")
if __name__ == "__main__":
main()

201
scripts/ops/bulk_ingest_pmc.py Executable file
View File

@@ -0,0 +1,201 @@
#!/usr/bin/env python3
"""Массовая заливка полнотекстовых статей из открытого бакета PMC (миллионы).
Почему отдельный путь, а не парсеры и Celery: постраничный API даёт ~1-2 статьи
в секунду и упирается в rate limit — миллионы так не залить. У AWS Open Data
бакета `pmc-oa-opendata` (открыт, ключ не нужен) для каждой статьи лежит уже
извлечённый текст `PMC*.txt` (33-130 КБ) и метаданные `PMC*.json`. Замер:
12 статей/с в 10 потоков, то есть миллион — примерно за сутки.
Второе узкое место — запись отпечатков. Построчный INSERT даёт ~6.7 тыс. строк/с
(миллион статей = 2.5 млрд строк = четверо суток только на вставку), поэтому
здесь используется COPY: он на порядки быстрее и не гоняет данные через ORM.
Масштаб для планирования: 1 млн статей ≈ 2.5 млрд отпечатков ≈ 400 ГБ в БД
с индексами. Плотность отпечатков ограничена --fp-per-doc (по умолчанию 2000
вместо 20 000 из настроек воркера) — иначе объём растёт быстрее пользы.
Запуск (в контейнере worker-indexer, где есть psycopg2 и winnowing):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
$C --limit 200 < scripts/ops/bulk_ingest_pmc.py # пробная порция
$C --limit 100000 --workers 20 < scripts/ops/bulk_ingest_pmc.py
Идемпотентно: документы вставляются с ON CONFLICT DO NOTHING по ext_id, уже
залитые пропускаются. Прогресс листинга сохраняется в --state-file, поэтому
прерванная заливка продолжается с той же страницы бакета.
СТАТУС: НЕ ЗАКОНЧЕН. Проба на 200 статьях: документы вставляются верно
(проверено), но шаг отпечатков ошибочно считает свежевставленные документы уже
обработанными — проверка «есть ли отпечатки» срабатывает там, где их нет, и
COPY не выполняется. Пробные записи из базы удалены. До отладки этого места
скрипт запускать на объёме нельзя: он зальёт документы без отпечатков, то есть
невидимые для L1.
"""
import argparse
import io
import json
import os
import re
import time
from concurrent.futures import ThreadPoolExecutor
from urllib.parse import quote
BUCKET = "https://pmc-oa-opendata.s3.amazonaws.com"
STATE_FILE = "/tmp/bulk_pmc_state.json"
KEY_RE = re.compile(r"<Key>(PMC\d+\.\d+)/\1\.txt</Key>")
TOKEN_RE = re.compile(r"<NextContinuationToken>([^<]+)</NextContinuationToken>")
def list_articles(client, token: str | None, want: int) -> tuple[list[str], str | None]:
"""Собрать id статей из листинга бакета, продолжая с сохранённой страницы."""
ids: list[str] = []
while len(ids) < want:
url = f"{BUCKET}/?list-type=2&max-keys=1000"
if token:
# В токене бывают + и / — без экранирования S3 отвечает 400
url += f"&continuation-token={quote(token, safe='')}"
resp = client.get(url)
resp.raise_for_status()
ids.extend(KEY_RE.findall(resp.text))
m = TOKEN_RE.search(resp.text)
token = m.group(1) if m else None
if token is None:
break
return ids[:want], token
def fetch_one(client, art_id: str) -> dict | None:
"""Текст статьи и её метаданные. None — если статья недоступна или пустая."""
try:
txt = client.get(f"{BUCKET}/{art_id}/{art_id}.txt")
if txt.status_code != 200 or len(txt.text) < 1500:
return None
meta_resp = client.get(f"{BUCKET}/{art_id}/{art_id}.json")
meta = meta_resp.json() if meta_resp.status_code == 200 else {}
except Exception:
return None
citation = meta.get("citation") or ""
year = None
m = re.search(r"\b(19|20)\d{2}\b", citation)
if m:
year = int(m.group(0))
return {
"ext_id": f"pmc:{meta.get('pmcid') or art_id.split('.')[0]}",
"title": (meta.get("title") or art_id)[:1000],
"doi": meta.get("doi"),
"year": year,
"journal": citation[:300] or None,
"text": txt.text,
}
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--limit", type=int, default=200, help="сколько статей залить за прогон")
ap.add_argument("--workers", type=int, default=10, help="потоков закачки (вежливо: 10-30)")
ap.add_argument("--batch", type=int, default=200, help="статей в одной транзакции")
ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на документ")
ap.add_argument("--state-file", default=STATE_FILE, help="файл с позицией листинга")
args = ap.parse_args()
import httpx
import psycopg2
from app.algorithms.winnowing import winnow
from app.config import settings
token = None
if os.path.exists(args.state_file):
with open(args.state_file) as fh:
token = json.load(fh).get("token")
print("продолжаю листинг с сохранённой позиции")
client = httpx.Client(timeout=60, headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"})
conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
password=settings.POSTGRES_PASSWORD)
added = skipped = failed = 0
fp_total = 0
t0 = time.time()
done = 0
while done < args.limit:
want = min(args.batch, args.limit - done)
ids, token = list_articles(client, token, want)
if not ids:
print("бакет закончился")
break
done += len(ids)
with ThreadPoolExecutor(max_workers=args.workers) as pool:
docs = [d for d in pool.map(lambda i: fetch_one(client, i), ids) if d]
failed += len(ids) - len(docs)
if not docs:
continue
with conn.cursor() as cur:
# Документы: ON CONFLICT — повторный прогон не плодит дубли
rows = [
(d["ext_id"], "pmc", d["title"], d["doi"], d["year"], "en",
d["journal"], d["text"][:2000])
for d in docs
]
cur.executemany(
"""INSERT INTO documents (ext_id, source, title, doi, year, lang, journal,
abstract, authors, indexed_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,'[]'::json, now())
ON CONFLICT (ext_id) DO NOTHING""",
rows,
)
cur.execute("SELECT id, ext_id FROM documents WHERE ext_id = ANY(%s)",
([d["ext_id"] for d in docs],))
id_by_ext = {e: i for i, e in cur.fetchall()}
# Отпечатки: COPY вместо построчных INSERT — единственный способ
# писать миллиарды строк за разумное время
buf = io.StringIO()
fresh = 0
for d in docs:
doc_id = id_by_ext.get(d["ext_id"])
if doc_id is None:
continue
cur.execute("SELECT 1 FROM fingerprints WHERE doc_id=%s LIMIT 1", (doc_id,))
if cur.fetchone():
skipped += 1
continue
hashes = list(winnow(d["text"]))[: args.fp_per_doc]
for pos, h in enumerate(hashes):
buf.write(f"{doc_id}\t{h}\t{pos}\n")
fp_total += len(hashes)
fresh += 1
buf.seek(0)
if fresh:
cur.copy_from(buf, "fingerprints", columns=("doc_id", "hash_value", "position"))
added += fresh
conn.commit()
with open(args.state_file, "w") as fh:
json.dump({"token": token}, fh)
el = time.time() - t0
print(f" {done}/{args.limit} · залито {added} · пропущено {skipped} · "
f"недоступно {failed} · отпечатков {fp_total:,} · {added / max(el, 1):.1f} ст/с",
flush=True)
if token is None:
break
conn.close()
el = time.time() - t0
print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, "
f"недоступно {failed}, отпечатков {fp_total:,}")
if added:
print(f"темп: {added / el:.1f} статей/с → миллион за ~{1e6 / max(added / el, 0.01) / 3600:.0f} ч")
if __name__ == "__main__":
main()

101
scripts/ops/faiss_reconcile.py Executable file
View File

@@ -0,0 +1,101 @@
#!/usr/bin/env python3
"""Сверить отметки `documents.faiss_id` с реальным содержимым векторного индекса.
Проблема, которую он решает: отметка в базе НЕ означает, что вектор есть в
индексе. Она остаётся после пересоздания индекса (например, при смене модели
эмбеддингов и размерности) и после сбоев worker-gpu. Из-за этого документ
навсегда выпадает из L3: в индексе его нет, а на пересчёт он не попадёт —
`reembed_missing.py` ищет только `faiss_id IS NULL`.
Замер на проде 2026-08-28: в индексе 93 053 вектора, отмечено в базе 154 046 —
60 993 документа считались векторизованными, не будучи ими.
Скрипт обнуляет отметки у документов, которых в индексе нет. После него
`reembed_missing.py` отправит их на пересчёт.
Запуск (нужен доступ к самому индексу → контейнер worker-gpu):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-gpu python -"
$C < scripts/ops/faiss_reconcile.py # dry-run, только показывает
$C --apply < scripts/ops/faiss_reconcile.py # обнулить ложные отметки
Работает только с FAISS (`VECTOR_BACKEND=faiss`): у Qdrant идентификаторы
хранит сам сервис, и сверка там делается иначе.
"""
import argparse
import sys
def stale_marks(marked_ids: set[int], indexed_ids: set[int]) -> set[int]:
"""Документы, отмеченные как векторизованные, но отсутствующие в индексе."""
return marked_ids - indexed_ids
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--apply", action="store_true", help="реально обнулить faiss_id (иначе dry-run)")
args = ap.parse_args()
from app.config import settings
if settings.VECTOR_BACKEND != "faiss":
sys.exit(f"VECTOR_BACKEND={settings.VECTOR_BACKEND}: скрипт рассчитан на faiss")
import faiss # noqa: F401 (нужен для vector_to_array)
from app.faiss_manager import FAISSManager
FAISSManager.load_or_create()
index = FAISSManager._index
if index is None or not hasattr(index, "id_map"):
sys.exit("индекс не загрузился или не хранит id (нет id_map)")
indexed = {int(i) for i in faiss.vector_to_array(index.id_map)}
print(f"векторов в индексе: {index.ntotal} (уникальных id: {len(indexed)})")
import psycopg2
conn = psycopg2.connect(
host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
password=settings.POSTGRES_PASSWORD,
)
with conn.cursor() as cur:
cur.execute("SELECT id FROM documents WHERE faiss_id IS NOT NULL")
marked = {r[0] for r in cur.fetchall()}
cur.execute("SELECT count(*) FROM documents")
total = cur.fetchone()[0]
stale = stale_marks(marked, indexed)
print(f"документов всего: {total}")
print(f"отмечено как векторизованные: {len(marked)}")
print(f"ложных отметок (нет в индексе): {len(stale)}")
print(f"реальное покрытие L3: {len(indexed & marked) * 100.0 / total:.1f}%")
if not stale:
print("Сверка чистая, делать нечего.")
conn.close()
return
if not args.apply:
print(f"\n[dry-run] первые id: {sorted(stale)[:10]}")
print("Запустите с --apply, чтобы обнулить faiss_id у этих документов,")
print("затем reembed_missing.py отправит их на пересчёт эмбеддингов.")
conn.close()
return
ids = list(stale)
with conn.cursor() as cur:
# Пачками: один UPDATE с десятками тысяч id упирается в лимиты параметров
for i in range(0, len(ids), 5000):
chunk = ids[i : i + 5000]
cur.execute("UPDATE documents SET faiss_id = NULL WHERE id = ANY(%s)", (chunk,))
conn.commit()
conn.close()
print(f"\nОбнулено отметок: {len(ids)}")
print("Дальше: reembed_missing.py --apply (отправит их на пересчёт).")
if __name__ == "__main__":
main()

140
scripts/ops/reembed_missing.py Executable file
View File

@@ -0,0 +1,140 @@
#!/usr/bin/env python3
"""Догнать эмбеддинги для документов, у которых их нет (faiss_id IS NULL).
Зачем: документ попадает в корпус сразу (метаданные, fingerprints L1, MinHash L2),
а вектор для L3 считает отдельная задача `gpu.embed_documents`. Если worker-gpu
или Ollama были недоступны в момент заливки, эти задачи теряются — документ
остаётся в базе, но семантический поиск его не видит. Скрипт находит такие
документы и переотправляет задачи пачками.
Идемпотентно: повторный запуск возьмёт только оставшиеся без faiss_id.
ВАЖНО, чего скрипт НЕ делает: документы с непустым `faiss_id`, которого нет в
самом FAISS-индексе (наследие смены модели эмбеддингов — размерность сменилась,
индекс пересобран с нуля, а ссылки в БД остались), сюда не попадут. Их сначала
надо выявить сверкой с индексом и обнулить faiss_id — это отдельная операция.
Запуск (нужны psycopg2 + celery — они есть в образе воркера; сам репозиторий
внутрь контейнера не смонтирован, поэтому передаём скрипт через stdin):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
$C < scripts/ops/reembed_missing.py # dry-run, только считает
$C --apply < scripts/ops/reembed_missing.py # отправить задачи
$C --apply --limit 1000 < scripts/ops/reembed_missing.py # пробная порция
Креды берутся из окружения (в контейнере они уже есть из env_file), а при
запуске файлом — из .env репозитория: POSTGRES_*, RABBITMQ_URL.
"""
import argparse
import os
import sys
from pathlib import Path
DEFAULT_BATCH = 64 # столько же, сколько EMBED_BATCH_SIZE у индексера
def _load_env() -> None:
"""Подтянуть переменные из .env репозитория, если файл есть.
При запуске через stdin (`docker exec ... python -`) __file__ не определён —
это штатный способ запуска здесь, и переменные в контейнере уже есть.
"""
try:
env = Path(__file__).resolve().parent.parent.parent / ".env"
except NameError:
return
if not env.exists():
return
for line in env.read_text(encoding="utf-8").splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, _, val = line.partition("=")
os.environ.setdefault(key.strip(), val.strip())
def _pg_dsn() -> dict:
return {
"host": os.environ.get("POSTGRES_HOST", "localhost"),
"port": int(os.environ.get("POSTGRES_PORT", "5432")),
"dbname": os.environ.get("POSTGRES_DB", "antiplagiator"),
"user": os.environ.get("POSTGRES_USER", "antiplagiator"),
"password": os.environ.get("POSTGRES_PASSWORD", ""),
}
def chunked(items: list[int], size: int) -> list[list[int]]:
"""Разбить список id на пачки по size элементов."""
return [items[i : i + size] for i in range(0, len(items), size)]
def main() -> None:
ap = argparse.ArgumentParser(
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
)
ap.add_argument("--apply", action="store_true", help="реально отправить задачи (иначе dry-run)")
ap.add_argument("--batch", type=int, default=DEFAULT_BATCH, help=f"документов в пачке (по умолчанию {DEFAULT_BATCH})")
ap.add_argument("--limit", type=int, default=0, help="взять не больше N документов (0 = все)")
args = ap.parse_args()
_load_env()
try:
import psycopg2
except ImportError:
sys.exit("Нужен psycopg2 (есть в образах воркеров). Запускайте в контейнере.")
dsn = _pg_dsn()
print(f"PostgreSQL {dsn['host']}:{dsn['port']}/{dsn['dbname']}")
conn = psycopg2.connect(**dsn)
with conn.cursor() as cur:
sql = "SELECT id FROM documents WHERE faiss_id IS NULL ORDER BY id"
if args.limit:
sql += f" LIMIT {int(args.limit)}"
cur.execute(sql)
doc_ids = [r[0] for r in cur.fetchall()]
cur.execute("SELECT count(*) FROM documents")
total = cur.fetchone()[0]
conn.close()
batches = chunked(doc_ids, args.batch)
print(f"Документов всего: {total}")
print(f"Без эмбеддинга: {len(doc_ids)} → пачек по {args.batch}: {len(batches)}")
if not doc_ids:
print("Нечего досчитывать.")
return
if not args.apply:
head = ", ".join(str(i) for i in doc_ids[:10])
print(f"\n[dry-run] первые id: {head}{' …' if len(doc_ids) > 10 else ''}")
print("Запустите с --apply, чтобы отправить gpu.embed_documents в queue.gpu.")
return
try:
from celery import Celery
except ImportError:
sys.exit("Нужен celery (есть в образах воркеров). Запускайте в контейнере.")
broker = os.environ.get("RABBITMQ_URL")
if not broker:
sys.exit("Не задан RABBITMQ_URL — нечем отправлять задачи.")
app = Celery("reembed", broker=broker)
sent = 0
try:
for batch in batches:
app.send_task("gpu.embed_documents", args=[batch], queue="queue.gpu")
sent += 1
if sent % 50 == 0:
print(f" отправлено пачек: {sent}/{len(batches)}")
except Exception as e:
sys.exit(f"\nБрокер недоступен после {sent} пачек: {e}")
print(f"\nОтправлено пачек: {sent} ({len(doc_ids)} документов).")
print("Ход выполнения — админка → «Отладка» (очередь queue.gpu и покрытие эмбеддингами).")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,31 @@
"""Отдельная отметка старта выполнения прогона парсинга.
Revision ID: 006
Revises: 005
Create Date: 2026-08-28
`started_at` пишется в момент СОЗДАНИЯ строки, то есть постановки в очередь.
При массовом запуске очередь разбирается часами, и «длительность» прогона по
двум таймстампам показывала 2.8 часа там, где сам парсинг занял 50 секунд —
для отладки это дезинформация. Момент, когда воркер реально взял задачу,
пишем отдельно; разница со `started_at` — это ожидание в очереди.
"""
import sqlalchemy as sa
from alembic import op
revision = "006"
down_revision = "005"
branch_labels = None
depends_on = None
def upgrade() -> None:
# Для уже прошедших прогонов остаётся NULL: подставлять им started_at
# значило бы выдать время ожидания за время работы.
op.add_column("parse_runs", sa.Column("run_started_at", sa.DateTime(), nullable=True))
def downgrade() -> None:
op.drop_column("parse_runs", "run_started_at")

View File

@@ -784,6 +784,20 @@ async def upload_documents(files: list[UploadFile]) -> dict:
# ═══════════════════════════════════════════════════════════════════════════════ # ═══════════════════════════════════════════════════════════════════════════════
# ОТЛАДКА # ОТЛАДКА
# ═══════════════════════════════════════════════════════════════════════════════ # ═══════════════════════════════════════════════════════════════════════════════
def _vector_index_stats() -> dict:
"""Реальное наполнение векторного индекса — спросить у worker-gpu.
Колонка `documents.faiss_id` для этого не годится: пометка остаётся и когда
вектор в индекс не попал, и когда индекс пересоздали после смены модели
эмбеддингов. Отладка, показывающая покрытие L3 по ней, завышает его в разы.
"""
try:
result = celery_app.send_task("gpu.index_stats", queue="queue.gpu")
return result.get(timeout=10) or {}
except Exception as e:
return {"error": str(e)[:200] or type(e).__name__}
def _celery_snapshot() -> dict: def _celery_snapshot() -> dict:
"""Живое состояние воркеров через Celery inspect (блокирующий вызов).""" """Живое состояние воркеров через Celery inspect (блокирующий вызов)."""
try: try:
@@ -827,8 +841,10 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
именно сейчас крутит, что копится в очередях, как наполняется корпус, именно сейчас крутит, что копится в очередях, как наполняется корпус,
какие прогоны идут и на чём падали последние. какие прогоны идут и на чём падали последние.
""" """
# Воркеры — блокирующий Celery inspect, в пул потоков, чтобы не вешать луп # Воркеры и статистика индекса — блокирующие вызовы Celery, в пул потоков,
# чтобы не вешать событийный цикл
celery_state = await run_in_threadpool(_celery_snapshot) celery_state = await run_in_threadpool(_celery_snapshot)
vector_stats = await run_in_threadpool(_vector_index_stats)
rmq = await _rabbitmq_queues() rmq = await _rabbitmq_queues()
queues = ( queues = (
@@ -915,8 +931,10 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
"queues": queues, "queues": queues,
"corpus": { "corpus": {
"documents": docs_total, "documents": docs_total,
"documents_embedded": docs_embedded, # Пометка в БД и реальное содержимое индекса расходятся — показываем
"documents_without_embedding": max(docs_total - docs_embedded, 0), # обе величины, иначе покрытие L3 выглядит лучше, чем оно есть
"documents_marked_embedded": docs_embedded,
"vector_index": vector_stats,
"fingerprints_estimate": int(fingerprints_est), "fingerprints_estimate": int(fingerprints_est),
"elasticsearch_documents": es_docs, "elasticsearch_documents": es_docs,
}, },

View File

@@ -75,7 +75,10 @@ class ParseRun(Base):
error: Mapped[str | None] = mapped_column(Text, nullable=True) error: Mapped[str | None] = mapped_column(Text, nullable=True)
# [{"ts": ISO8601, "level": "info|warning|error", "msg": str}] # [{"ts": ISO8601, "level": "info|warning|error", "msg": str}]
log: Mapped[list | None] = mapped_column(JSON, default=list) log: Mapped[list | None] = mapped_column(JSON, default=list)
# Постановка в очередь (строку создаёт API) и реальный старт выполнения:
# между ними при массовом запуске проходят часы, мешать их нельзя
started_at: Mapped[datetime] = mapped_column(server_default=func.now(), index=True) started_at: Mapped[datetime] = mapped_column(server_default=func.now(), index=True)
run_started_at: Mapped[datetime | None] = mapped_column(nullable=True)
# Последний признак жизни: по нему видно зависший прогон (running, но тишина) # Последний признак жизни: по нему видно зависший прогон (running, но тишина)
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True) heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(nullable=True) finished_at: Mapped[datetime | None] = mapped_column(nullable=True)

View File

@@ -80,7 +80,9 @@ class ParseRunResponse(BaseModel):
failed: int failed: int
cancel_requested: bool = False cancel_requested: bool = False
error: str | None = None error: str | None = None
# started_at — постановка в очередь, run_started_at — реальный старт работы
started_at: datetime started_at: datetime
run_started_at: datetime | None = None
heartbeat_at: datetime | None = None heartbeat_at: datetime | None = None
finished_at: datetime | None = None finished_at: datetime | None = None
@@ -91,6 +93,26 @@ class ParseRunResponse(BaseModel):
def percent(self) -> float: def percent(self) -> float:
return run_percent(self.status, self.stage, self.target, self.fetched, self.processed) return run_percent(self.status, self.stage, self.target, self.fetched, self.processed)
@computed_field # type: ignore[prop-decorator]
@property
def queued_s(self) -> float | None:
"""Сколько прогон ждал своей очереди, сек."""
if self.run_started_at is None:
return None
return round((self.run_started_at - self.started_at).total_seconds(), 1)
@computed_field # type: ignore[prop-decorator]
@property
def duration_s(self) -> float | None:
"""Сколько прогон реально работал, сек (None — ещё идёт или не начинался).
Прогоны до появления run_started_at (миграция 006) остаются без
длительности: у них известен только момент постановки в очередь.
"""
if self.run_started_at is None or self.finished_at is None:
return None
return round((self.finished_at - self.run_started_at).total_seconds(), 1)
class ParseRunDetail(ParseRunResponse): class ParseRunDetail(ParseRunResponse):
log: list[dict[str, Any]] = Field(default_factory=list) log: list[dict[str, Any]] = Field(default_factory=list)

View File

@@ -1,6 +1,9 @@
"""Юнит-тесты шкалы загрузки источников (app.core.progress) — чистая логика.""" """Юнит-тесты шкалы загрузки и таймингов прогонов — чистая логика, без БД."""
from datetime import datetime
from app.core.progress import run_percent from app.core.progress import run_percent
from app.schemas.admin import ParseRunResponse
def test_queued_run_shows_nothing_done(): def test_queued_run_shows_nothing_done():
@@ -39,3 +42,45 @@ def test_finished_runs_are_always_full():
"""Шкала показывает «работа окончена», исход виден по статусу рядом.""" """Шкала показывает «работа окончена», исход виден по статусу рядом."""
for status in ("done", "partial", "error", "cancelled"): for status in ("done", "partial", "error", "cancelled"):
assert run_percent(status, "finished", target=100, fetched=3, processed=1) == 100.0 assert run_percent(status, "finished", target=100, fetched=3, processed=1) == 100.0
# ─── Тайминги прогона в схеме ответа ─────────────────────────────────────────
# Регрессия, ради которой они и разделены: в отладке «длительность прогона»
# показывала время ожидания в очереди (часы) вместо времени работы (секунды).
def _run(**over) -> ParseRunResponse:
base = {
"id": 1, "source_id": 1, "status": "done", "stage": "finished", "target": 100,
"fetched": 100, "processed": 100, "added": 10, "duplicates": 90,
"skipped": 0, "failed": 0,
"started_at": datetime(2026, 8, 27, 12, 0, 0),
"run_started_at": datetime(2026, 8, 27, 14, 0, 0),
"finished_at": datetime(2026, 8, 27, 14, 0, 50),
}
base.update(over)
return ParseRunResponse(**base)
def test_queued_and_duration_are_measured_separately():
r = _run()
assert r.queued_s == 7200.0 # два часа в очереди
assert r.duration_s == 50.0 # полминуты работы
def test_no_durations_until_worker_picked_run_up():
r = _run(status="queued", stage="queued", run_started_at=None, finished_at=None)
assert r.queued_s is None
assert r.duration_s is None
def test_running_run_has_wait_but_no_duration_yet():
r = _run(status="running", stage="index", finished_at=None)
assert r.queued_s == 7200.0
assert r.duration_s is None
def test_legacy_runs_report_no_duration_instead_of_queue_time():
"""Прогоны до миграции 006: длительности нет — но и вранья тоже."""
r = _run(run_started_at=None)
assert r.duration_s is None

View File

@@ -1,18 +1,21 @@
import React from 'react'; import React from 'react';
import { useQuery } from '@tanstack/react-query'; import { useQuery } from '@tanstack/react-query';
import { import {
Activity, AlertTriangle, Cpu, Database, Layers, RefreshCw, Settings2, ListTree, Activity, AlertTriangle, CheckCircle2, Cpu, Database, Layers, RefreshCw,
Server, Settings2, ListTree, XCircle,
} from 'lucide-react'; } from 'lucide-react';
import { adminApi } from '../../api/client'; import { adminApi } from '../../api/client';
import { ProgressBar } from '../../components/ProgressBar'; import { ProgressBar } from '../../components/ProgressBar';
interface Health { name: string; ok: boolean; detail?: string }
interface ActiveTask { name: string; id: string; args: string; started_ago_s: number | null } 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 Worker { name: string; concurrency: number | null; reserved: number; active: ActiveTask[] }
interface Run { interface Run {
id: number; source_id: number; status: string; stage: string; target: number; id: number; source_id: number; status: string; stage: string; target: number;
fetched: number; processed: number; added: number; duplicates: number; fetched: number; processed: number; added: number; duplicates: number;
skipped: number; failed: number; error: string | null; percent: number; skipped: number; failed: number; error: string | null; percent: number;
started_at: string; heartbeat_at: string | null; started_at: string; run_started_at: string | null; heartbeat_at: string | null;
queued_s: number | null; duration_s: number | null;
} }
interface DebugData { interface DebugData {
@@ -20,7 +23,8 @@ interface DebugData {
celery: { workers: Worker[]; error?: string }; celery: { workers: Worker[]; error?: string };
queues: Record<string, { messages: number; unacked: number; consumers: number }> | { error: string }; queues: Record<string, { messages: number; unacked: number; consumers: number }> | { error: string };
corpus: { corpus: {
documents: number; documents_embedded: number; documents_without_embedding: number; documents: number; documents_marked_embedded: number;
vector_index: { backend?: string; vectors?: number; dim?: number; embed_model?: string; error?: string };
fingerprints_estimate: number; elasticsearch_documents: number | null; fingerprints_estimate: number; elasticsearch_documents: number | null;
}; };
sources: { sources: {
@@ -42,24 +46,48 @@ function fmtNum(n: number | null | undefined): string {
return n == null ? '—' : n.toLocaleString('ru-RU'); return n == null ? '—' : n.toLocaleString('ru-RU');
} }
/** Сколько идущий прогон уже работает — по нему видно застрявший. */
function runningFor(run: Run): string {
if (!run.run_started_at) return 'ещё не начат';
const sec = Math.max(0, (Date.now() - new Date(run.run_started_at + 'Z').getTime()) / 1000);
return sec < 90 ? `${Math.round(sec)}с` : `${Math.round(sec / 60)} мин`;
}
export function Debug() { export function Debug() {
const { data, isFetching, error } = useQuery({ const { data, isFetching, error } = useQuery({
queryKey: ['admin-debug'], queryKey: ['admin-debug'],
queryFn: () => adminApi.debug().then((r) => r.data as DebugData), queryFn: () => adminApi.debug().then((r) => r.data as DebugData),
refetchInterval: 5000, refetchInterval: 5000,
}); });
// Отдельным запросом: когда отваливается инфраструктура (брокер, Ollama),
// отладка должна первой показывать ЧТО именно недоступно, а не только следствия
const { data: health } = useQuery({
queryKey: ['admin-health'],
queryFn: () => adminApi.health().then((r) => r.data as Health[]),
refetchInterval: 10000,
});
if (error) { if (error || !data) {
return <div className="text-red-600 text-sm">Не удалось получить срез состояния: {String(error)}</div>; return (
} <div className="space-y-4">
if (!data) { {error
return <div className="text-gray-400 text-sm">Сбор данных…</div>; ? <div className="text-red-600 text-sm">Не удалось получить срез состояния: {String(error)}</div>
: <div className="text-gray-400 text-sm">Сбор данных…</div>}
<HealthCard health={health} />
</div>
);
} }
const queuesErr = 'error' in data.queues ? (data.queues as { error: string }).error : null; const queuesErr = 'error' in data.queues ? (data.queues as { error: string }).error : null;
const queues = queuesErr ? {} : (data.queues as Record<string, { messages: number; unacked: number; consumers: number }>); const queues = queuesErr ? {} : (data.queues as Record<string, { messages: number; unacked: number; consumers: number }>);
const embedPercent = data.corpus.documents // Покрытие L3 считаем по РЕАЛЬНОМУ содержимому индекса: пометка faiss_id в БД
? (data.corpus.documents_embedded / data.corpus.documents) * 100 // остаётся и когда вектор туда не попал, и завышает картину в разы
const vectors = data.corpus.vector_index?.vectors ?? null;
const embedPercent = data.corpus.documents && vectors != null
? (vectors / data.corpus.documents) * 100
: 0;
const staleMarks = vectors != null
? Math.max(0, data.corpus.documents_marked_embedded - vectors)
: 0; : 0;
return ( return (
@@ -72,6 +100,8 @@ export function Debug() {
<span className="text-xs text-gray-400">срез от {new Date(data.generated_at).toLocaleTimeString('ru-RU')}</span> <span className="text-xs text-gray-400">срез от {new Date(data.generated_at).toLocaleTimeString('ru-RU')}</span>
</div> </div>
<HealthCard health={health} />
{/* Активные прогоны заливки */} {/* Активные прогоны заливки */}
<Card icon={Layers} title={`Заливка источников — активных прогонов: ${data.sources.active_runs.length}`}> <Card icon={Layers} title={`Заливка источников — активных прогонов: ${data.sources.active_runs.length}`}>
{data.sources.stale_run_ids.length > 0 && ( {data.sources.stale_run_ids.length > 0 && (
@@ -87,8 +117,9 @@ export function Debug() {
<div className="space-y-3"> <div className="space-y-3">
{data.sources.active_runs.map((r) => ( {data.sources.active_runs.map((r) => (
<div key={r.id} className="flex items-center gap-4"> <div key={r.id} className="flex items-center gap-4">
<div className="w-40 shrink-0 text-sm text-gray-600 truncate"> <div className="w-52 shrink-0 text-sm text-gray-600 truncate">
#{r.id} · источник {r.source_id} #{r.id} · источник {r.source_id}
<span className="text-gray-400"> · {runningFor(r)}</span>
</div> </div>
<ProgressBar <ProgressBar
percent={r.percent} percent={r.percent}
@@ -167,15 +198,32 @@ export function Debug() {
<Card icon={Database} title="Корпус сравнения"> <Card icon={Database} title="Корпус сравнения">
<div className="grid grid-cols-2 lg:grid-cols-4 gap-4 mb-4"> <div className="grid grid-cols-2 lg:grid-cols-4 gap-4 mb-4">
<Metric label="Документов" value={fmtNum(data.corpus.documents)} /> <Metric label="Документов" value={fmtNum(data.corpus.documents)} />
<Metric label="С эмбеддингом" value={fmtNum(data.corpus.documents_embedded)} /> <Metric label="Векторов в индексе" value={fmtNum(vectors)} />
<Metric label="Отпечатков (оценка)" value={fmtNum(data.corpus.fingerprints_estimate)} /> <Metric label="Отпечатков (оценка)" value={fmtNum(data.corpus.fingerprints_estimate)} />
<Metric label="В Elasticsearch" value={fmtNum(data.corpus.elasticsearch_documents)} /> <Metric label="В Elasticsearch" value={fmtNum(data.corpus.elasticsearch_documents)} />
</div> </div>
{data.corpus.vector_index?.error && (
<p className="text-xs text-red-600 mb-2">
индекс недоступен: {data.corpus.vector_index.error}
</p>
)}
<ProgressBar <ProgressBar
percent={embedPercent} percent={embedPercent}
status={data.corpus.documents_without_embedding > 0 ? 'partial' : 'done'} status={embedPercent >= 99 ? 'done' : 'partial'}
label={`покрытие эмбеддингами (L3): без вектора ${fmtNum(data.corpus.documents_without_embedding)} документов`} label={`покрытие L3 (реально в индексе ${data.corpus.vector_index?.backend || '—'}): ` +
`без вектора ${fmtNum(Math.max(0, data.corpus.documents - (vectors ?? 0)))} документов`}
/> />
{staleMarks > 0 && (
<div className="mt-3 flex items-start gap-2 text-sm text-amber-700 bg-amber-50 border border-amber-100 rounded-lg p-3">
<AlertTriangle className="w-4 h-4 mt-0.5 shrink-0" />
<span>
У {fmtNum(staleMarks)} документов в базе стоит отметка faiss_id, но вектора в индексе нет —
обычно это след пересоздания индекса после смены модели эмбеддингов.
Такие документы не участвуют в семантическом поиске и не будут пересчитаны,
пока отметку не сбросить: <code>scripts/ops/faiss_reconcile.py</code>.
</span>
</div>
)}
</Card> </Card>
{/* Проблемные прогоны */} {/* Проблемные прогоны */}
@@ -189,6 +237,7 @@ export function Debug() {
<span className={r.status === 'error' ? 'text-red-600' : 'text-amber-600'}>{r.status}</span> <span className={r.status === 'error' ? 'text-red-600' : 'text-amber-600'}>{r.status}</span>
<span className="text-gray-500"> <span className="text-gray-500">
получено {r.fetched}/{r.target}, добавлено {r.added}, ошибок {r.failed} получено {r.fetched}/{r.target}, добавлено {r.added}, ошибок {r.failed}
{r.duration_s != null && ` · работал ${Math.round(r.duration_s)}с`}
</span> </span>
{r.error && <span className="w-full text-xs font-mono text-red-600 break-all">{r.error}</span>} {r.error && <span className="w-full text-xs font-mono text-red-600 break-all">{r.error}</span>}
</div> </div>
@@ -234,6 +283,31 @@ export function Debug() {
); );
} }
/** Что из инфраструктуры доступно прямо сейчас — первый вопрос при разборе аварии. */
function HealthCard({ health }: { health?: Health[] }) {
return (
<Card icon={Server} title="Инфраструктура">
{health?.length ? (
<div className="grid grid-cols-1 sm:grid-cols-2 lg:grid-cols-3 gap-x-6 gap-y-1">
{health.map((h) => (
<div key={h.name} className="flex items-center justify-between gap-3 py-1 border-b border-gray-50">
<span className="flex items-center gap-2 text-sm text-gray-700">
{h.ok
? <CheckCircle2 className="w-4 h-4 text-emerald-500 shrink-0" />
: <XCircle className="w-4 h-4 text-red-500 shrink-0" />}
{h.name}
</span>
<span className={`text-xs truncate ${h.ok ? 'text-gray-400' : 'text-red-500'}`} title={h.detail}>
{h.detail || (h.ok ? 'OK' : 'недоступен')}
</span>
</div>
))}
</div>
) : <p className="text-sm text-gray-400">Опрос сервисов…</p>}
</Card>
);
}
function Card({ icon: Icon, title, children }: { icon: React.ElementType; title: string; children: React.ReactNode }) { function Card({ icon: Icon, title, children }: { icon: React.ElementType; title: string; children: React.ReactNode }) {
return ( return (
<div className="bg-white rounded-xl border border-gray-100 p-5"> <div className="bg-white rounded-xl border border-gray-100 p-5">

View File

@@ -13,8 +13,9 @@ interface Run {
target: number; fetched: number; processed: number; target: number; fetched: number; processed: number;
added: number; duplicates: number; skipped: number; failed: number; added: number; duplicates: number; skipped: number; failed: number;
cancel_requested: boolean; error: string | null; cancel_requested: boolean; error: string | null;
started_at: string; heartbeat_at: string | null; finished_at: string | null; started_at: string; run_started_at: string | null;
percent: number; heartbeat_at: string | null; finished_at: string | null;
percent: number; queued_s: number | null; duration_s: number | null;
} }
interface LogEntry { ts: string; elapsed: number; level: string; msg: string } interface LogEntry { ts: string; elapsed: number; level: string; msg: string }
@@ -360,6 +361,14 @@ export function Sources() {
); );
} }
/** Секунды → «45с» / «12 мин» / «2 ч 5 мин»: в отладке важен порядок, не точность. */
function fmtDuration(seconds: number): string {
if (seconds < 90) return `${Math.round(seconds)}с`;
const min = Math.round(seconds / 60);
if (min < 90) return `${min} мин`;
return `${Math.floor(min / 60)} ч ${min % 60} мин`;
}
const LOG_COLORS: Record<string, string> = { const LOG_COLORS: Record<string, string> = {
error: 'text-red-600', error: 'text-red-600',
warning: 'text-amber-600', warning: 'text-amber-600',
@@ -389,7 +398,9 @@ function RunLog({ runId, live }: { runId?: number; live: boolean }) {
<Chip label="дублей" value={String(data.duplicates)} /> <Chip label="дублей" value={String(data.duplicates)} />
{data.skipped > 0 && <Chip label="без метаданных" value={String(data.skipped)} />} {data.skipped > 0 && <Chip label="без метаданных" value={String(data.skipped)} />}
{data.failed > 0 && <Chip label="ошибок" value={String(data.failed)} />} {data.failed > 0 && <Chip label="ошибок" value={String(data.failed)} />}
<Chip label="начат" value={new Date(data.started_at).toLocaleString('ru-RU')} /> <Chip label="поставлен в очередь" value={new Date(data.started_at).toLocaleString('ru-RU')} />
{data.queued_s != null && <Chip label="ждал очереди" value={fmtDuration(data.queued_s)} />}
{data.duration_s != null && <Chip label="работал" value={fmtDuration(data.duration_s)} />}
{data.finished_at && <Chip label="завершён" value={new Date(data.finished_at).toLocaleString('ru-RU')} />} {data.finished_at && <Chip label="завершён" value={new Date(data.finished_at).toLocaleString('ru-RU')} />}
</div> </div>

View File

@@ -9,7 +9,7 @@ celery_app = Celery(
"worker_gpu", "worker_gpu",
broker=settings.RABBITMQ_URL, broker=settings.RABBITMQ_URL,
backend=settings.REDIS_URL, backend=settings.REDIS_URL,
include=["app.tasks.search", "app.tasks.plagiarism"], include=["app.tasks.search", "app.tasks.plagiarism", "app.tasks.stats"],
) )
celery_app.conf.update( celery_app.conf.update(

View File

@@ -0,0 +1,39 @@
"""Служебная задача: состояние векторного индекса для панели отладки.
Индекс живёт в памяти и на диске worker-gpu, у API к нему доступа нет. Без этой
задачи админка показывала покрытие L3 по колонке `documents.faiss_id` — а она
врёт: пометка остаётся и тогда, когда вектор в индекс не попал (или индекс был
пересоздан после смены модели эмбеддингов). Здесь возвращается то, что в индексе
есть на самом деле.
"""
from typing import Any
from celery.utils.log import get_task_logger
from app.celery_app import celery_app
from app.config import settings
from app.vector_store import get_backend
logger = get_task_logger(__name__)
@celery_app.task(name="gpu.index_stats")
def index_stats() -> dict[str, Any]:
"""Реальное число векторов в активном бэкенде (FAISS или Qdrant)."""
backend = get_backend()
try:
# У FAISS индекс ленивый: без обращения ntotal вернёт 0 на холодном воркере
if hasattr(backend, "_ensure"):
backend._ensure()
total = backend.total_vectors()
except Exception as e:
logger.warning(f"index_stats: не удалось прочитать индекс: {e}")
return {"backend": settings.VECTOR_BACKEND, "error": str(e)[:200]}
return {
"backend": settings.VECTOR_BACKEND,
"vectors": total,
"dim": settings.EMBED_DIM,
"embed_model": settings.EMBED_MODEL,
}

View File

@@ -101,6 +101,7 @@ class ParseRun(Base):
error: Mapped[str | None] = mapped_column(Text, nullable=True) error: Mapped[str | None] = mapped_column(Text, nullable=True)
log: Mapped[list | None] = mapped_column(JSON, default=list) log: Mapped[list | None] = mapped_column(JSON, default=list)
started_at: Mapped[datetime] = mapped_column(server_default=func.now()) started_at: Mapped[datetime] = mapped_column(server_default=func.now())
run_started_at: Mapped[datetime | None] = mapped_column(nullable=True)
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True) heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(nullable=True) finished_at: Mapped[datetime | None] = mapped_column(nullable=True)

View File

@@ -347,6 +347,53 @@ def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[
return {"status": "indexed", "doc_id": doc_id} return {"status": "indexed", "doc_id": doc_id}
def store_full_text(doc_id: int, text: str) -> dict[str, Any]:
"""Сохранить полный текст документа и переиндексировать его по нему.
Общая часть для всех путей получения полного текста: скачанный PDF
(enrich_full_text) и текст, пришедший прямо из API источника (бэкфилл PMC —
scripts/ops/). Провизорные fingerprints, посчитанные по аннотации,
заменяются на посчитанные по полному тексту — ради этого всё и делается:
L1 начинает видеть тело статьи, а не только её краткое описание.
Args:
doc_id: документ в PostgreSQL
text: полный текст статьи
Returns:
dict со статусом, объёмом текста и числом отпечатков
"""
from sqlalchemy import delete
from app.models import Document, Fingerprint
minio = get_minio()
key = f"corpus/{doc_id}.txt"
data = text.encode("utf-8")
minio.put_object(
settings.MINIO_BUCKET_DOCS, key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8",
)
hashes = list(winnow(text))[: settings.MAX_FINGERPRINTS_PER_DOC]
with db_session() as session:
doc = session.get(Document, doc_id)
if doc is None:
return {"status": "doc_gone", "doc_id": doc_id}
doc.minio_key = key
session.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id))
for i, hash_val in enumerate(hashes):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
session.commit()
# MinHash LSH (L2) тоже должен считаться по полному тексту
add_to_lsh(f"doc:{doc_id}", text)
logger.info(f"store_full_text: doc={doc_id} {len(text)} симв., fingerprints={len(hashes)}")
return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)}
@celery_app.task( @celery_app.task(
name="index.enrich_full_text", name="index.enrich_full_text",
bind=True, bind=True,
@@ -366,51 +413,17 @@ def enrich_full_text(self, doc_id: int, url: str) -> dict[str, Any]:
Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext. Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext.
""" """
from app.fulltext import fetch_full_text from app.fulltext import fetch_full_text
from app.models import Document, Fingerprint
text = fetch_full_text(url) text = fetch_full_text(url)
if not text: if not text:
return {"status": "no_fulltext", "doc_id": doc_id} return {"status": "no_fulltext", "doc_id": doc_id}
# Сохранить полный текст в MinIO
try: try:
minio = get_minio() return store_full_text(doc_id, text)
key = f"corpus/{doc_id}.txt"
data = text.encode("utf-8")
minio.put_object(
settings.MINIO_BUCKET_DOCS, key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8",
)
except Exception as exc: except Exception as exc:
logger.error(f"enrich_full_text: не удалось сохранить текст в MinIO для {doc_id}: {exc}") logger.error(f"enrich_full_text: не удалось сохранить текст для {doc_id}: {exc}")
raise self.retry(exc=exc, countdown=120) from exc raise self.retry(exc=exc, countdown=120) from exc
# Пересчитать fingerprints по полному тексту
doc_fp = winnow(text)
hashes = list(doc_fp)[: settings.MAX_FINGERPRINTS_PER_DOC]
from sqlalchemy import delete
with db_session() as session:
doc = session.get(Document, doc_id)
if doc is None:
return {"status": "doc_gone", "doc_id": doc_id}
doc.minio_key = key
# Удалить провизорные fingerprints и записать новые
session.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id))
for i, hash_val in enumerate(hashes):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
session.commit()
# Обновить MinHash LSH по полному тексту
add_to_lsh(f"doc:{doc_id}", text)
logger.info(
f"enrich_full_text: doc={doc_id} полный текст {len(text)} симв., "
f"fingerprints={len(hashes)}"
)
return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)}
@celery_app.task( @celery_app.task(
name="index.auto_approve_submission", name="index.auto_approve_submission",
@@ -625,7 +638,15 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any
prog.stage = "fetch" prog.stage = "fetch"
prog.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}") prog.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}")
_write_run( _write_run(
run_id, status="running", celery_task_id=self.request.id, error=None, **prog.snapshot() run_id,
status="running",
celery_task_id=self.request.id,
error=None,
# Отдельно от started_at (постановка в очередь): при массовом запуске
# между ними часы ожидания, и без этой отметки «длительность прогона»
# в отладке показывала очередь, а не работу
run_started_at=datetime.now(UTC),
**prog.snapshot(),
) )
cancelled = False cancelled = False