Compare commits
6 Commits
f1a8d07cb3
...
2d38dfd1f4
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2d38dfd1f4 | ||
|
|
d7c004af76 | ||
|
|
d8630ca0f9 | ||
|
|
ea62f10829 | ||
|
|
4008b5019f | ||
|
|
99bf14fe6a |
@@ -232,7 +232,7 @@ 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` по сервисам: 132 теста на ядро детекции, скоринга,
|
||||
2. **Юнит-тесты** — `pytest` по сервисам: 150 тестов на ядро детекции, скоринга,
|
||||
парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры
|
||||
(БД/Redis/GPU/Ollama замоканы либо не нужны).
|
||||
|
||||
@@ -265,7 +265,7 @@ make test-one SVC=worker-gost # тесты одного сервиса
|
||||
| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 19 |
|
||||
| OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 |
|
||||
| Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 |
|
||||
| Шкала загрузки источников | `api/app/core/progress.py` | 7 |
|
||||
| Шкала загрузки и тайминги прогонов | `api/app/core/progress.py`, `schemas/admin.py` | 11 |
|
||||
|
||||
## Лицензия
|
||||
|
||||
|
||||
@@ -81,7 +81,10 @@
|
||||
`last_run_id` — ссылка на последний прогон.
|
||||
- **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/
|
||||
added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал
|
||||
событий (JSON). Из них админка рисует шкалу загрузки, см. §12.
|
||||
событий (JSON). Тайминги раздельные: `started_at` — постановка в очередь,
|
||||
`run_started_at` — реальный старт работы воркером (при массовом запуске между
|
||||
ними часы ожидания), `finished_at` — конец. Из них админка рисует шкалу
|
||||
загрузки, см. §12.
|
||||
- **staged_works** — пользовательские загрузки на модерацию перед добавлением в корпус.
|
||||
- **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.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.notify` | `notify.send_task_done`, `notify.send_verification` | worker-notifier |
|
||||
|
||||
@@ -218,7 +221,10 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
|
||||
Кнопки «Запустить всё» / «Остановить всё» — массовый старт и кооперативная отмена
|
||||
(флаг `cancel_requested`, воркер останавливается сам на ближайшем тике; уже
|
||||
начатый прогон не рвём посреди записи в базу).
|
||||
- **Панель отладки** (админка → «Отладка», `GET /api/admin/debug`): один срез —
|
||||
- **Панель отладки** (админка → «Отладка», `GET /api/admin/debug` + `/health`):
|
||||
один срез — доступность инфраструктуры (PostgreSQL/Redis/MinIO/ES/Ollama/
|
||||
брокер: этот блок показывается даже когда сам срез не собирается, потому что
|
||||
при аварии первый вопрос — что именно отвалилось),
|
||||
живые воркеры Celery и что именно они крутят, глубина очередей RabbitMQ
|
||||
(в т.ч. `unacked` и число потребителей), покрытие корпуса эмбеддингами,
|
||||
активные и проблемные прогоны, зависшие прогоны (нет heartbeat >10 мин),
|
||||
@@ -226,6 +232,18 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
|
||||
`VECTOR_BACKEND`).
|
||||
- **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) —
|
||||
защита от краш-лупа по `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. Безопасность
|
||||
|
||||
|
||||
@@ -55,6 +55,17 @@ Redis у нас — кэш/rate-limits/LSH-индекс (префикс `antipla
|
||||
| Векторный индекс | было | `VECTOR_BACKEND=qdrant` снимает (см. README) ✅ |
|
||||
| RabbitMQ .82 | да | мониторинг ловит падение ✅; кластер — по потребности |
|
||||
| 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 долгие таски
|
||||
|
||||
|
||||
@@ -1,13 +1,45 @@
|
||||
# Наполнение корпуса — 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` (проверенные пользователями работы, не источник для сравнения
|
||||
сами с собой — см. 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`),
|
||||
англоязычные научные статьи открытого доступа.
|
||||
- Массовое расширение по дисциплинам теперь двумя сидерами: `seed_ru_sources.py`
|
||||
@@ -61,6 +93,18 @@ python scripts/seed_ru_sources.py --apply --limit 500
|
||||
- **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и
|
||||
упавшие прогоны (см. 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` — это не ошибка: сработал бюджет времени
|
||||
(`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout`
|
||||
RabbitMQ. Остаток добирается повторным запуском источника.
|
||||
@@ -82,6 +126,6 @@ SELECT source, count(*) FROM documents WHERE source='cyberleninka'; -- > 0
|
||||
заголовки с меткой ru).
|
||||
- Для миллионов — **bulk** (снапшот OpenAlex на S3), а не постраничный API.
|
||||
- На масштабе обязателен `VECTOR_BACKEND=qdrant` (FAISS flat не тянет), а таблица
|
||||
`fingerprints` (уже ~88M строк на 165K доков — партиционирование стоит планировать
|
||||
`fingerprints` (уже ~113M строк на 177K доков — партиционирование стоит планировать
|
||||
заранее, не постфактум) потребует партиционирования. См.
|
||||
[ARCHITECTURE.md](ARCHITECTURE.md) и [DR-HA.md](DR-HA.md).
|
||||
|
||||
112
scripts/ops/backfill_cyberleninka_pdf.py
Executable file
112
scripts/ops/backfill_cyberleninka_pdf.py
Executable 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()
|
||||
124
scripts/ops/backfill_pmc_fulltext.py
Executable file
124
scripts/ops/backfill_pmc_fulltext.py
Executable 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
201
scripts/ops/bulk_ingest_pmc.py
Executable 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
101
scripts/ops/faiss_reconcile.py
Executable 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
140
scripts/ops/reembed_missing.py
Executable 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()
|
||||
31
services/api/alembic/versions/006_parse_run_started_at.py
Normal file
31
services/api/alembic/versions/006_parse_run_started_at.py
Normal 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")
|
||||
@@ -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:
|
||||
"""Живое состояние воркеров через Celery inspect (блокирующий вызов)."""
|
||||
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)
|
||||
vector_stats = await run_in_threadpool(_vector_index_stats)
|
||||
|
||||
rmq = await _rabbitmq_queues()
|
||||
queues = (
|
||||
@@ -915,8 +931,10 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
|
||||
"queues": queues,
|
||||
"corpus": {
|
||||
"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),
|
||||
"elasticsearch_documents": es_docs,
|
||||
},
|
||||
|
||||
@@ -75,7 +75,10 @@ class ParseRun(Base):
|
||||
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)
|
||||
# Постановка в очередь (строку создаёт API) и реальный старт выполнения:
|
||||
# между ними при массовом запуске проходят часы, мешать их нельзя
|
||||
started_at: Mapped[datetime] = mapped_column(server_default=func.now(), index=True)
|
||||
run_started_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
# Последний признак жизни: по нему видно зависший прогон (running, но тишина)
|
||||
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
|
||||
@@ -80,7 +80,9 @@ class ParseRunResponse(BaseModel):
|
||||
failed: int
|
||||
cancel_requested: bool = False
|
||||
error: str | None = None
|
||||
# started_at — постановка в очередь, run_started_at — реальный старт работы
|
||||
started_at: datetime
|
||||
run_started_at: datetime | None = None
|
||||
heartbeat_at: datetime | None = None
|
||||
finished_at: datetime | None = None
|
||||
|
||||
@@ -91,6 +93,26 @@ class ParseRunResponse(BaseModel):
|
||||
def percent(self) -> float:
|
||||
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):
|
||||
log: list[dict[str, Any]] = Field(default_factory=list)
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
"""Юнит-тесты шкалы загрузки источников (app.core.progress) — чистая логика."""
|
||||
"""Юнит-тесты шкалы загрузки и таймингов прогонов — чистая логика, без БД."""
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from app.core.progress import run_percent
|
||||
from app.schemas.admin import ParseRunResponse
|
||||
|
||||
|
||||
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"):
|
||||
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
|
||||
|
||||
@@ -1,18 +1,21 @@
|
||||
import React from 'react';
|
||||
import { useQuery } from '@tanstack/react-query';
|
||||
import {
|
||||
Activity, AlertTriangle, Cpu, Database, Layers, RefreshCw, Settings2, ListTree,
|
||||
Activity, AlertTriangle, CheckCircle2, Cpu, Database, Layers, RefreshCw,
|
||||
Server, Settings2, ListTree, XCircle,
|
||||
} from 'lucide-react';
|
||||
import { adminApi } from '../../api/client';
|
||||
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 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;
|
||||
started_at: string; run_started_at: string | null; heartbeat_at: string | null;
|
||||
queued_s: number | null; duration_s: number | null;
|
||||
}
|
||||
|
||||
interface DebugData {
|
||||
@@ -20,7 +23,8 @@ interface DebugData {
|
||||
celery: { workers: Worker[]; error?: string };
|
||||
queues: Record<string, { messages: number; unacked: number; consumers: number }> | { error: string };
|
||||
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;
|
||||
};
|
||||
sources: {
|
||||
@@ -42,24 +46,48 @@ function fmtNum(n: number | null | undefined): string {
|
||||
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() {
|
||||
const { data, isFetching, error } = useQuery({
|
||||
queryKey: ['admin-debug'],
|
||||
queryFn: () => adminApi.debug().then((r) => r.data as DebugData),
|
||||
refetchInterval: 5000,
|
||||
});
|
||||
// Отдельным запросом: когда отваливается инфраструктура (брокер, Ollama),
|
||||
// отладка должна первой показывать ЧТО именно недоступно, а не только следствия
|
||||
const { data: health } = useQuery({
|
||||
queryKey: ['admin-health'],
|
||||
queryFn: () => adminApi.health().then((r) => r.data as Health[]),
|
||||
refetchInterval: 10000,
|
||||
});
|
||||
|
||||
if (error) {
|
||||
return <div className="text-red-600 text-sm">Не удалось получить срез состояния: {String(error)}</div>;
|
||||
}
|
||||
if (!data) {
|
||||
return <div className="text-gray-400 text-sm">Сбор данных…</div>;
|
||||
if (error || !data) {
|
||||
return (
|
||||
<div className="space-y-4">
|
||||
{error
|
||||
? <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 queues = queuesErr ? {} : (data.queues as Record<string, { messages: number; unacked: number; consumers: number }>);
|
||||
const embedPercent = data.corpus.documents
|
||||
? (data.corpus.documents_embedded / data.corpus.documents) * 100
|
||||
// Покрытие L3 считаем по РЕАЛЬНОМУ содержимому индекса: пометка faiss_id в БД
|
||||
// остаётся и когда вектор туда не попал, и завышает картину в разы
|
||||
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;
|
||||
|
||||
return (
|
||||
@@ -72,6 +100,8 @@ export function Debug() {
|
||||
<span className="text-xs text-gray-400">срез от {new Date(data.generated_at).toLocaleTimeString('ru-RU')}</span>
|
||||
</div>
|
||||
|
||||
<HealthCard health={health} />
|
||||
|
||||
{/* Активные прогоны заливки */}
|
||||
<Card icon={Layers} title={`Заливка источников — активных прогонов: ${data.sources.active_runs.length}`}>
|
||||
{data.sources.stale_run_ids.length > 0 && (
|
||||
@@ -87,8 +117,9 @@ export function Debug() {
|
||||
<div className="space-y-3">
|
||||
{data.sources.active_runs.map((r) => (
|
||||
<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}
|
||||
<span className="text-gray-400"> · {runningFor(r)}</span>
|
||||
</div>
|
||||
<ProgressBar
|
||||
percent={r.percent}
|
||||
@@ -167,15 +198,32 @@ export function Debug() {
|
||||
<Card icon={Database} title="Корпус сравнения">
|
||||
<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_embedded)} />
|
||||
<Metric label="Векторов в индексе" value={fmtNum(vectors)} />
|
||||
<Metric label="Отпечатков (оценка)" value={fmtNum(data.corpus.fingerprints_estimate)} />
|
||||
<Metric label="В Elasticsearch" value={fmtNum(data.corpus.elasticsearch_documents)} />
|
||||
</div>
|
||||
{data.corpus.vector_index?.error && (
|
||||
<p className="text-xs text-red-600 mb-2">
|
||||
индекс недоступен: {data.corpus.vector_index.error}
|
||||
</p>
|
||||
)}
|
||||
<ProgressBar
|
||||
percent={embedPercent}
|
||||
status={data.corpus.documents_without_embedding > 0 ? 'partial' : 'done'}
|
||||
label={`покрытие эмбеддингами (L3): без вектора ${fmtNum(data.corpus.documents_without_embedding)} документов`}
|
||||
status={embedPercent >= 99 ? 'done' : 'partial'}
|
||||
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>
|
||||
|
||||
{/* Проблемные прогоны */}
|
||||
@@ -189,6 +237,7 @@ export function Debug() {
|
||||
<span className={r.status === 'error' ? 'text-red-600' : 'text-amber-600'}>{r.status}</span>
|
||||
<span className="text-gray-500">
|
||||
получено {r.fetched}/{r.target}, добавлено {r.added}, ошибок {r.failed}
|
||||
{r.duration_s != null && ` · работал ${Math.round(r.duration_s)}с`}
|
||||
</span>
|
||||
{r.error && <span className="w-full text-xs font-mono text-red-600 break-all">{r.error}</span>}
|
||||
</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 }) {
|
||||
return (
|
||||
<div className="bg-white rounded-xl border border-gray-100 p-5">
|
||||
|
||||
@@ -13,8 +13,9 @@ interface Run {
|
||||
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;
|
||||
started_at: string; run_started_at: string | null;
|
||||
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 }
|
||||
@@ -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> = {
|
||||
error: 'text-red-600',
|
||||
warning: 'text-amber-600',
|
||||
@@ -389,7 +398,9 @@ function RunLog({ runId, live }: { runId?: number; live: boolean }) {
|
||||
<Chip label="дублей" value={String(data.duplicates)} />
|
||||
{data.skipped > 0 && <Chip label="без метаданных" value={String(data.skipped)} />}
|
||||
{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')} />}
|
||||
</div>
|
||||
|
||||
|
||||
@@ -9,7 +9,7 @@ celery_app = Celery(
|
||||
"worker_gpu",
|
||||
broker=settings.RABBITMQ_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(
|
||||
|
||||
39
services/worker-gpu/app/tasks/stats.py
Normal file
39
services/worker-gpu/app/tasks/stats.py
Normal 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,
|
||||
}
|
||||
@@ -101,6 +101,7 @@ class ParseRun(Base):
|
||||
error: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
log: Mapped[list | None] = mapped_column(JSON, default=list)
|
||||
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)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
|
||||
|
||||
@@ -347,6 +347,53 @@ def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[
|
||||
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(
|
||||
name="index.enrich_full_text",
|
||||
bind=True,
|
||||
@@ -366,51 +413,17 @@ def enrich_full_text(self, doc_id: int, url: str) -> dict[str, Any]:
|
||||
Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext.
|
||||
"""
|
||||
from app.fulltext import fetch_full_text
|
||||
from app.models import Document, Fingerprint
|
||||
|
||||
text = fetch_full_text(url)
|
||||
if not text:
|
||||
return {"status": "no_fulltext", "doc_id": doc_id}
|
||||
|
||||
# Сохранить полный текст в MinIO
|
||||
try:
|
||||
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",
|
||||
)
|
||||
return store_full_text(doc_id, text)
|
||||
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
|
||||
|
||||
# Пересчитать 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(
|
||||
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.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}")
|
||||
_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
|
||||
|
||||
Reference in New Issue
Block a user