diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 621a376..f1be38e 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -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** — одноразовые коды входа в админку. @@ -226,6 +229,12 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica `VECTOR_BACKEND`). - **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) — защита от краш-лупа по `consumer_timeout` RabbitMQ, см. [DR-HA.md](DR-HA.md) §6. +- **Добор эмбеддингов** — [`scripts/ops/reembed_missing.py`](../scripts/ops/reembed_missing.py): + документ попадает в корпус сразу, а вектор для L3 считает отдельная задача + `gpu.embed_documents`; если worker-gpu или Ollama были недоступны, эти задачи + теряются и документ остаётся невидимым для семантического поиска. Скрипт + находит `faiss_id IS NULL` и переотправляет задачи пачками (dry-run по + умолчанию). Покрытие видно в панели отладки. ## 13. Безопасность diff --git a/docs/DR-HA.md b/docs/DR-HA.md index eab7c09..3174195 100644 --- a/docs/DR-HA.md +++ b/docs/DR-HA.md @@ -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 долгие таски diff --git a/docs/INGESTION.md b/docs/INGESTION.md index baf0fa3..a241321 100644 --- a/docs/INGESTION.md +++ b/docs/INGESTION.md @@ -1,13 +1,18 @@ # Наполнение корпуса — 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` или новые темы, а не повторный запуск. +- OpenAlex при массовом запуске упирается в лимит вежливого пула (10 req/s + на mailto, общий для всех воркеров) — пауза между страницами поднята до 1с. - Добавлен 4-й парсер — **PMC** (PubMed Central, `scripts/parsers/pmc.py`), англоязычные научные статьи открытого доступа. - Массовое расширение по дисциплинам теперь двумя сидерами: `seed_ru_sources.py` @@ -61,6 +66,11 @@ python scripts/seed_ru_sources.py --apply --limit 500 - **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и упавшие прогоны (см. ARCHITECTURE.md §12). +После большой заливки стоит свериться с покрытием L3 (панель отладки, строка +«без вектора»): эмбеддинги считаются отдельной задачей на worker-gpu, и если он +или Ollama были недоступны, документы останутся без вектора. Догнать — +`scripts/ops/reembed_missing.py` (dry-run по умолчанию, `--apply` отправляет). + Прогон со статусом `partial` — это не ошибка: сработал бюджет времени (`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout` RabbitMQ. Остаток добирается повторным запуском источника. @@ -82,6 +92,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). diff --git a/scripts/ops/reembed_missing.py b/scripts/ops/reembed_missing.py new file mode 100755 index 0000000..30033f4 --- /dev/null +++ b/scripts/ops/reembed_missing.py @@ -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() diff --git a/services/api/alembic/versions/006_parse_run_started_at.py b/services/api/alembic/versions/006_parse_run_started_at.py new file mode 100644 index 0000000..4458895 --- /dev/null +++ b/services/api/alembic/versions/006_parse_run_started_at.py @@ -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") diff --git a/services/api/app/models/admin.py b/services/api/app/models/admin.py index 6dabc26..cad4802 100644 --- a/services/api/app/models/admin.py +++ b/services/api/app/models/admin.py @@ -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) diff --git a/services/api/app/schemas/admin.py b/services/api/app/schemas/admin.py index abbdcc5..ff433f0 100644 --- a/services/api/app/schemas/admin.py +++ b/services/api/app/schemas/admin.py @@ -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) diff --git a/services/frontend/src/pages/admin/Sources.tsx b/services/frontend/src/pages/admin/Sources.tsx index 21acd65..6f72564 100644 --- a/services/frontend/src/pages/admin/Sources.tsx +++ b/services/frontend/src/pages/admin/Sources.tsx @@ -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 = { error: 'text-red-600', warning: 'text-amber-600', @@ -389,7 +398,9 @@ function RunLog({ runId, live }: { runId?: number; live: boolean }) { {data.skipped > 0 && } {data.failed > 0 && } - + + {data.queued_s != null && } + {data.duration_s != null && } {data.finished_at && } diff --git a/services/worker-indexer/app/models/__init__.py b/services/worker-indexer/app/models/__init__.py index e4a5a86..1143a1a 100644 --- a/services/worker-indexer/app/models/__init__.py +++ b/services/worker-indexer/app/models/__init__.py @@ -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) diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index 0eb629d..ac8a81f 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -625,7 +625,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