Compare commits

...

18 Commits

Author SHA1 Message Date
jze9
d137074ed5 fix(core): ходить в CORE через sing-box — с российских адресов ключ не работает
All checks were successful
Deploy / test (push) Successful in 3m15s
Deploy / deploy (push) Successful in 3m43s
Заливка встала после двух порций: прогоны стали заканчиваться за две минуты
с нулём документов. Выглядело как поломка сети, и на ложные следы ушло время —
MTU в норме (1472 байта проходят), IPv6 ни при чём, Cloudflare из того же
контейнера качается на 2.2 МБ/с, ключ и квота целы (с домашней машины тот же
запрос отвечает за 6с, лимит нетронут).

Разница оказалась в адресе. С прода:
  напрямую      — код 000, обрыв на 25с (соединение есть, тело не приходит)
  через sing-box — код 200 за 1.7с, 207 КБ
Анонимные короткие запросы проходят и напрямую — поэтому блокировка и
маскировалась под сетевой сбой.

CORE ведёт себя как OpenRouter, ради которого singbox-proxy и заводили,
поэтому решение то же: httpx получает proxy из CORE_PROXY_URL (умолчание —
socks5://singbox-proxy:1080, пусто = напрямую для локальных прогонов).
В образ индексатора добавлен socksio: без него httpx не умеет SOCKS5.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-17 15:48:43 +05:00
jze9
9123844a87 fix(ops): не хоронить источник после одного пустого прогона
All checks were successful
Deploy / test (push) Successful in 11m54s
Deploy / deploy (push) Successful in 8s
Сторож считал источник исчерпанным, если прошлый прогон не добавил и не
выбрал ни одной статьи. Но ровно так же выглядит обрыв связи с API и прогон
на устаревшем коде парсера: сегодня CORE после единственной неудачи был
помечен исчерпанным и больше не запускался — молча, без ошибки в интерфейсе.

Теперь нужно два пустых прогона подряд. Разовый сбой переживём, а реально
кончившийся источник остановится всего на один прогон позже.

Заодно таймаут запроса к CORE снижен с 60 до 30 секунд: три попытки по
минуте отъедали 180 секунд из 1500 бюджета на одной залипшей странице —
та же грабля, что чинили в pmc_bulk. Обычный ответ приходит за 7 секунд.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 18:34:14 +05:00
jze9
263a521d63 fix(ops): перезапускать worker-indexer, когда менялись парсеры
Some checks failed
Deploy / deploy (push) Has been cancelled
Deploy / test (push) Has been cancelled
Парсеры не вшиты в образ, а смонтированы в worker-indexer, поэтому деплой
их не пересобирал — и не перезапускал контейнер, раз `services/` не менялся.
Но воркер держит модуль парсера в памяти с прошлого прогона: новый файл
лежит в контейнере, а работает старый код.

Поймано вживую на CORE: заливка после деплоя ушла ноль документов за 208
секунд, и только тело запросов в логах показало, что она всё ещё ходит по
старой схеме — с огромным OR-запросом, который на глубоком смещении трижды
упал по таймауту. Со стороны выглядело как «источник исчерпан», и сторож
чуть не выключил источник насовсем.

Теперь при изменениях в `scripts/parsers/` деплой перезапускает
worker-indexer — если образ и так не пересобирается на этом прогоне.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 18:32:11 +05:00
jze9
7adccd0362 fix(core): резать выборку по архивам, а не по годам — AND в API не работает
All checks were successful
Deploy / test (push) Successful in 7m41s
Deploy / deploy (push) Successful in 4s
Первый прогон на проде выдал документы 2008-2019 годов, хотя запрос просил
`yearPublished:2026`. Проверка показала: условия в запросе CORE через `AND`
не связываются. По отдельности каждое фильтрует честно (`yearPublished:2026`
— ровно 2026, `repositories.id:1298` — ровно этот архив), а вместе
`(архивы) AND yearPublished:2026` отдаёт 2.3 млн работ вперемешку, то есть
условия объединяются по «или», и год работает лишь подсказкой ранжированию.

Значит нарезка по годам не нарезала ничего: каждый «год» перебирал один и тот
же набор, а в выдачу подмешивались посторонние работы нужного года — включая
англоязычные, ради ухода от которых источник и заводился.

Теперь в запросе ровно одно условие — номер архива, и каждый архив
опрашивается отдельно. Позиция продолжения стала «архив:смещение». Это ещё и
честнее по потолку: `offset` упирается в 100 000, а самый крупный из наших
архивов содержит 63 749 работ, то есть влезает целиком.

Годы, если заданы в источнике, отсекаются теперь на нашей стороне. Список
архивов парсер принимает и простым списком номеров, и прежним выражением
`(repositories.id:N OR ...)` — настройку источника менять не нужно.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 18:15:38 +05:00
jze9
7348051c7b chore(ci): перезапустить деплой после чистки сервера Gitea
All checks were successful
Deploy / test (push) Successful in 8m7s
Deploy / deploy (push) Successful in 24s
Прогон №50 упал не на коде: сервер Gitea подмешивал в поток git-протокола
вывод постороннего скрипта, и клонирование рвалось с `fatal: early EOF`.
Причина устранена на сервере, клонирование проверено — пустой коммит нужен
только чтобы дать раннеру новое задание.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 17:56:43 +05:00
jze9
3d4fcecbb2 feat(core): подключить CORE как источник полных текстов
Some checks failed
Deploy / test (push) Failing after 6s
Deploy / deploy (push) Has been skipped
Оба массовых источника исчерпаны — wikipedia_ru и pmc_bulk отдают ноль,
сторож честно пишет «источник исчерпан, пропускаем», и корпус стоит.
Брать новые статьи неоткуда, а главная дыра прежняя: полный текст есть у
11% корпуса, по-русски — почти нигде (CyberLeninka отдаёт только HTML
страницы, тела статьи там нет вовсе).

CORE кладёт готовый текст прямо в выдачу поиска: его не надо ни качать
отдельно, ни извлекать из PDF. Парсер массовый, как Википедия и PMC —
отдаёт статьи генератором, пишется пачками через COPY, помнит позицию.

Что выяснено живьём и учтено в коде:

- у /v3/search/works обязателен слэш на конце, иначе 301, а на редиректе
  теряется заголовок Authorization и запрос уходит анонимным;
- offset упирается в 100 000 (под капотом Azure Search), поэтому выборка
  режется по годам, а позиция продолжения — строка «год:смещение»;
- fullText не фильтруемое поле, _exists_:fullText отвечает 500 — статьи с
  текстом приходится отбирать на своей стороне;
- CORE отдаёт одну статью под разными id из разных репозиториев, а корпус
  дедуплицируется только по ext_id: такие пары легли бы отдельными
  документами и ловились бы как заимствование друг у друга. Отсеиваем по
  хешу начала текста в пределах прогона.

Позиция указывает на страницу, из которой пришла статья, а не на
следующую: пачка может прерваться на середине страницы, и тогда прогон
перечитает её целиком. Лишний запрос дешевле потерянных статей, дубли
отсекает ON CONFLICT.

Замер на живом API: страница в 100 записей приходит за ~7с. В общем
потоке CORE кириллицы нет вовсе (0 из 48 полных текстов за 2019), язык
размечен негодно — сплошь None и «zz». Зато работает отбор по архиву:
запрос по 41 вузовскому репозиторию России и Беларуси даёт 79 полных
текстов из 100 записей, 78 из них кириллические. Список архивов живёт в
поле query источника, а не в коде — правится из админки без пересборки.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 17:05:34 +05:00
jze9
20c1d00443 fix(pmc): не давать залипшему запросу съесть весь бюджет прогона
Таймаут запроса к бакету был 60с при бюджете таска 1500с — один
залипший GET отъедал почти весь прогон. Плюс map() отдаёт результаты
строго по порядку отправки, поэтому один медленный запрос блокировал
все 11 уже готовых потоков.

Таймаут снижен до 12с (обычный GET укладывается в доли секунды, 12с —
запас на джиттер), map() заменён на as_completed: готовые результаты
отдаются сразу, не дожидаясь залипшего соседа.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-09-15 13:39:45 +05:00
jze9
5ceec1e2e3 fix(ops): сторож не снимает живой прогон на долгом COPY
Молчащий heartbeat сторож трактовал как смерть прогона и снимал его по
порогу 10 минут. Но запись большой пачки отпечатков в базу на HDD идёт
COPY'ем 6-15 минут без единого тика heartbeat — и сторож убивал вполне
живой прогон, тут же запуская дубль по тем же статьям. Корпус часами
топтался на месте: из лога видно десятки перезапусков подряд с нулевым
приростом.

Теперь перед снятием сторож спрашивает воркеров queue.index через
inspect().active(), выполняется ли ещё таск прогона. Снимаются только
настоящие зомби — те, чьего celery_task_id нет ни на одном воркере.
Если хоть один воркер не ответил, живость не проверить и не снимается
ничего (fail-safe). Проверено на живом прогоне: его таск попал в набор
активных, прогон помечен защищённым.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-09-15 13:39:36 +05:00
jze9
d86bb606c7 fix(wikipedia): позиция — байтовое смещение, а не номер статьи
All checks were successful
Deploy / test (push) Successful in 3m8s
Deploy / deploy (push) Successful in 30s
Заливка Википедии останавливалась сама собой: позиция хранилась как номер
статьи, и каждый прогон перечитывал дамп с начала. На 20 тысячах это стоило
6 минут из 25 доступных, на 27 тысячах — уже около десяти, а на сотне тысяч
съело бы весь бюджет и заливка встала бы совсем. Ровно это и наблюдалось:
прогон висел с нулём полученных статей, воркер на 100% CPU.

Переход на multistream-вариант дампа: он состоит из независимых bz2-блоков по
~95 статей, к нему прилагается индекс со смещениями. Позиция продолжения стала
байтовым смещением, прогон стартует мгновенно с нужного места.

Проверено на реальных файлах, а не по предположению: формат индекса
(offset:page_id:title), 21 уникальное смещение на 2001 статью, прыжок seek на
смещение из середины файла даёт валидный XML со страницами.

wikipedia_resume_offset.py — разовый пересчёт позиции при переходе: находит по
индексу блок с максимальным залитым page_id (получилось 273 281 821).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 23:26:16 +05:00
jze9
d3106efeee feat(ops): заливка идёт сама — автоперезапуск порций по cron
All checks were successful
Deploy / test (push) Successful in 3m8s
Deploy / deploy (push) Successful in 4s
Массовый источник за прогон берёт порцию, сохраняет позицию и останавливается
по бюджету времени. Без внешнего толчка заливка шла рывками — ровно столько,
сколько раз кто-то нажмёт «Запустить», и всё это время корпус стоял на месте.

keep_ingesting.py раз в 5 минут проверяет, не простаивают ли массовые
источники, и запускает следующую порцию. Идущие прогоны не трогает.
Останавливается сам, когда источник исчерпан (прогон закрылся с нулём
полученных статей), чтобы не крутить пустые запуски вечно.

Поставлен в cron на app-хосте рядом с монитором.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 21:56:48 +05:00
jze9
eb2dee9093 perf(ops): начинать листинг PMC после последней залитой статьи
All checks were successful
Deploy / test (push) Successful in 3m4s
Deploy / deploy (push) Successful in 29s
Первый прогон нового конвейера показал слабое место: листинг бакета шёл с
начала, и заливка часами перемалывала уже существующие статьи как дубли —
из 300 полученных 300 оказались дублями.

S3 умеет start-after, и стартовый ключ выводится прямо из базы: ext_id вида
`pmc:PMC10000000` соответствует ключу `PMC10000000.1`, бакет отдаётся
лексикографически. Теперь при отсутствии сохранённого токена (первый прогон
после ручных заливок) листинг начинается после максимального залитого.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 21:53:07 +05:00
jze9
a76e46c561 feat(ops): замер качества детекции — первые честные цифры
All checks were successful
Deploy / deploy (push) Successful in 4s
Deploy / test (push) Successful in 2m53s
О качестве проверки мы до сих пор знали только «механизм жив»: находит
подброшенный фрагмент. Продукт при этом продаёт процент заимствований, за
который никто не ручался — неизвестно было ни сколько списываний система
пропускает, ни как часто обвиняет невиновных.

Бенчмарк делает из документов корпуса «студенческие работы» четырёх видов и
гоняет их через настоящий путь L1. Первый замер на проде:

  дословно        100% найдено
  лёгкий рерайт   100% найдено
  сильный рерайт    0% — это работа L3/L4, не L1
  оригинал          0% ложных обвинений

То есть основа работает как задумано. Показательно другое: из 20 взятых
документов в замер попали 8 — у остальных нет полного текста нужной длины.
Узкое место не алгоритм, а глубина корпуса.

Запускать после изменения порогов и параметров winnowing — иначе непонятно,
улучшение сделано или ухудшение.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 20:45:59 +05:00
jze9
fed4b54cfc fix(ops): мониторинг под реальную топологию + проверка воркеров
Some checks failed
Deploy / test (push) Successful in 2m52s
Deploy / deploy (push) Has been cancelled
Монитор лежал только на сервере, не версионировался и месяц проверял
конфигурацию, которой уже нет: Ollama на CT 108 (.20.163), остановленном
27.08 при переходе на OpenRouter. Итог — вечный ложный DOWN, на фоне которого
теряются настоящие аварии, и полная слепота к реальному серверу эмбеддингов
(.1.40) и к воркерам.

- скрипт переехал в репозиторий и теперь деплоится вместе с кодом;
- адреса читаются из прод-.env, а не зашиты: сервер эмбеддингов переезжал
  дважды, и каждый раз монитор оставался со старым адресом;
- добавлена проверка воркеров Celery — 05.09 брокер лежал час, воркеры молчали,
  а монитор рапортовал, что всё хорошо, потому что смотрел только TCP-порт;
- разбор URL чинит падение на кредах в REDIS_URL/RABBITMQ_URL.

Проверено на проде: все восемь проверок отдают OK, включая Embeddings и
CeleryWorkers.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 20:43:01 +05:00
jze9
26de1b9de1 refactor(ops): массовая заливка — в общий конвейер вместо отдельных скриптов
All checks were successful
Deploy / deploy (push) Successful in 2m28s
Deploy / test (push) Successful in 2m49s
Заливку Википедии и PMC я сделал отдельными скриптами мимо существующей
инфраструктуры: запуск руками через ssh, состояние в файле в /tmp, никакой
видимости. Результат предсказуем — за неделю обе умерли молча (обрыв базы на
20 008 статьях из 2 млн и таймаут сети на 85 тыс. из 100 тыс.), прогресс
потерялся, а узнали мы об этом через неделю. При том что рядом лежит готовый
механизм: parse_sources, прогоны со шкалой, журнал, кнопки, ретраи Celery.

Теперь это обычные типы источника — wikipedia_ru и pmc_bulk:

- заводятся и запускаются из админки, как OpenAlex или КиберЛенинка;
- показывают ту же шкалу, журнал и кнопку остановки;
- падение воркера больше не теряет прогресс: позиция продолжения хранится в
  parse_sources.resume_token (номер статьи в дампе / токен страницы бакета),
  повторный запуск берёт следующую порцию;
- укладываются в бюджет времени таска — заливка идёт порциями, а не одним
  многосуточным процессом.

Чего не хватало конвейеру для миллионов и что добавлено:
- парсеры отдают генератор, а не список: 2 млн статей в память не влезают;
- app/bulk_writer.py — запись пачками через COPY (21 тыс. строк/с против
  6.7 тыс. построчно) с переподключением к базе при обрыве;
- эмбеддинги при массовой заливке не диспатчатся: они на порядок медленнее и
  стали бы узким местом, вектора досчитываются отдельно (reembed_missing.py).

scripts/ops/bulk_ingest_*.py удалены — их работу делает конвейер.

Проверено на проде: оба парсера отдают документы, прогон через run_parser
завершается штатно, позиция продолжения сдвигается (300 → 600), повторный
запуск продолжает с неё.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 20:32:53 +05:00
jze9
1b1ca3b7e3 fix(admin): починить панель отладки — она отдавала 500
All checks were successful
Deploy / test (push) Successful in 2m59s
Deploy / deploy (push) Successful in 2m15s
Моя же правка про зависшие проверки уронила /admin/debug: колонки tasks в
PostgreSQL — timestamptz, а модель объявляет их без зоны, и asyncpg отверг
переданное из Python время («can't subtract offset-naive and offset-aware»).

Возраст задачи теперь считает сама база (now() - coalesce(updated_at,
created_at)), так что расхождение объявления и реального типа колонки роли не
играет. Проверено на живых данных: находит все 4 задачи, висящие с 30 мая.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 20:17:28 +05:00
jze9
39951d6a0f feat(admin): показывать зависшие проверки в панели отладки
All checks were successful
Deploy / test (push) Successful in 3m44s
Deploy / deploy (push) Successful in 8m39s
Нашлись 4 задачи в статусе processing, висящие с 30 мая: пользователь видит
вечное «обрабатывается», а в системе никаких следов. Теперь задачи без
движения дольше 2 часов попадают в /admin/debug рядом с зависшими прогонами
заливки, с указанием, сколько часов они стоят.

tasks.created_at/updated_at — timestamptz, поэтому сравнение идёт с
осведомлённым о зоне временем: naive utcnow() дал бы TypeError в рантайме.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
jze9
f4092e5331 fix(ops): заливки переживают обрывы базы и сети
За неделю обе массовые заливки умерли молча и не возобновились: Википедия на
20 008 статьях из ~2 млн («server closed the connection»), PMC примерно на
85 тыс. из 100 тыс. (таймаут SSL-рукопожатия). Ни ретраев, ни возобновления.

- запись пачки повторяется с переподключением к базе (5 попыток с паузой);
- обрыв листинга бакета стоит паузы, а не всей заливки;
- у Википедии сохраняется позиция в дампе — после сбоя продолжаем с неё,
  а не проматываем с нуля, полагаясь на дедуп по ext_id;
- батч уменьшен (500 статей × 2000 отпечатков = миллион строк в одной
  транзакции — вероятная причина обрыва);
- обрезка отпечатков переведена на равномерную выборку.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
jze9
476682741e fix(L1): обрезать отпечатки равномерно по тексту, а не произвольно
winnow() возвращает set, поэтому list(fp)[:LIMIT] брал случайное подмножество:
у длинного документа целые куски оставались без отпечатков, и списывание
именно из них не находилось. Обнаружено при разборе того, почему фрагмент
статьи PMC не искался.

Добавлены winnow_ordered() — отпечатки в порядке появления в тексте, и
sample_evenly() — выборка каждого n-го элемента вместо первых N. Применено в
add_document и store_full_text.

7 тестов на главное свойство: выборка растянута по всей длине документа, шаг
ровный, порядок сохранён.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
26 changed files with 2070 additions and 492 deletions

View File

@@ -152,6 +152,10 @@ services:
depends_on: depends_on:
elasticsearch: elasticsearch:
condition: service_healthy condition: service_healthy
# CORE не отвечает на ключ с российских адресов — парсер ходит через
# sing-box (см. scripts/parsers/core.py)
singbox-proxy:
condition: service_started
worker-notifier: worker-notifier:
build: build:

View File

@@ -78,7 +78,9 @@
- **fingerprints** — `doc_id`, `hash_value` (BIGINT, Winnowing), `position` — для L1. - **fingerprints** — `doc_id`, `hash_value` (BIGINT, Winnowing), `position` — для L1.
- **usage_logs** — `user_id`, `action` — учёт лимитов по тарифу. - **usage_logs** — `user_id`, `action` — учёт лимитов по тарифу.
- **parse_sources** — задания парсеров (админка): тип, query, годы, лимит, статус, - **parse_sources** — задания парсеров (админка): тип, query, годы, лимит, статус,
`last_run_id` — ссылка на последний прогон. `last_run_id` — ссылка на последний прогон, `resume_token` — позиция
продолжения для массовых источников (номер статьи в дампе, токен страницы
бакета).
- **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). Тайминги раздельные: `started_at` — постановка в очередь, событий (JSON). Тайминги раздельные: `started_at` — постановка в очередь,
@@ -119,10 +121,20 @@
`app.bibliography.build_bibliography` (сортировка кириллица→латиница, нумерация, `app.bibliography.build_bibliography` (сортировка кириллица→латиница, нумерация,
формат 7.1/7.0.5) → результат. формат 7.1/7.0.5) → результат.
**Наполнение корпуса.** админка/CLI → `index.run_parser` (OpenAlex/arXiv/PMC/ **Наполнение корпуса.** админка → `index.run_parser` — один конвейер для всех
КиберЛенинка, фильтр `is_oa`) → `index.add_document` (дедуп по `ext_id`, источников, со шкалой, журналом и кнопкой остановки. Внутри два пути записи:
fingerprints, MinHash, эмбеддинг) → `index.enrich_full_text` (скачать OA-PDF →
MinIO → переиндексация). Ход заливки пишется в `parse_runs` (§12). - *обычные источники* (OpenAlex/arXiv/PMC/КиберЛенинка, фильтр `is_oa`) —
`index.add_document` на каждый документ: дедуп по `ext_id`, fingerprints,
MinHash, Elasticsearch, эмбеддинг, затем `index.enrich_full_text` (скачать
OA-PDF → MinIO → переиндексация);
- *массовые источники* (`wikipedia_ru`, `pmc_bulk`) — парсер отдаёт генератор,
запись идёт пачками через `COPY` (`app/bulk_writer.py`), позиция продолжения
хранится в `parse_sources.resume_token`. Эмбеддинги там не считаются: они
медленнее заливки на порядок и стали бы её узким местом, вектора
досчитываются отдельно (§12).
Ход заливки в обоих случаях пишется в `parse_runs` (§12).
Второй путь наполнения — ручная загрузка файлов админом: api сохраняет их в Второй путь наполнения — ручная загрузка файлов админом: api сохраняет их в
MinIO (`corpus-upload/`) → `index.ingest_upload` (извлечь текст → `add_document`, MinIO (`corpus-upload/`) → `index.ingest_upload` (извлечь текст → `add_document`,
`source=manual_upload`). Это не проверка на плагиат: файл сразу становится `source=manual_upload`). Это не проверка на плагиат: файл сразу становится
@@ -202,6 +214,24 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
## 11. Качество и тесты ## 11. Качество и тесты
**Замер качества детекции** — [`scripts/ops/detection_benchmark.py`](../scripts/ops/detection_benchmark.py).
Делает из документов корпуса «студенческие работы» (дословная копия, лёгкий и
сильный рерайт, плюс заведомо оригинальный текст) и прогоняет через настоящий
путь L1. Первый замер, 05.09.2026:
| Случай | Доля найденных |
|--------|---------------:|
| дословно | 100% |
| лёгкий рерайт (выброшено каждое 10-е слово) | 100% |
| сильный рерайт (каждое 3-е слово) | 0% — задача L3/L4, не L1 |
| оригинальный текст | 0% ложных обвинений |
Запускать после изменения порогов (`EXACT_FRAGMENT_THRESHOLD`), параметров
winnowing и плотности отпечатков: без этих чисел непонятно, улучшение сделано
или ухудшение. Ограничение замера: в выборку попадают только документы с
реальной глубиной (>800 отпечатков) — по документам с одной аннотацией мерить
нечего, и это само по себе показатель состояния корпуса.
- **113 юнит-тестов** (pytest, per-service) на чистую логику L1-L4/скоринг/фрагменты/ - **113 юнит-тестов** (pytest, per-service) на чистую логику L1-L4/скоринг/фрагменты/
ГОСТ/библиография; инфра замокана или не нужна. Запуск: `make test`. ГОСТ/библиография; инфра замокана или не нужна. Запуск: `make test`.
- **Гейты CI**: ruff (весь Python) + mypy (доменная логика) + тесты — блокируют деплой. - **Гейты CI**: ruff (весь Python) + mypy (доменная логика) + тесты — блокируют деплой.
@@ -211,8 +241,13 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
## 12. Наблюдаемость и эксплуатация ## 12. Наблюдаемость и эксплуатация
- **Мониторинг**: `antiplag_monitor.py` (cron 5 мин, 7 сервисов, email-алерт при смене - **Мониторинг**: [`scripts/ops/antiplag_monitor.py`](../scripts/ops/antiplag_monitor.py)
статуса). Опционально — Prometheus+Grafana (профиль `observability`, метрики API + Flower). (cron 5 мин на 1.32, email при смене статуса). Адреса берутся из прод-`.env`, а не
зашиты в код: прежняя версия месяц проверяла остановленный CT 108 и не видела
реального сервера эмбеддингов — постоянный ложный DOWN заглушал настоящие аварии.
Проверяются в том числе **воркеры Celery**: живой брокер ещё не значит работающую
систему (05.09 воркеры час простаивали при «зелёном» брокере).
Опционально — Prometheus+Grafana (профиль `observability`, метрики API + Flower).
- **Бэкапы**: `pg_dump→gzip→MinIO`, cron 03:00, ротация 14. Проверка восстановления — - **Бэкапы**: `pg_dump→gzip→MinIO`, cron 03:00, ротация 14. Проверка восстановления —
`scripts/ops/pg_restore_verify.sh`. HA/DR — [DR-HA.md](DR-HA.md). `scripts/ops/pg_restore_verify.sh`. HA/DR — [DR-HA.md](DR-HA.md).
- **Шкала загрузки источников** (админка → «Источники»): каждый запуск создаёт строку - **Шкала загрузки источников** (админка → «Источники»): каждый запуск создаёт строку

View File

@@ -62,7 +62,7 @@
| Источник | Объём | Годен для массовой заливки | | Источник | Объём | Годен для массовой заливки |
|----------|-------|----------------------------| |----------|-------|----------------------------|
| **Википедия ru** | дамп 5.6 ГБ, ~2 млн статей | **да** — `bulk_ingest_wikipedia_ru.py`, потоком без блокировок, ~919 отпечатков на статью. Требует осмысленный User-Agent, иначе 403 | | **Википедия ru** | дамп 5.6 ГБ, ~2 млн статей | **да** — тип источника `wikipedia_ru`, ~919 отпечатков на статью. Дамп качается заранее (Wikimedia отдаёт 403 без осмысленного User-Agent и рвёт долгие потоковые соединения) |
| КиберЛенинка | ~3 млн статей | нет — блокирует выкачку PDF после ~130 запросов | | КиберЛенинка | ~3 млн статей | нет — блокирует выкачку PDF после ~130 запросов |
| OpenAlex `language:ru` + OA | заявлено 395 410 | практически нет — по прямым `pdf_url` скачалось 2 из 10 (остальное 403 издателей), а метка языка ненадёжна: в выдаче попадаются англоязычные журналы | | OpenAlex `language:ru` + OA | заявлено 395 410 | практически нет — по прямым `pdf_url` скачалось 2 из 10 (остальное 403 издателей), а метка языка ненадёжна: в выдаче попадаются англоязычные журналы |
| eLIBRARY.RU (РИНЦ) | ~40 млн | только по договору, публичной выгрузки нет | | eLIBRARY.RU (РИНЦ) | ~40 млн | только по договору, публичной выгрузки нет |
@@ -72,19 +72,36 @@
Википедия формально не научный источник, но студенты копируют из неё чаще всего, Википедия формально не научный источник, но студенты копируют из неё чаще всего,
а по объёму связного русского текста ей нет альтернативы среди доступного. а по объёму связного русского текста ей нет альтернативы среди доступного.
### Массовая заливка: bulk вместо API ### Массовые источники — тот же конвейер, что и обычные
Постраничные API дают 1-2 статьи в секунду и упираются в rate limit — миллионы Постраничные API дают 1-2 статьи в секунду и упираются в rate limit — миллионы
так не залить. Для объёма есть `scripts/ops/bulk_ingest_pmc.py`: открытый бакет так не залить. Для объёма есть два типа источника, которые читают дамп или
`pmc-oa-opendata` (ключ не нужен) отдаёт у каждой статьи готовый извлечённый бакет потоком:
текст. Проверено 31.08: 100 статей, 183 090 отпечатков — 1831 на статью, то есть
настоящая глубина, а не аннотация. | Тип | Откуда | Особенность |
|-----|--------|-------------|
| `wikipedia_ru` | локальный дамп `scripts/parsers/ruwiki.xml.bz2` | позиция продолжения — номер статьи в дампе |
| `pmc_bulk` | бакет `pmc-oa-opendata` (открыт, ключ не нужен) | позиция — токен страницы бакета; у каждой статьи готовый извлечённый текст |
Заводятся и запускаются они как любой другой источник — через админку, и точно
так же показывают шкалу, журнал и кнопку остановки. Отличия внутри:
- парсер отдаёт **генератор**, а не список: миллионы статей в память не влезут;
- запись идёт пачками через `COPY` (`worker-indexer/app/bulk_writer.py`) —
построчная вставка даёт 6.7 тыс. строк/с против 21 тыс. у COPY, а миллион
статей это ~2 млрд отпечатков;
- **эмбеддинги при заливке не считаются**: они медленнее заливки и сделали бы её
узким местом. L1 и L2 работают сразу, векторы досчитываются потом
(`scripts/ops/reembed_missing.py`);
- позиция продолжения хранится в `parse_sources.resume_token`, поэтому источник
запускается повторно до исчерпания — каждый прогон берёт следующую порцию и
укладывается в бюджет времени таска.
Планировать объём: 1 млн статей ≈ 2 млрд отпечатков ≈ 200 ГБ в базе с индексами.
При таком росте индексы перестают помещаться в память сервера БД — проверено на
практике: поиск L1 деградировал с 31 мс до 484 мс, пока не увеличили RAM и
`shared_buffers` (см. DR-HA.md).
Ограничение по темпу: 0.3 статьи/с, миллион в один поток — около 40 суток.
Скачивание тут не узкое место (12 статей/с в 12 потоков), упирается в
последовательную обработку документа; для миллионов её надо распараллелить по
ядрам воркера. Планировать объём: 1 млн статей ≈ 2.5 млрд отпечатков ≈ 400 ГБ
в базе с индексами (сейчас 113 млн строк занимают 12 ГБ).
## Что подготовлено ## Что подготовлено
- **Парсер CyberLeninka починен** (`scripts/parsers/cyberleninka.py`): раньше слал GET - **Парсер CyberLeninka починен** (`scripts/parsers/cyberleninka.py`): раньше слал GET

View File

@@ -43,6 +43,7 @@ fi
REBUILD="" REBUILD=""
FRONTEND_CHANGED=0 FRONTEND_CHANGED=0
PARSERS_CHANGED=0
if [ "$FULL" = 1 ]; then if [ "$FULL" = 1 ]; then
REBUILD="$ALL_BACKEND" REBUILD="$ALL_BACKEND"
FRONTEND_CHANGED=1 FRONTEND_CHANGED=1
@@ -51,6 +52,12 @@ else
echo "$CHANGED" | grep -qE "^services/$svc/" && REBUILD="$REBUILD $svc" echo "$CHANGED" | grep -qE "^services/$svc/" && REBUILD="$REBUILD $svc"
done done
echo "$CHANGED" | grep -qE '^services/frontend/' && FRONTEND_CHANGED=1 echo "$CHANGED" | grep -qE '^services/frontend/' && FRONTEND_CHANGED=1
# Парсеры не вшиты в образ, а смонтированы в worker-indexer (см. compose),
# поэтому пересобирать нечего. Но процесс воркера держит модуль парсера в
# памяти с прошлого прогона: без перезапуска правка молча не применяется —
# заливка продолжает ходить по старому коду, и это не видно ниоткуда, кроме
# тела запросов в логах. Ловилось на CORE 16.09.2026.
echo "$CHANGED" | grep -qE '^scripts/parsers/' && PARSERS_CHANGED=1
fi fi
REBUILD="$(echo "$REBUILD" | xargs || true)" REBUILD="$(echo "$REBUILD" | xargs || true)"
@@ -111,6 +118,13 @@ echo "==> [4/6] Применение секретов к неизменивши
# shellcheck disable=SC2086 # shellcheck disable=SC2086
$COMPOSE up -d $ALL_BACKEND $COMPOSE up -d $ALL_BACKEND
# Правка парсера доезжает монтированием, но модуль уже загружен в память
# воркера — перезапускаем, если образ и так не пересобирался выше.
if [ "$PARSERS_CHANGED" = 1 ] && ! echo " $REBUILD " | grep -q " worker-indexer "; then
echo "==> [4.5/6] Парсеры менялись — перезапуск worker-indexer"
$COMPOSE restart worker-indexer
fi
echo "==> [5/6] Миграции БД (с ретраем)" echo "==> [5/6] Миграции БД (с ретраем)"
for i in $(seq 1 10); do for i in $(seq 1 10); do
$COMPOSE exec -T api alembic upgrade head && break $COMPOSE exec -T api alembic upgrade head && break

View File

@@ -0,0 +1,204 @@
#!/usr/bin/env python3
"""Health-монитор anti-plagiarism: письмо при СМЕНЕ статуса сервиса, без спама.
Запускается по cron на app-хосте (1.32) каждые 5 минут:
*/5 * * * * /usr/bin/python3 /home/user/antiplag_monitor.py >> ~/.antiplag-monitor/run.log 2>&1
Живёт в репозитории намеренно: прошлая версия лежала только на сервере, ничем
не версионировалась и месяц проверяла топологию, которой уже нет — вечный
ложный DOWN на остановленном CT 108 и полная слепота к реальному серверу
эмбеддингов. Постоянный ложный сигнал хуже отсутствия сигнала: на его фоне
теряется настоящая авария.
Что проверяется и почему именно это:
- API, PostgreSQL, Redis, RabbitMQ, MinIO, Elasticsearch — без них сервис не
принимает и не обрабатывает работы;
- Ollama на 1.40 — эмбеддинги для L3 (адрес брать из .env, а не хардкодить:
он уже дважды переезжал);
- воркеры Celery — самое важное дополнение: 05.09 брокер лежал час, воркеры
молчали, а прежний монитор рапортовал, что всё хорошо, потому что смотрел
только на TCP-порты.
OpenRouter (L4) намеренно не проверяется: это внешний сервис за VPN-прокси, его
недоступность не ломает основную проверку — L4 лишь уточняет вердикт.
"""
import json
import os
import smtplib
import socket
import ssl
import subprocess
import urllib.request
from datetime import datetime
from email.mime.text import MIMEText
APP_DIR = os.environ.get("APP_DIR", "/home/user/anti-plagiarism")
STATE_FILE = os.path.expanduser("~/.antiplag-monitor/state.json")
SMTP_HOST, SMTP_PORT = "mail.jze9mail.ru", 587
SMTP_USER = os.environ.get("MONITOR_SMTP_USER", "noreply")
SMTP_PASS = os.environ.get("MONITOR_SMTP_PASS", "kGy3vgt6EuiUBMNDmV01Ng")
SMTP_FROM = "noreply@jze9mail.ru"
ALERT_TO = os.environ.get("MONITOR_ALERT_TO", "jze9programer@gmail.com")
def env_value(key: str, default: str = "") -> str:
"""Значение из прод-.env (он генерируется из Infisical на каждом деплое).
Адреса инфраструктуры читаем оттуда, а не хардкодим: сервер эмбеддингов уже
переезжал дважды, и монитор каждый раз оставался со старым адресом.
"""
path = os.path.join(APP_DIR, ".env")
try:
with open(path, encoding="utf-8") as fh:
for line in fh:
line = line.strip()
if line.startswith(f"{key}="):
return line.split("=", 1)[1].strip().strip("'\"")
except Exception:
pass
return default
def host_port_from_url(url: str, default_port: int) -> tuple[str, int]:
"""Разобрать URL в пару (host, port), отбросив логин с паролем.
В REDIS_URL и RABBITMQ_URL креды идут перед адресом
(`redis://:пароль@host:port/0`), и без их отсечения разбор падает.
"""
rest = url.split("//", 1)[-1].split("/", 1)[0]
rest = rest.rpartition("@")[2] or rest
if ":" in rest:
host, _, port = rest.partition(":")
return host, int(port or default_port)
return rest, default_port
def check_tcp(host: str, port: int) -> bool:
try:
with socket.create_connection((host, port), timeout=5):
return True
except Exception:
return False
def check_http(url: str) -> bool:
try:
with urllib.request.urlopen(url, timeout=8) as r:
return r.status == 200
except Exception:
return False
def check_in_api_container(command: str) -> bool:
"""Проверка изнутри контейнера — так же, как это видит приложение."""
try:
out = subprocess.run(
["docker", "exec", "antiplagiator-api", "sh", "-c", command],
capture_output=True, text=True, timeout=20,
)
return out.stdout.strip() == "200"
except Exception:
return False
def check_celery_workers() -> bool:
"""Отвечают ли воркеры на ping.
Живой брокер ещё не значит работающую систему: воркер может не суметь к
нему подключиться и молча простаивать — именно так и было 05.09.
"""
try:
out = subprocess.run(
["docker", "exec", "antiplagiator-api", "python", "-c",
"from app.core.celery_app import celery_app;"
"print(len(celery_app.control.inspect(timeout=5).ping() or {}))"],
capture_output=True, text=True, timeout=30,
)
return int(out.stdout.strip() or 0) >= 4 # api ждёт 4 воркера
except Exception:
return False
def build_checks() -> list[tuple[str, callable]]:
"""Список проверок; адреса — из прод-.env, чтобы не разъезжаться с реальностью."""
pg_host = env_value("POSTGRES_HOST", "192.168.1.38")
redis_host, redis_port = host_port_from_url(
env_value("REDIS_URL", "redis://192.168.1.35:6379/0"), 6379)
minio_host, minio_port = host_port_from_url(
"http://" + env_value("MINIO_ENDPOINT", "192.168.1.21:9000"), 9000)
ollama_url = env_value("OLLAMA_URL", "http://192.168.1.40:11434")
rabbit_host, rabbit_port = host_port_from_url(
env_value("RABBITMQ_URL", "amqp://guest:guest@192.168.20.82:5672/"), 5672)
return [
("API", lambda: check_http("http://localhost:8000/health")),
("PostgreSQL", lambda: check_tcp(pg_host, 5432)),
("Redis", lambda: check_tcp(redis_host, redis_port)),
("RabbitMQ", lambda: check_tcp(rabbit_host, rabbit_port)),
("MinIO", lambda: check_tcp(minio_host, minio_port)),
("Elasticsearch", lambda: check_in_api_container(
"curl -s -o /dev/null -w '%{http_code}' http://elasticsearch:9200")),
("Embeddings", lambda: check_http(f"{ollama_url}/api/version")),
("CeleryWorkers", check_celery_workers),
]
def send_alert(subject: str, body: str) -> None:
msg = MIMEText(body, "plain", "utf-8")
msg["Subject"] = subject
msg["From"] = SMTP_FROM
msg["To"] = ALERT_TO
ctx = ssl.create_default_context()
with smtplib.SMTP(SMTP_HOST, SMTP_PORT, timeout=25) as s:
s.ehlo()
s.starttls(context=ctx)
s.ehlo()
s.login(SMTP_USER, SMTP_PASS)
s.send_message(msg)
def main() -> None:
results = {name: fn() for name, fn in build_checks()}
prev: dict = {}
if os.path.exists(STATE_FILE):
try:
with open(STATE_FILE) as fh:
prev = json.load(fh)
except Exception:
prev = {}
changes = []
for name, ok in results.items():
was = prev.get(name, True) # первый запуск считаем нормой
if was and not ok:
changes.append(f"УПАЛ: {name}")
elif not was and ok:
changes.append(f"восстановился: {name}")
if changes:
ts = datetime.now().strftime("%Y-%m-%d %H:%M")
status = "\n".join(f" {'OK ' if v else 'DOWN'} {k}" for k, v in results.items())
body = (f"Изменения статуса сервисов anti-plagiarism ({ts}):\n\n"
+ "\n".join(changes) + "\n\nТекущий статус:\n" + status)
subject = f"[anti-plagiarism] {changes[0]}"
if len(changes) > 1:
subject += f" (+{len(changes) - 1})"
try:
send_alert(subject, body)
except Exception as e:
print(f"не удалось отправить алерт: {e}")
os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True)
with open(STATE_FILE, "w") as fh:
json.dump(results, fh)
line = " ".join(f"{k}={'OK' if v else 'DOWN'}" for k, v in results.items())
print(datetime.now().strftime("%H:%M"), "|", line)
if __name__ == "__main__":
main()

View File

@@ -1,201 +0,0 @@
#!/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, поэтому
прерванная заливка продолжается с той же страницы бакета.
Проверено на проде 2026-08-31: 100 статей залито, 183 090 отпечатков (1831 на
статью — реальная глубина, а не аннотация). Темп — 0.3 статьи/с, то есть
миллион в один поток занял бы около 40 суток: скачивание тут не узкое место
(12 статей/с в 12 потоков), упирается в последовательную обработку одного
документа. Для миллионов обработку нужно распараллелить по ядрам воркера —
это следующий шаг, до него скрипт годится для порций в десятки тысяч.
"""
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()

View File

@@ -1,217 +0,0 @@
#!/usr/bin/env python3
"""Массовая заливка русской Википедии в корпус сравнения.
Зачем именно она: студенты копируют из Википедии чаще, чем из научных статей, а
в корпусе её нет вовсе. При этом она единственный русскоязычный источник такого
объёма, который отдаётся без блокировок — КиберЛенинка режет выкачку после ~130
запросов, eLIBRARY требует договора. Дамп `ruwiki-latest-pages-articles.xml.bz2`
(5.6 ГБ, ~2 млн статей) скачивается свободно.
Дамп не сохраняется на диск: читается потоком и распаковывается на лету —
на app-хосте всего 22 ГБ свободно.
Что отбрасывается: перенаправления, служебные пространства имён (оставляем
только статьи) и короткие тексты — для детекции нужен объём, а не заготовки.
Запуск (в контейнере worker-indexer):
cd /home/user/anti-plagiarism
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
$C --limit 100 < scripts/ops/bulk_ingest_wikipedia_ru.py # проба
$C --limit 500000 < scripts/ops/bulk_ingest_wikipedia_ru.py # объём
Идемпотентно: ext_id вида `wikipedia_ru:<pageid>`, повторная заливка пропускает
уже существующие статьи (ON CONFLICT DO NOTHING).
"""
import argparse
import bz2
import io
import re
import time
DUMP_URL = "https://dumps.wikimedia.org/ruwiki/latest/ruwiki-latest-pages-articles.xml.bz2"
PAGE_RE = re.compile(r"<page>(.*?)</page>", re.DOTALL)
TITLE_RE = re.compile(r"<title>(.*?)</title>", re.DOTALL)
ID_RE = re.compile(r"<id>(\d+)</id>")
NS_RE = re.compile(r"<ns>(\d+)</ns>")
TEXT_RE = re.compile(r'<text[^>]*>(.*?)</text>', re.DOTALL)
REDIRECT_RE = re.compile(r"<redirect ")
# Разметку чистим регулярками: для отпечатков нужен связный текст, а не точное
# восстановление вёрстки — полноценный парсер вики-разметки тут неоправдан
CLEAN_RULES = [
(re.compile(r"\{\{[^{}]*\}\}"), " "), # шаблоны (в несколько проходов)
(re.compile(r"\[\[[^\]|]*\|"), ""), # [[ссылка|текст → текст
(re.compile(r"\[\[|\]\]"), ""), # остатки скобок
(re.compile(r"<ref[^>]*>.*?</ref>", re.DOTALL), " "),
(re.compile(r"<[^>]+>"), " "), # html-теги
(re.compile(r"^[*#:;|!].*$", re.MULTILINE), " "), # списки и таблицы
(re.compile(r"^=+.*?=+$", re.MULTILINE), " "), # == заголовки разделов ==
(re.compile(r"'{2,}"), ""), # ''курсив''
(re.compile(r"&[a-z]+;"), " "),
(re.compile(r"[ \t]+"), " "),
(re.compile(r"\n{2,}"), "\n"),
]
def clean_wikitext(raw: str) -> str:
text = raw
for _ in range(3): # шаблоны бывают вложенными
text = CLEAN_RULES[0][0].sub(CLEAN_RULES[0][1], text)
for rx, repl in CLEAN_RULES[1:]:
text = rx.sub(repl, text)
return text.strip()
def _pages_from_chunks(chunks, min_chars: int):
"""Разобрать поток распакованного XML на статьи основного пространства имён."""
buf = ""
for raw in chunks:
buf += raw.decode("utf-8", errors="replace")
while True:
m = PAGE_RE.search(buf)
if not m:
break
page, buf = m.group(1), buf[m.end():]
if REDIRECT_RE.search(page):
continue
ns = NS_RE.search(page)
if not ns or ns.group(1) != "0": # только статьи
continue
tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page)
if not (tm and im and xm):
continue
text = clean_wikitext(xm.group(1))
if len(text) < min_chars:
continue
yield im.group(1), tm.group(1), text
if len(buf) > 20 * 1024 * 1024: # страховка от разбухания
buf = buf[-1024 * 1024:]
def iter_pages(min_chars: int, dump_file: str | None):
"""Статьи из локального дампа либо, если файла нет, прямо из сети.
Локальный файл предпочтителен: Wikimedia рвёт долгие потоковые соединения
(проверено — обрыв на 32 МБ из 5.9 ГБ), а скачать файл можно с докачкой
(`curl -C -`) и потом читать сколько угодно.
"""
decomp = bz2.BZ2Decompressor()
if dump_file:
def chunks():
with open(dump_file, "rb") as fh:
while True:
part = fh.read(4 * 1024 * 1024)
if not part:
return
try:
yield decomp.decompress(part)
except EOFError:
return
yield from _pages_from_chunks(chunks(), min_chars)
return
import httpx
# Wikimedia отдаёт 403 без осмысленного User-Agent — по их правилам он должен
# называть приложение и давать контакт
headers = {"User-Agent": "AcademicHelper/1.0 (https://academic.jze9.ru; noreply@jze9.ru)"}
def net_chunks():
with httpx.stream("GET", DUMP_URL, timeout=120, follow_redirects=True,
headers=headers) as resp:
resp.raise_for_status()
for chunk in resp.iter_bytes(4 * 1024 * 1024):
try:
yield decomp.decompress(chunk)
except EOFError:
return
yield from _pages_from_chunks(net_chunks(), min_chars)
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--limit", type=int, default=100, help="сколько статей залить")
ap.add_argument("--batch", type=int, default=500, help="статей в транзакции")
ap.add_argument("--min-chars", type=int, default=2000, help="минимальная длина текста")
ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на статью")
ap.add_argument("--dump-file", default=None,
help="путь к заранее скачанному дампу (надёжнее, чем читать из сети)")
args = ap.parse_args()
import psycopg2
from app.algorithms.winnowing import winnow
from app.config import settings
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 = 0
fp_total = 0
t0 = time.time()
batch: list[tuple[str, str, str]] = []
def flush(batch):
nonlocal added, skipped, fp_total
if not batch:
return
with conn.cursor() as cur:
cur.executemany(
"""INSERT INTO documents (ext_id, source, title, lang, abstract, url,
authors, indexed_at)
VALUES (%s, 'wikipedia_ru', %s, 'ru', %s, %s, '[]'::json, now())
ON CONFLICT (ext_id) DO NOTHING""",
[(f"wikipedia_ru:{pid}", title[:1000], text[:2000],
f"https://ru.wikipedia.org/?curid={pid}") for pid, title, text in batch],
)
cur.execute("SELECT id, ext_id FROM documents WHERE ext_id = ANY(%s)",
([f"wikipedia_ru:{pid}" for pid, _, _ in batch],))
id_by_ext = {e: i for i, e in cur.fetchall()}
buf = io.StringIO()
fresh = 0
for pid, _, text in batch:
doc_id = id_by_ext.get(f"wikipedia_ru:{pid}")
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(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()
for pid, title, text in iter_pages(args.min_chars, args.dump_file):
batch.append((pid, title, text))
if len(batch) >= args.batch:
flush(batch)
batch = []
el = time.time() - t0
print(f" залито {added} · пропущено {skipped} · отпечатков {fp_total:,} · "
f"{added / max(el, 1):.1f} ст/с", flush=True)
if added >= args.limit:
break
flush(batch)
conn.close()
el = time.time() - t0
print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, "
f"отпечатков {fp_total:,}")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,180 @@
#!/usr/bin/env python3
"""Замер качества детекции на реальном корпусе: что находим и что выдумываем.
Зачем: до сих пор о качестве проверки мы знали только «механизм жив» — находит
подброшенный фрагмент. Этого мало. Продукт продаёт процент заимствований, и
нужно знать две вещи: сколько списываний система пропускает (полнота) и как
часто обвиняет невиновных (точность). Без этих чисел непонятно, что улучшать —
пороги, глубину корпуса или алгоритм.
Как устроен замер. Берём документы из самого корпуса и делаем из них
«студенческие работы» четырёх видов:
дословно — фрагмент скопирован как есть (обязан находиться);
лёгкий рерайт — выброшено каждое 10-е слово (обязан находиться);
сильный рерайт— выброшено каждое 3-е слово (L1 не обязан; это работа L3/L4);
оригинал — текст, которого в корпусе нет (обязан НЕ находиться).
Считается путь L1 — тот же, что в index.extract_and_check: фрагментация,
winnowing по каждому фрагменту, поиск по отпечаткам с порогом
EXACT_FRAGMENT_THRESHOLD. Уровни L3/L4 сюда не входят намеренно: они требуют
GPU-воркера и на порядок медленнее, а мерить нужно в первую очередь основу.
Запуск (в контейнере worker-indexer):
cd /home/user/anti-plagiarism
docker compose -f docker-compose.prod.yml exec -T worker-indexer python - \\
< scripts/ops/detection_benchmark.py
... python - --docs 40 < scripts/ops/detection_benchmark.py
Ничего не пишет в базу — только читает.
"""
import argparse
import random
def make_cases(text: str, words_taken: int) -> dict[str, str]:
"""Сделать из текста источника четыре «студенческие работы»."""
words = text.split()
start = len(words) // 3
chunk = words[start : start + words_taken]
return {
"дословно": " ".join(chunk),
"лёгкий рерайт": " ".join(w for i, w in enumerate(chunk) if i % 10),
"сильный рерайт": " ".join(w for i, w in enumerate(chunk) if i % 3),
}
def original_text(rnd: random.Random, words_taken: int) -> str:
"""Текст, которого заведомо нет в корпусе — проверка на ложные обвинения."""
topics = [
"вчера вечером я варил гречку с грибами и вспоминал поездку на озеро",
"мой сосед купил подержанный велосипед и красит его в оранжевый цвет",
"бабушка печёт пироги с капустой по субботам и зовёт всех соседей",
"кот запрыгнул на подоконник и уронил горшок с геранью на пол",
"в субботу мы чинили забор и обсуждали цены на доски и гвозди",
]
out: list[str] = []
while len(out) < words_taken:
out.extend(rnd.choice(topics).split())
return " ".join(out[:words_taken])
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--docs", type=int, default=25, help="сколько документов взять для замера")
ap.add_argument("--words", type=int, default=400, help="сколько слов «списывать»")
ap.add_argument("--seed", type=int, default=20260905, help="зерно случайности (для повторяемости)")
args = ap.parse_args()
from app.algorithms.winnowing import winnow
from app.config import settings
from app.db import db_session, get_minio
from app.fragments import split_into_fragments
from app.models import Document, Fingerprint
from sqlalchemy import func, select
from sqlalchemy import text as sql_text
rnd = random.Random(args.seed)
# Берём документы с реальной глубиной: по аннотации в 30 отпечатков мерить
# нечего, и такой замер сказал бы больше о корпусе, чем об алгоритме
with db_session() as s:
rows = s.execute(sql_text("""
SELECT d.id, d.source, d.minio_key, d.abstract
FROM documents d
JOIN (SELECT doc_id, count(*) c FROM fingerprints
GROUP BY doc_id HAVING count(*) > 800) f ON f.doc_id = d.id
WHERE d.source <> 'user_submission'
ORDER BY random() LIMIT :n
"""), {"n": args.docs}).all()
print(f"документов в замере: {len(rows)}, «списываем» по {args.words} слов")
print(f"порог срабатывания фрагмента: {settings.EXACT_FRAGMENT_THRESHOLD}%\n")
def find_source(work_text: str) -> set[int]:
"""Прогнать «работу» через L1 и вернуть найденные документы-источники."""
found: set[int] = set()
frags = split_into_fragments(
work_text,
window=settings.FRAGMENT_WINDOW_WORDS,
overlap=settings.FRAGMENT_OVERLAP_WORDS,
)
with db_session() as s:
for frag in frags:
fp = winnow(frag["text"])
if not fp:
continue
hit = s.execute(
select(Fingerprint.doc_id, func.count(Fingerprint.id).label("c"))
.join(Document, Document.id == Fingerprint.doc_id)
.where(
Fingerprint.hash_value.in_(list(fp)),
Document.source != "user_submission",
)
.group_by(Fingerprint.doc_id)
.order_by(func.count(Fingerprint.id).desc())
.limit(1)
).first()
if hit and hit[1] / len(fp) * 100 >= settings.EXACT_FRAGMENT_THRESHOLD:
found.add(hit[0])
return found
stats: dict[str, dict[str, int]] = {}
minio = get_minio()
for doc_id, _source, minio_key, abstract in rows:
# Полный текст, если он сохранён; иначе то, что есть в базе
body = abstract or ""
if minio_key:
try:
obj = minio.get_object(settings.MINIO_BUCKET_DOCS, minio_key)
body = obj.read().decode("utf-8", errors="replace")
obj.close()
obj.release_conn()
except Exception:
pass
if len(body.split()) < args.words * 2:
continue
for case_name, work in make_cases(body, args.words).items():
found = find_source(work)
bucket = stats.setdefault(case_name, {"всего": 0, "нашли верно": 0, "не нашли": 0})
bucket["всего"] += 1
if doc_id in found:
bucket["нашли верно"] += 1
else:
bucket["не нашли"] += 1
# Контроль на ложные обвинения
found = find_source(original_text(rnd, args.words))
bucket = stats.setdefault("оригинал", {"всего": 0, "ложно обвинён": 0, "чисто": 0})
bucket["всего"] += 1
if found:
bucket["ложно обвинён"] += 1
else:
bucket["чисто"] += 1
print(f"{'случай':16} {'всего':>6} {'найден':>8} {'пропущен':>10} доля найденных")
for case in ("дословно", "лёгкий рерайт", "сильный рерайт"):
b = stats.get(case)
if not b:
continue
share = b["нашли верно"] / b["всего"] * 100 if b["всего"] else 0
print(f"{case:16} {b['всего']:>6} {b['нашли верно']:>8} {b['не нашли']:>10} {share:5.1f}%")
b = stats.get("оригинал")
if b:
fp_rate = b["ложно обвинён"] / b["всего"] * 100 if b["всего"] else 0
print(f"\nоригинальные тексты: {b['всего']}, ложных обвинений {b['ложно обвинён']} "
f"({fp_rate:.1f}%)")
print("\nЧитать так: «дословно» и «лёгкий рерайт» — обязанность L1, доля должна быть "
"близка к 100%. «Сильный рерайт» L1 брать не обязан, это работа L3/L4. "
"Ложные обвинения должны быть нулевыми.")
if __name__ == "__main__":
main()

160
scripts/ops/keep_ingesting.py Executable file
View File

@@ -0,0 +1,160 @@
#!/usr/bin/env python3
"""Держать массовую заливку идущей: перезапускать источники, когда они закончили.
Массовый источник за один прогон берёт столько статей, сколько влезает в бюджет
времени таска, сохраняет позицию и останавливается. Без внешнего толчка заливка
идёт рывками — ровно столько, сколько раз кто-то нажмёт «Запустить». Этот скрипт
и есть толчок: раз в несколько минут проверяет, не простаивает ли источник, и
запускает следующую порцию.
Ставится в cron на app-хосте:
*/5 * * * * cd /home/user/anti-plagiarism && /usr/bin/docker compose \
-f docker-compose.prod.yml exec -T worker-indexer python - \
< scripts/ops/keep_ingesting.py >> ~/.antiplag-monitor/ingest.log 2>&1
Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи,
прогон закрывается с нулём добавленных, и скрипт больше его не трогает
(--stop-when-empty, по умолчанию включено).
Перед подсчётом «занято ли» снимает зависшие прогоны: если worker-indexer упал
или перезапустился, прогон навсегда остаётся в статусе running с замолчавшим
heartbeat — без этой чистки сторож видит «занято» бесконечно и не запускает
вообще ничего, для всех источников сразу (порог тот же, что и в панели
отладки — services/api/app/api/admin.py, stale_cutoff).
Молчащий heartbeat сам по себе ещё не значит «мёртв»: пачка COPY на большой
базе идёт десятки минут без единого тика. Поэтому прогон снимается, только
если его таска нет среди выполняющихся на воркерах queue.index; не ответил
хоть один такой воркер — не снимается ничего. Иначе сторож снимал живой
медленный прогон и тут же запускал его дубль по тем же статьям.
"""
import argparse
BULK_TYPES = ("wikipedia_ru", "pmc_bulk", "core")
STALE_MINUTES = 10
def live_task_ids(celery_app) -> set[str] | None:
"""ID тасков, выполняющихся на воркерах queue.index; None — кто-то из них не ответил."""
queues = celery_app.control.inspect(timeout=10).active_queues() or {}
workers = [w for w, qs in queues.items() if any(q["name"] == "queue.index" for q in qs)]
if not workers:
return None
active = celery_app.control.inspect(destination=workers, timeout=10).active()
if not active or set(active) != set(workers):
return None
return {t["id"] for tasks in active.values() for t in tasks}
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--types", default=",".join(BULK_TYPES),
help="типы источников через запятую")
ap.add_argument("--keep-going-empty", action="store_true",
help="запускать даже если прошлый прогон ничего не добавил")
args = ap.parse_args()
from app.celery_app import celery_app
from app.db import db_session
from app.models import ParseRun, ParseSource
from sqlalchemy import text
types = tuple(t.strip() for t in args.types.split(",") if t.strip())
live = live_task_ids(celery_app)
with db_session() as s:
if live is None:
print("воркеры queue.index не ответили — живость прогонов не проверить, не снимаю")
reaped = []
else:
reaped = s.execute(text(f"""
UPDATE parse_runs SET status='error', stage='finished',
error='зависший прогон снят автоматически: таск не выполняется, heartbeat молчал дольше {STALE_MINUTES} минут',
finished_at=now()
WHERE status='running'
AND (heartbeat_at IS NULL OR heartbeat_at < now() - interval '{STALE_MINUTES} minutes')
AND (celery_task_id IS NULL OR NOT (celery_task_id = ANY(CAST(:live AS text[]))))
RETURNING id, source_id
"""), {"live": sorted(live)}).all()
# queued без heartbeat — тот же зомби другого вида: send_task не дошёл
# или упал между двумя commit при постановке в очередь (celery_task_id
# так и остался пустым), воркер о такой задаче никогда не узнает
reaped += s.execute(text(f"""
UPDATE parse_runs SET status='error', stage='finished',
error='зависший прогон снят автоматически: висел в очереди дольше {STALE_MINUTES} минут',
finished_at=now()
WHERE status='queued'
AND celery_task_id IS NULL
AND started_at < now() - interval '{STALE_MINUTES} minutes'
RETURNING id, source_id
""")).all()
if reaped:
for source_id in {row.source_id for row in reaped}:
s.execute(text("""
UPDATE parse_sources SET last_status='error',
last_error='зависший прогон снят автоматически'
WHERE id=:sid AND last_status='running'
"""), {"sid": source_id})
s.commit()
print(f"снято зависших прогонов: {len(reaped)} ({', '.join(str(r.id) for r in reaped)})")
active = s.execute(text(
"SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')"
)).scalar_one()
if active:
print(f"уже идёт прогонов: {active} — ждём")
return
sources = s.execute(text("""
SELECT id, source_type, last_run_id FROM parse_sources
WHERE enabled AND source_type = ANY(:types) ORDER BY id
"""), {"types": list(types)}).all()
started = []
for source_id, source_type, last_run_id in sources:
# Источник исчерпан, если прогоны перестали добавлять статьи: дальше
# в дампе/бакете/выдаче для нас ничего нет.
#
# Но ОДИН пустой прогон — ещё не приговор: ровно так же выглядят
# обрыв связи с API и прогон на устаревшем коде парсера. На CORE
# 16.09.2026 из-за этого источник был помечен исчерпанным после
# единственной неудачи и больше не запускался — молча, без ошибки.
# Поэтому ждём два пустых прогона подряд: разовый сбой переживём,
# а реально кончившийся источник остановится всего на один прогон
# позже.
if last_run_id and not args.keep_going_empty:
last_two = s.execute(text("""
SELECT status, added, fetched FROM parse_runs
WHERE source_id = :sid ORDER BY id DESC LIMIT 2
"""), {"sid": source_id}).all()
if len(last_two) == 2 and all(
r.status == "done" and r.added == 0 and r.fetched == 0 for r in last_two
):
print(f"{source_type}: источник исчерпан (два пустых прогона), пропускаем")
continue
src = s.get(ParseSource, source_id)
run = ParseRun(source_id=source_id, status="queued", stage="queued",
target=src.limit, log=[])
s.add(run)
s.flush()
src.last_status = "running"
src.last_error = None
src.last_run_id = run.id
s.commit()
result = celery_app.send_task("index.run_parser", args=[source_id, run.id],
queue="queue.index")
run.celery_task_id = result.id
s.commit()
started.append(f"{source_type}(прогон {run.id})")
print("запущено:", ", ".join(started) if started else "нечего запускать")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,79 @@
#!/usr/bin/env python3
"""Вычислить байтовое смещение в дампе Википедии по уже залитым статьям.
Нужен один раз при переходе на multistream-дамп: позиция продолжения сменила
смысл — раньше это был номер статьи (и прогон перечитывал дамп с начала),
теперь байтовое смещение блока.
Как считается: берём максимальный page_id среди залитых статей и находим в
индексе дампа блок, которому он принадлежит. Страницы в дампе идут по
возрастанию page_id, поэтому всё, что дальше этого блока, ещё не залито.
Запуск (в контейнере worker-indexer):
docker compose -f docker-compose.prod.yml exec -T worker-indexer python - \\
< scripts/ops/wikipedia_resume_offset.py # показать
... python - --apply < scripts/ops/wikipedia_resume_offset.py # записать
"""
import argparse
import bz2
import os
DEFAULT_INDEX = "/parsers/ruwiki-index.txt.bz2"
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--apply", action="store_true", help="записать смещение в источник")
ap.add_argument("--index", default=DEFAULT_INDEX, help="путь к индексу дампа")
args = ap.parse_args()
from app.db import db_session
from sqlalchemy import text
if not os.path.exists(args.index):
raise SystemExit(f"индекс не найден: {args.index}")
with db_session() as s:
row = s.execute(text("""
SELECT max(split_part(ext_id, ':', 2)::bigint)
FROM documents WHERE source = 'wikipedia_ru'
""")).first()
max_pageid = int(row[0]) if row and row[0] else 0
print(f"максимальный page_id среди залитых: {max_pageid}")
if not max_pageid:
print("залитых статей нет — начинать с нуля")
return
# Индекс отсортирован по смещению, page_id внутри растут: ищем блок,
# в котором лежит наш максимальный, и берём его смещение
offset = 0
found_title = ""
with bz2.open(args.index, "rt", encoding="utf-8") as fh:
for line in fh:
parts = line.split(":", 2)
if len(parts) < 3:
continue
off, pid = int(parts[0]), int(parts[1])
if pid > max_pageid:
break
offset, found_title = off, parts[2].strip()
print(f"блок с этой статьёй начинается на байте {offset} ({found_title[:50]!r})")
if not args.apply:
print("\nзапустите с --apply, чтобы записать смещение в источник")
return
with db_session() as s:
s.execute(text("""
UPDATE parse_sources SET resume_token = :off WHERE source_type = 'wikipedia_ru'
"""), {"off": str(offset)})
s.commit()
print(f"смещение {offset} записано в источник")
if __name__ == "__main__":
main()

282
scripts/parsers/core.py Normal file
View File

@@ -0,0 +1,282 @@
"""Парсер CORE (core.ac.uk) — полные тексты открытых репозиториев.
Зачем нужен: полный текст есть у малой доли корпуса, и по-русски дыра самая
большая — CyberLeninka отдаёт только HTML-страницу статьи, тела не даёт. CORE
агрегирует открытые репозитории (в том числе вузовские русско- и
белорусскоязычные) и кладёт готовый текст прямо в выдачу поиска — его не надо
ни скачивать отдельно, ни извлекать из PDF.
Замер на живом API 16.09.2026: страница в 100 записей приходит за ~7с и
содержит 47 статей с текстом длиннее 1500 символов, медиана — 12.7 тыс.
символов. То есть один запрос ≈ 47 документов корпуса.
Чем сузить выборку до русских текстов: в общем потоке CORE кириллицы нет
вовсе (проверено: 0 из 48 полных текстов за 2019). Язык у работ размечен
плохо — `language` сплошь None или "zz", фильтровать по нему нельзя. Зато
работает отбор по архиву-поставщику: запрос вида
`(repositories.id:1298 OR repositories.id:21908 OR ...)` по вузовским
репозиториям России и Беларуси даёт 79 полных текстов из 100 записей, 78 из
них — кириллица. Сам список архивов живёт в поле `query` источника, а не в
коде: его правят из админки, не пересобирая образ.
Массовый парсер (`bulk = True`): отдаёт статьи генератором и пишется пачками
через COPY, как Википедия и PMC (см. bulk_writer.py). Статьи без текста
пропускаются молча — они уже есть у нас как метаданные из OpenAlex, а
обогащать корпус нечем.
Грабли API, все проверены живьём:
- **с российских адресов ключ не работает**: прямой запрос с прода виснет без
ошибки (соединение есть, тело ответа не приходит), через `singbox-proxy` тот
же запрос отвечает за 1.7с. Анонимные короткие запросы проходят и напрямую,
из-за чего поломка выглядит как сетевая. Адрес прокси — `CORE_PROXY_URL`;
- у `/v3/search/works` обязателен слэш на конце: иначе 301, а при редиректе
теряется заголовок Authorization и запрос уходит анонимным;
- `offset` упирается в 100 000 (под капотом Azure Search, глубже — 400),
поэтому выборка режется на части — по одному архиву-поставщику на запрос;
- **условия через `AND` не связываются**: `(repositories.id:...) AND yearPublished:2026`
отдаёт 2.3 млн работ вперемешку по годам, то есть условия объединяются по
«или» и год работает лишь как подсказка ранжированию. По отдельности каждое
условие фильтрует честно (`yearPublished:2026` — ровно 2026, `repositories.id:1298`
— ровно этот архив). Поэтому в запрос идёт РОВНО ОДНО условие: номер архива.
Годы, если заданы, отсекаются уже на нашей стороне;
- `fullText` — не фильтруемое поле, `_exists_:fullText` отвечает 500.
Отбирать статьи с текстом приходится на своей стороне, оплачивая трафиком;
- CORE агрегирует репозитории и отдаёт одну статью несколько раз под разными
id (проверено: 90810208 и 354263199 — один текст). Дедупликация корпуса
идёт только по ext_id, поэтому такие пары легли бы отдельными документами
и потом ловились бы как заимствование друг у друга. Отсеиваем по хешу
текста в пределах прогона: копии приходят в соседних строках выдачи.
"""
import hashlib
import logging
import os
import re
import time
from collections.abc import Iterator
from typing import Any
import httpx
from base import BaseParser, ProgressCallback
logger = logging.getLogger(__name__)
API_URL = "https://api.core.ac.uk/v3/search/works/" # слэш обязателен, см. докстринг
PAGE = 100 # максимум записей за запрос
MAX_OFFSET = 100_000 # потолок глубины у Azure Search
MIN_CHARS = 1500 # короче — обрывок или аннотация, а не статья
RETRIES = 3
# CORE отдаёт данные по ключу только за пределами РФ: с прода прямой запрос
# молча виснет (TCP есть, тело ответа не приходит), а тот же запрос через
# sing-box возвращается за 1.7с. Анонимные мелкие запросы проходят и напрямую,
# поэтому со стороны это выглядит как поломка сети, а не блокировка. Тот же
# приём, что для OpenRouter в worker-gpu. Пусто = ходить напрямую.
DEFAULT_PROXY = "socks5://singbox-proxy:1080"
class COREParser(BaseParser):
"""Полные тексты из CORE. Отдаёт статьи потоком, режет выборку по годам."""
source_name = "core"
bulk = True
def __init__(self, api_key: str | None = None, proxy_url: str | None = None) -> None:
super().__init__()
self.api_key = api_key or os.environ.get("CORE_API_KEY", "")
# Вне compose-сети (локальный прогон, тесты) имени singbox-proxy нет —
# там прокси отключают, выставив CORE_PROXY_URL пустым
self.proxy_url = self._proxy(proxy_url)
if not self.api_key:
logger.warning("CORE: ключ не задан (CORE_API_KEY) — API ответит 401")
# Страница весит ~5 МБ, обычный ответ ~7с. 30с — запас на джиттер в
# десять раз, но без риска съесть бюджет прогона: три попытки по 60с
# отъедали 180с из 1500с на одной залипшей странице (поймано 16.09.2026,
# та же грабля, что чинили в pmc_bulk)
self.client = httpx.Client(
timeout=30,
proxy=self.proxy_url or None,
headers={
"Authorization": f"Bearer {self.api_key}",
"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)",
},
)
self.last_token: str | None = None
@staticmethod
def _proxy(proxy_url: str | None) -> str:
"""Адрес прокси: явный аргумент, иначе CORE_PROXY_URL, иначе умолчание."""
if proxy_url is not None:
return proxy_url
return os.environ.get("CORE_PROXY_URL", DEFAULT_PROXY)
def fetch( # type: ignore[override]
self,
limit: int = 1000,
query: str = "",
year_from: int | None = None,
year_to: int | None = None,
resume_token: str | None = None,
progress_cb: ProgressCallback | None = None,
**_ignored: Any,
) -> Iterator[dict[str, Any]]:
"""Статьи с полным текстом, архив за архивом.
Args:
limit: сколько статей с текстом отдать за прогон
query: список архивов — номера через запятую/пробел или выражение
`(repositories.id:N OR ...)`. Пустой — обычный поиск по словам
year_from: отсекать статьи старше этого года (необязательно)
year_to: отсекать статьи новее этого года (необязательно)
resume_token: позиция вида "1298:4200" — архив и смещение в нём
progress_cb: см. base.ProgressCallback
"""
parts = self._repos(query)
start_part, start_offset = self._parse_token(resume_token)
if start_part in parts:
parts = parts[parts.index(start_part):]
else:
start_offset = 0
given = 0
seen: set[str] = set()
for part in parts:
offset = start_offset
start_offset = 0 # смещение относится только к архиву из токена
while given < limit and offset < MAX_OFFSET:
results = self._page(part, query, offset)
if results is None: # API не ответил — прогон закончен
return
if not results:
logger.info("CORE: архив %s исчерпан на смещении %d", part, offset)
break
token = f"{part}:{offset}"
for work in results:
text = work.get("fullText") or ""
if len(text) < MIN_CHARS:
continue
year = work.get("yearPublished")
if year and ((year_from and year < year_from)
or (year_to and year > year_to)):
continue
# Дубли CORE: один текст под разными id (см. докстринг).
# Хеша начала текста хватает — совпадение первых 5 тыс.
# символов у разных статей практически исключено
digest = hashlib.sha1(text[:5000].encode()).hexdigest()
if digest in seen:
continue
seen.add(digest)
given += 1
# Токен указывает на страницу, из которой пришла статья, а
# не на следующую: пачка может прерваться на её середине, и
# тогда следующий прогон перечитает страницу целиком.
# Лишний запрос дешевле потерянных статей, дубли отсекает
# ON CONFLICT на ext_id
work["resume_token"] = token
yield work
if given >= limit:
break
offset += PAGE
self.last_token = f"{part}:{offset}"
if progress_cb and not progress_cb(given):
logger.info("CORE: выборка остановлена по запросу (%d)", given)
return
logger.info("CORE: пройдено архивов: %d, отдано %d", len(parts), given)
@staticmethod
def _repos(query: str) -> list[str]:
"""Номера архивов из настройки источника.
Принимает и выражение `(repositories.id:1298 OR ...)`, и простой список
«1298, 21908». Если номеров нет вовсе — единственная часть с именем
`q`: тогда парсер просто ищет по словам запроса.
"""
ids = re.findall(r"repositories\.id:(\d+)", query or "")
if not ids:
ids = re.findall(r"\b\d{2,7}\b", query or "")
return ids or ["q"]
@staticmethod
def _parse_token(token: str | None) -> tuple[str | None, int]:
"""Разобрать позицию "архив:смещение"; мусор — начать сначала."""
if not token or ":" not in token:
return None, 0
part, _, offset = token.rpartition(":")
try:
return part, int(offset)
except ValueError:
logger.warning("CORE: непонятная позиция %r, начинаем сначала", token)
return None, 0
def _page(self, part: str, query: str, offset: int) -> list[dict[str, Any]] | None:
"""Одна страница выдачи; None — API не отвечает, прогон пора кончать.
В запросе РОВНО одно условие: связка через AND в этом API не работает
(см. докстринг модуля), поэтому годы отсекаются уже после выборки.
"""
q = f"repositories.id:{part}" if part != "q" else (query or "*")
params = {"q": q, "limit": PAGE, "offset": offset}
for attempt in range(RETRIES):
try:
resp = self.client.get(API_URL, params=params)
if resp.status_code == 429:
# Лимит тарифа: подождать и повторить, а не ронять прогон
wait = int(resp.headers.get("retry-after", 30))
logger.warning("CORE: лимит запросов, ждём %dс", wait)
time.sleep(min(wait, 60))
continue
resp.raise_for_status()
return resp.json().get("results") or []
except Exception as e:
logger.warning(
"CORE: страница %s:%d не удалась (%d/%d): %s",
part, offset, attempt + 1, RETRIES, e,
)
time.sleep(2 * (attempt + 1))
logger.error("CORE: страница %s:%d не далась за %d попыток", part, offset, RETRIES)
return None
def transform(self, raw: dict[str, Any]) -> dict[str, Any]:
"""Привести статью CORE к унифицированному формату корпуса."""
ext = raw.get("id")
text = raw.get("fullText") or ""
if not ext or len(text) < MIN_CHARS:
return {}
lang = raw.get("language") or {}
code = lang.get("code") if isinstance(lang, dict) else lang
journals = raw.get("journals") or []
journal = journals[0].get("title") if journals and isinstance(journals[0], dict) else None
return {
"source": self.source_name,
"ext_id": f"core:{ext}",
"title": (raw.get("title") or f"CORE {ext}")[:1000],
"authors": self.normalize_authors(self._authors(raw)),
"doi": raw.get("doi"),
"year": raw.get("yearPublished"),
# zz у CORE значит «язык не определён» — лучше пусто, чем мусор
"lang": code if code and code != "zz" else None,
"journal": journal[:300] if journal else None,
"url": raw.get("downloadUrl") or f"https://core.ac.uk/works/{ext}",
"abstract": (raw.get("abstract") or text[:2000])[:2000],
"text": text,
"resume_token": raw.get("resume_token"),
}
@staticmethod
def _authors(raw: dict[str, Any]) -> list[dict[str, str]]:
"""Авторы CORE приходят одной строкой "Фамилия, Имя Отчество"."""
out = []
for author in raw.get("authors") or []:
name = author.get("name") if isinstance(author, dict) else author
if not name:
continue
last, _, first = str(name).partition(",")
out.append({"last_name": last.strip(), "first_name": first.strip()})
return out

164
scripts/parsers/pmc_bulk.py Normal file
View File

@@ -0,0 +1,164 @@
"""Парсер PubMed Central из открытого бакета AWS — массовый путь.
Отличие от `pmc.py`: тот ходит в API E-utilities и годится для точечной
докачки по теме (1-2 статьи в секунду, есть rate limit). Этот берёт статьи из
бакета `pmc-oa-opendata`, где у каждой лежит **уже извлечённый текст**, и
качает их пачками в несколько потоков — путь к сотням тысяч и миллионам.
Ключ для доступа не нужен, бакет открыт. Замер: 12 статей/с в 10 потоков.
Как и парсер Википедии, отдаёт статьи генератором: складывать миллионы
документов в список нельзя (см. bulk_writer.py).
"""
import logging
import re
from collections.abc import Iterator
from concurrent.futures import ThreadPoolExecutor, as_completed
from typing import Any
from urllib.parse import quote
import httpx
from base import BaseParser, ProgressCallback
logger = logging.getLogger(__name__)
BUCKET = "https://pmc-oa-opendata.s3.amazonaws.com"
KEY_RE = re.compile(r"<Key>(PMC\d+\.\d+)/\1\.txt</Key>")
TOKEN_RE = re.compile(r"<NextContinuationToken>([^<]+)</NextContinuationToken>")
YEAR_RE = re.compile(r"\b(19|20)\d{2}\b")
MIN_CHARS = 1500 # меньше — обрывок, а не статья
class PMCBulkParser(BaseParser):
"""PMC Open Access из бакета AWS. Отдаёт статьи потоком."""
source_name = "pmc" # тот же источник в корпусе, что и у API-парсера
bulk = True
def __init__(self, workers: int = 12) -> None:
super().__init__()
self.workers = workers
# 60с на файл при бюджете задачи в 1500с — один залипший запрос съедал
# почти весь бюджет. Обычный GET сюда укладывается в доли секунды,
# 12с — с большим запасом на джиттер, но без риска съесть весь прогон
self.client = httpx.Client(
timeout=12,
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"},
)
# Позиция листинга: бакет отдаётся страницами, и продолжать прогон
# нужно с той же страницы, а не с начала
self.last_token: str | None = None
def fetch( # type: ignore[override]
self,
limit: int = 1000,
resume_token: str | None = None,
start_after: str | None = None,
progress_cb: ProgressCallback | None = None,
**_ignored: Any,
) -> Iterator[dict[str, Any]]:
"""Статьи из бакета: листинг страницами, скачивание в потоках.
Args:
limit: сколько статей отдать
resume_token: продолжить листинг с сохранённой страницы бакета
start_after: ключ, после которого начинать листинг. Нужен, когда
токена нет (первый прогон после ручных заливок): без него
листинг пошёл бы с начала бакета и часами перемалывал уже
залитые статьи как дубли
progress_cb: см. base.ProgressCallback
"""
token = resume_token
given = 0
while given < limit:
try:
ids, token = self._list_page(
token, want=min(200, limit - given),
start_after=None if token else start_after,
)
except Exception as e:
logger.error("PMC bulk: листинг не удался: %s", e)
return
if not ids:
logger.info("PMC bulk: бакет закончился")
return
# as_completed вместо map(): map() отдаёт результаты строго по
# порядку отправки, поэтому один залипший запрос блокирует все
# уже готовые — даже если остальные 11 потоков давно отработали
with ThreadPoolExecutor(max_workers=self.workers) as pool:
futures = {pool.submit(self._fetch_article, i): i for i in ids}
for future in as_completed(futures):
raw = future.result()
if raw is None:
continue
given += 1
raw["resume_token"] = token
yield raw
if given >= limit:
break
self.last_token = token
if progress_cb and not progress_cb(given):
logger.info("PMC bulk: выборка остановлена по запросу (%d)", given)
return
if token is None:
return
def _list_page(
self, token: str | None, want: int, start_after: str | None = None
) -> tuple[list[str], str | None]:
"""Одна страница листинга бакета: id статей и токен следующей страницы."""
url = f"{BUCKET}/?list-type=2&max-keys=1000"
if token:
# В токене бывают + и / — без экранирования S3 отвечает 400
url += f"&continuation-token={quote(token, safe='')}"
elif start_after:
url += f"&start-after={quote(start_after, safe='')}"
resp = self.client.get(url)
resp.raise_for_status()
ids = KEY_RE.findall(resp.text)[:want]
m = TOKEN_RE.search(resp.text)
return ids, (m.group(1) if m else None)
def _fetch_article(self, art_id: str) -> dict[str, Any] | None:
"""Текст статьи и метаданные; None — статья недоступна или пустая."""
try:
txt = self.client.get(f"{BUCKET}/{art_id}/{art_id}.txt")
if txt.status_code != 200 or len(txt.text) < MIN_CHARS:
return None
meta_resp = self.client.get(f"{BUCKET}/{art_id}/{art_id}.json")
meta = meta_resp.json() if meta_resp.status_code == 200 else {}
except Exception:
return None
return {"art_id": art_id, "meta": meta, "text": txt.text}
def transform(self, raw: dict[str, Any]) -> dict[str, Any]:
"""Привести статью к унифицированному формату документов корпуса."""
meta = raw.get("meta") or {}
art_id = raw.get("art_id", "")
pmcid = meta.get("pmcid") or art_id.split(".")[0]
if not pmcid:
return {}
citation = meta.get("citation") or ""
year_match = YEAR_RE.search(citation)
text = raw.get("text", "")
return {
"source": self.source_name,
"ext_id": f"pmc:{pmcid}",
"title": (meta.get("title") or pmcid)[:1000],
"authors": [],
"doi": meta.get("doi"),
"year": int(year_match.group(0)) if year_match else None,
"lang": "en",
"journal": citation[:300] or None,
"url": f"https://www.ncbi.nlm.nih.gov/pmc/articles/{pmcid}/",
"abstract": text[:2000],
"text": text,
"resume_token": raw.get("resume_token"),
}

View File

@@ -0,0 +1,156 @@
"""Юнит-тесты парсера CORE — чистая логика, без сети.
Стерегут то, ради чего парсер написан иначе остальных: отбор статей с полным
текстом, отсев дублей CORE (одна статья под разными id), нарезку выборки по
архивам-поставщикам и позицию продолжения «архив:смещение».
Нарезка по архивам — не украшение: связка условий через `AND` в API CORE не
работает (запрос по архивам вместе с годом отдаёт годы вперемешку), поэтому в
запрос идёт ровно одно условие, а годы отсекаются уже на нашей стороне.
"""
from core import DEFAULT_PROXY, MIN_CHARS, COREParser
LONG = "слово " * 400 # заведомо длиннее MIN_CHARS
SAMPLE = {
"id": 123456,
"title": "Нейронные сети в медицине",
"authors": [{"name": "Кузьмин, Ярослав Вадимович"}, {"name": "Савенок А."}],
"doi": "10.1234/x",
"yearPublished": 2019,
"language": {"code": "ru"},
"journals": [{"title": "Вестник"}],
"downloadUrl": "https://core.ac.uk/download/1.pdf",
"abstract": "Аннотация",
"fullText": LONG,
"resume_token": "1298:100",
}
def test_repos_from_or_expression():
q = "(repositories.id:1298 OR repositories.id:21908 OR repositories.id:949)"
assert COREParser._repos(q) == ["1298", "21908", "949"]
def test_repos_from_plain_list():
assert COREParser._repos("1298, 21908 949") == ["1298", "21908", "949"]
def test_repos_without_numbers_falls_back_to_word_search():
# Номеров нет — единственная часть «q»: обычный поиск по словам
assert COREParser._repos("нейронные сети") == ["q"]
assert COREParser._repos("") == ["q"]
def test_parse_token():
assert COREParser._parse_token("1298:4200") == ("1298", 4200)
assert COREParser._parse_token("q:300") == ("q", 300)
assert COREParser._parse_token(None) == (None, 0)
assert COREParser._parse_token("мусор") == (None, 0)
assert COREParser._parse_token("архив:смещение") == (None, 0)
def test_transform_maps_fields():
t = COREParser(proxy_url="").transform(SAMPLE)
assert t["ext_id"] == "core:123456"
assert t["source"] == "core"
assert t["year"] == 2019 and t["lang"] == "ru"
assert t["journal"] == "Вестник"
assert t["text"] == LONG
assert t["resume_token"] == "1298:100"
# Автор приходит одной строкой «Фамилия, Имя Отчество»
assert t["authors"][0]["last_name"] == "Кузьмин"
assert t["authors"][0]["initials"] == "Я.В."
def test_transform_skips_short_text():
assert COREParser(proxy_url="").transform(dict(SAMPLE, fullText="коротко")) == {}
assert COREParser(proxy_url="").transform(dict(SAMPLE, fullText=None)) == {}
assert COREParser(proxy_url="").transform(dict(SAMPLE, id=None)) == {}
def test_transform_drops_undefined_language():
# zz у CORE значит «язык не определён» — в корпус такое класть незачем
assert COREParser(proxy_url="").transform(dict(SAMPLE, language={"code": "zz"}))["lang"] is None
assert COREParser(proxy_url="").transform(dict(SAMPLE, language=None))["lang"] is None
def _parser_with_pages(pages):
"""Парсер, у которого выдача подменена заранее заготовленными страницами."""
p = COREParser(api_key="test", proxy_url="")
calls = []
def fake_page(part, query, offset):
calls.append((part, offset))
return pages.pop(0) if pages else []
p._page = fake_page # type: ignore[method-assign]
p.calls = calls # type: ignore[attr-defined]
return p
def test_fetch_drops_core_duplicates():
# Одна и та же статья под разными id — CORE так отдаёт всегда
page = [
{"id": 1, "fullText": LONG},
{"id": 2, "fullText": LONG},
{"id": 3, "fullText": LONG + "иное"},
{"id": 4, "fullText": "коротко"},
]
p = _parser_with_pages([page])
got = list(p.fetch(limit=10, query="repositories.id:1298"))
assert [w["id"] for w in got] == [1, 3]
def test_fetch_position_points_at_own_page():
# Токен обязан указывать на страницу, откуда пришла статья: пачка может
# прерваться на середине, и следующий прогон перечитает её целиком
pages = [[{"id": 1, "fullText": LONG}], [{"id": 2, "fullText": LONG + "два"}]]
p = _parser_with_pages(pages)
got = list(p.fetch(limit=10, query="repositories.id:1298"))
assert [w["resume_token"] for w in got] == ["1298:0", "1298:100"]
def test_fetch_walks_archives_and_resumes():
p = _parser_with_pages([[], []])
list(p.fetch(limit=10, query="(repositories.id:1298 OR repositories.id:949)",
resume_token="1298:300"))
# Начали с архива из токена и его смещения, пустой архив — переход к следующему
assert p.calls == [("1298", 300), ("949", 0)]
def test_fetch_ignores_token_of_unknown_archive():
p = _parser_with_pages([[]])
list(p.fetch(limit=10, query="repositories.id:1298", resume_token="99999:500"))
assert p.calls == [("1298", 0)]
def test_fetch_filters_years_on_our_side():
# Годы API не фильтрует, поэтому отсекаем сами — и только если заданы
page = [
{"id": 1, "fullText": LONG, "yearPublished": 2004},
{"id": 2, "fullText": LONG + "два", "yearPublished": 2015},
{"id": 3, "fullText": LONG + "три", "yearPublished": 2030},
]
p = _parser_with_pages([list(page)])
got = list(p.fetch(limit=10, query="repositories.id:1298", year_from=2010, year_to=2026))
assert [w["id"] for w in got] == [2]
p2 = _parser_with_pages([list(page)])
got2 = list(p2.fetch(limit=10, query="repositories.id:1298"))
assert [w["id"] for w in got2] == [1, 2, 3]
def test_min_chars_threshold_is_meaningful():
assert MIN_CHARS >= 1000
def test_proxy_choice(monkeypatch):
# CORE не отвечает на ключ с российских адресов — по умолчанию идём через
# sing-box, но вне compose-сети прокси отключают пустым значением
monkeypatch.delenv("CORE_PROXY_URL", raising=False)
assert COREParser._proxy(None) == DEFAULT_PROXY
monkeypatch.setenv("CORE_PROXY_URL", "socks5://иной:1080")
assert COREParser._proxy(None) == "socks5://иной:1080"
assert COREParser._proxy("") == "" # явное «без прокси» сильнее окружения

View File

@@ -0,0 +1,199 @@
"""Парсер русской Википедии из дампа Wikimedia.
Зачем в корпусе: студенты копируют из Википедии чаще, чем из научных статей, а
по объёму связного русского текста ей нет альтернативы среди доступного —
КиберЛенинка блокирует выкачку, eLIBRARY требует договора.
Отличие от остальных парсеров: `fetch` возвращает **генератор**, а не список.
Дамп — 5.6 ГБ и около 2 млн статей, держать их в памяти нельзя, поэтому
`run_parser` читает результат лениво и пишет пачками (см. bulk_writer.py).
Используется **multistream**-вариант дампа: он состоит из независимых bz2-блоков
по 100 статей, и к нему прилагается индекс со смещениями. Это принципиально —
обычный дамп читается только с начала, поэтому каждый следующий прогон
перечитывал всё уже залитое: на 20 тысячах статей это стоило 6 минут из 25
доступных, а на 100 тысячах съело бы весь бюджет и заливка встала бы совсем.
С multistream позиция продолжения — байтовое смещение, и прогон стартует
мгновенно.
Файлы качаются заранее и кладутся туда, где их видит воркер:
B=https://dumps.wikimedia.org/ruwiki/latest
UA="AcademicHelper/1.0 (https://academic.jze9.ru; noreply@jze9.ru)"
curl -L -C - --retry 100 -A "$UA" -o scripts/parsers/ruwiki-multistream.xml.bz2 \\
$B/ruwiki-latest-pages-articles-multistream.xml.bz2
Читать дамп прямо из сети не выйдет: Wikimedia обрывает долгие соединения
(проверено — обрыв на 32 МБ из 5.9 ГБ).
"""
import bz2
import logging
import os
import re
from collections.abc import Iterator
from typing import Any
from base import BaseParser, ProgressCallback
logger = logging.getLogger(__name__)
DEFAULT_DUMP = "/parsers/ruwiki-multistream.xml.bz2"
PAGE_RE = re.compile(r"<page>(.*?)</page>", re.DOTALL)
TITLE_RE = re.compile(r"<title>(.*?)</title>", re.DOTALL)
ID_RE = re.compile(r"<id>(\d+)</id>")
NS_RE = re.compile(r"<ns>(\d+)</ns>")
TEXT_RE = re.compile(r"<text[^>]*>(.*?)</text>", re.DOTALL)
REDIRECT_RE = re.compile(r"<redirect ")
# Разметку чистим регулярками: для отпечатков нужен связный текст, а не точное
# восстановление вёрстки — полноценный парсер вики-разметки тут неоправдан
TEMPLATE_RE = re.compile(r"\{\{[^{}]*\}\}")
CLEAN_RULES = [
(re.compile(r"\[\[[^\]|]*\|"), ""), # [[ссылка|текст → текст
(re.compile(r"\[\[|\]\]"), ""),
(re.compile(r"<ref[^>]*>.*?</ref>", re.DOTALL), " "),
(re.compile(r"<[^>]+>"), " "),
(re.compile(r"^[*#:;|!].*$", re.MULTILINE), " "), # списки и таблицы
(re.compile(r"^=+.*?=+$", re.MULTILINE), " "), # == заголовки разделов ==
(re.compile(r"'{2,}"), ""),
(re.compile(r"&[a-z]+;"), " "),
(re.compile(r"[ \t]+"), " "),
(re.compile(r"\n{2,}"), "\n"),
]
def clean_wikitext(raw: str) -> str:
"""Убрать вики-разметку, оставив читаемый текст статьи."""
text = raw
for _ in range(3): # шаблоны бывают вложенными
text = TEMPLATE_RE.sub(" ", text)
for rx, repl in CLEAN_RULES:
text = rx.sub(repl, text)
return text.strip()
class WikipediaRuParser(BaseParser):
"""Русская Википедия из локального дампа. Отдаёт статьи потоком."""
source_name = "wikipedia_ru"
# Признак для run_parser: результат читается лениво и пишется пачками,
# а не собирается в список и не идёт через add_document по одному
bulk = True
def fetch( # type: ignore[override]
self,
limit: int = 1000,
min_chars: int = 2000,
dump_path: str | None = None,
start_offset: int = 0,
progress_cb: ProgressCallback | None = None,
**_ignored: Any,
) -> Iterator[dict[str, Any]]:
"""Статьи основного пространства имён из multistream-дампа.
Args:
limit: сколько статей отдать
min_chars: минимальная длина текста — заготовки в корпусе бесполезны
dump_path: путь к дампу (по умолчанию /parsers/ruwiki-multistream.xml.bz2)
start_offset: байтовое смещение в файле, с которого продолжать.
Именно смещение, а не номер статьи: перечитывание дампа с начала
росло линейно и на сотне тысяч статей съедало весь бюджет прогона
progress_cb: см. base.ProgressCallback
"""
path = dump_path or os.environ.get("WIKIPEDIA_DUMP_PATH") or DEFAULT_DUMP
if not os.path.exists(path):
raise FileNotFoundError(
f"дамп Википедии не найден: {path} — скачайте его (см. модуль) "
f"или укажите WIKIPEDIA_DUMP_PATH"
)
given = 0
buf = ""
with open(path, "rb") as fh:
fh.seek(start_offset)
# Смещение блока, из которого пришли уже отданные статьи: его и
# сохраняем как позицию продолжения, чтобы ничего не потерять
block_offset = start_offset
decomp = bz2.BZ2Decompressor()
while given < limit:
part = fh.read(4 * 1024 * 1024)
if not part:
break
# В multistream-дампе потоки идут подряд: закончился один —
# начинаем следующий с того места, где предыдущий остановился
while part:
try:
raw = decomp.decompress(part)
except (OSError, EOFError):
return
if raw:
buf += raw.decode("utf-8", errors="replace")
if not decomp.eof:
break
part = decomp.unused_data
decomp = bz2.BZ2Decompressor()
while given < limit:
m = PAGE_RE.search(buf)
if not m:
break
page, buf = m.group(1), buf[m.end():]
doc = self._page_to_doc(page, min_chars)
if doc is None:
continue
given += 1
doc["dump_offset"] = block_offset
yield doc
if progress_cb and not progress_cb(given):
logger.info("Википедия: выборка остановлена по запросу (%d)", given)
return
# Всё разобранное отдано — следующая позиция продолжения здесь
block_offset = fh.tell() - len(buf.encode("utf-8", errors="ignore")) // 4
if len(buf) > 20 * 1024 * 1024: # страховка от разбухания
buf = buf[-1024 * 1024:]
def _page_to_doc(self, page: str, min_chars: int) -> dict[str, Any] | None:
"""Разобрать <page> в документ; None — страница нам не подходит."""
if REDIRECT_RE.search(page):
return None
ns = NS_RE.search(page)
if not ns or ns.group(1) != "0": # только статьи
return None
tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page)
if not (tm and im and xm):
return None
text = clean_wikitext(xm.group(1))
if len(text) < min_chars:
return None
return {"pageid": im.group(1), "title": tm.group(1), "text": text}
def transform(self, raw: dict[str, Any]) -> dict[str, Any]:
"""Привести статью к унифицированному формату документов корпуса."""
pid = raw.get("pageid")
if not pid:
return {}
return {
"source": self.source_name,
"ext_id": f"wikipedia_ru:{pid}",
"title": raw.get("title", ""),
"authors": [],
"year": None,
"lang": "ru",
"url": f"https://ru.wikipedia.org/?curid={pid}",
"abstract": (raw.get("text") or "")[:2000],
"text": raw.get("text", ""),
"dump_offset": raw.get("dump_offset"),
}

View File

@@ -0,0 +1,29 @@
"""Позиция продолжения массовой заливки в источнике.
Revision ID: 007
Revises: 006
Create Date: 2026-09-05
Массовые источники (дамп Википедии, бакет PMC) заливаются частями: за один
прогон берётся столько статей, сколько влезает в бюджет времени таска. Чтобы
следующий прогон продолжал с того же места, а не перечитывал залитое заново,
позиция сохраняется прямо в источнике: для Википедии это номер статьи в дампе,
для PMC — токен страницы бакета.
"""
import sqlalchemy as sa
from alembic import op
revision = "007"
down_revision = "006"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.add_column("parse_sources", sa.Column("resume_token", sa.Text(), nullable=True))
def downgrade() -> None:
op.drop_column("parse_sources", "resume_token")

View File

@@ -50,8 +50,10 @@ logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin", tags=["admin"], dependencies=[Depends(get_admin_user)]) router = APIRouter(prefix="/admin", tags=["admin"], dependencies=[Depends(get_admin_user)])
# Типы источников, для которых в worker-indexer есть парсер (index.run_parser) # Типы источников, для которых в worker-indexer есть парсер (index.run_parser).
SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc") # wikipedia_ru и pmc_bulk — массовые: читают дамп/бакет потоком и пишут пачками,
# продолжая с сохранённой позиции (parse_sources.resume_token)
SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc", "wikipedia_ru", "pmc_bulk")
# Прогон в этих статусах ещё может двигаться — повторно запускать источник нельзя # Прогон в этих статусах ещё может двигаться — повторно запускать источник нельзя
RUN_ACTIVE_STATUSES = ("queued", "running") RUN_ACTIVE_STATUSES = ("queued", "running")
@@ -925,6 +927,29 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
) )
).scalars().all() ).scalars().all()
# Зависшая проверка: числится в работе, но воркер давно её не трогал.
# Пользователь видит вечное «обрабатывается» — молча и без следов в логах,
# поэтому такие задачи должны быть видны в отладке (нашлись экземпляры,
# висевшие с мая).
#
# Возраст считает сама база: колонки в PostgreSQL — timestamptz, а модель
# объявляет их без зоны, и передача сюда Python-времени ломается на
# несовпадении (asyncpg отвергает aware-параметр, naive даёт неверный сдвиг).
stale_tasks = (
await db.execute(
text("""
SELECT public_id, type, created_at,
round(extract(epoch FROM now() - coalesce(updated_at, created_at))
/ 3600.0, 1) AS stuck_hours
FROM tasks
WHERE status = 'processing'
AND coalesce(updated_at, created_at) < now() - interval '2 hours'
ORDER BY created_at DESC
LIMIT 20
""")
)
).all()
return { return {
"generated_at": datetime.utcnow().isoformat(timespec="seconds"), "generated_at": datetime.utcnow().isoformat(timespec="seconds"),
"celery": celery_state, "celery": celery_state,
@@ -956,6 +981,15 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
} }
for t in failed_tasks for t in failed_tasks
], ],
"stale_tasks": [
{
"public_id": t.public_id,
"type": t.type,
"created_at": t.created_at,
"stuck_hours": float(t.stuck_hours),
}
for t in stale_tasks
],
# Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на # Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на
# поведение воркеров — при разборе «почему не так считает» нужны первыми # поведение воркеров — при разборе «почему не так считает» нужны первыми
"config": { "config": {

View File

@@ -35,6 +35,9 @@ class ParseSource(Base):
docs_added: Mapped[int] = mapped_column(default=0) docs_added: Mapped[int] = mapped_column(default=0)
# Последний (пока идёт — текущий) прогон, см. ParseRun # Последний (пока идёт — текущий) прогон, см. ParseRun
last_run_id: Mapped[int | None] = mapped_column(nullable=True) last_run_id: Mapped[int | None] = mapped_column(nullable=True)
# Позиция продолжения для массовых источников: номер статьи в дампе
# Википедии или токен страницы бакета PMC
resume_token: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(server_default=func.now()) created_at: Mapped[datetime] = mapped_column(server_default=func.now())
def __repr__(self) -> str: def __repr__(self) -> str:

View File

@@ -120,7 +120,9 @@ class ParseRunDetail(ParseRunResponse):
# ─── Источники парсинга ─────────────────────────────────────────────────────── # ─── Источники парсинга ───────────────────────────────────────────────────────
class ParseSourceCreate(BaseModel): class ParseSourceCreate(BaseModel):
source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") source_type: str = Field(
description="openalex / cyberleninka / arxiv / pmc / wikipedia_ru / pmc_bulk"
)
name: str = Field(min_length=1, max_length=255) name: str = Field(min_length=1, max_length=255)
query: str | None = None query: str | None = None
lang: str | None = None lang: str | None = None
@@ -164,7 +166,9 @@ class ParseSourceResponse(BaseModel):
class ParseSourceBulkCreate(BaseModel): class ParseSourceBulkCreate(BaseModel):
"""Пакетное добавление: один тип/лимит, много запросов (по строке на тему).""" """Пакетное добавление: один тип/лимит, много запросов (по строке на тему)."""
source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") source_type: str = Field(
description="openalex / cyberleninka / arxiv / pmc / wikipedia_ru / pmc_bulk"
)
queries: list[str] = Field(min_length=1, max_length=500) queries: list[str] = Field(min_length=1, max_length=500)
name_prefix: str = "" name_prefix: str = ""
lang: str | None = None lang: str | None = None

View File

@@ -32,7 +32,11 @@ const SOURCE_TYPES = [
{ value: 'openalex', label: 'OpenAlex' }, { value: 'openalex', label: 'OpenAlex' },
{ value: 'cyberleninka', label: 'КиберЛенинка' }, { value: 'cyberleninka', label: 'КиберЛенинка' },
{ value: 'arxiv', label: 'arXiv' }, { value: 'arxiv', label: 'arXiv' },
{ value: 'pmc', label: 'PubMed Central' }, { value: 'pmc', label: 'PubMed Central (по теме)' },
// Массовые: читают дамп/бакет потоком и продолжают с места остановки,
// поэтому запускаются повторно до исчерпания источника
{ value: 'wikipedia_ru', label: 'Википедия ru (дамп, массово)' },
{ value: 'pmc_bulk', label: 'PubMed Central (бакет, массово)' },
]; ];
const STATUS_COLORS: Record<string, string> = { const STATUS_COLORS: Record<string, string> = {

View File

@@ -85,6 +85,68 @@ def winnow(text: str, k: int = 5, window: int = 4) -> set[int]:
return fingerprint return fingerprint
def winnow_ordered(text: str, k: int = 5, window: int = 4) -> list[int]:
"""То же, что winnow, но с сохранением порядка появления отпечатков в тексте.
Нужно там, где отпечатки приходится обрезать по лимиту: множество не хранит
порядка, и `list(winnow(text))[:limit]` берёт произвольное подмножество —
у длинного документа целые куски остаются без покрытия, и списывание из них
не находится. Со списком в порядке текста обрезку можно делать равномерной
(см. sample_evenly).
Args:
text: Исходный текст
k: Размер k-граммы
window: Размер скользящего окна
Returns:
Отпечатки в порядке появления в тексте, без повторов
"""
tokens = text.lower().split()
if len(tokens) < k:
return []
hashes = [hash_ngram(ng) for ng in get_ngrams(tokens, k)]
if not hashes:
return []
ordered: list[int] = []
seen: set[int] = set()
prev_min_idx = -1
for i in range(len(hashes) - window + 1):
window_hashes = hashes[i : i + window]
min_val = min(window_hashes)
min_idx = i + window_hashes.index(min_val)
if min_idx != prev_min_idx:
if min_val not in seen:
seen.add(min_val)
ordered.append(min_val)
prev_min_idx = min_idx
return ordered
def sample_evenly(items: list[int], limit: int) -> list[int]:
"""Оставить не больше limit элементов, равномерно по всей длине списка.
Берём каждый n-й элемент, а не первые limit штук: обрезка «с начала»
оставила бы без отпечатков весь конец документа.
Args:
items: Отпечатки в порядке текста (winnow_ordered)
limit: Максимум отпечатков; 0 или меньше — не ограничивать
Returns:
Подсписок длиной не больше limit, сохраняющий порядок
"""
if limit <= 0 or len(items) <= limit:
return items
step = len(items) / limit
return [items[int(i * step)] for i in range(limit)]
def jaccard_similarity(fp_a: set[int], fp_b: set[int]) -> float: def jaccard_similarity(fp_a: set[int], fp_b: set[int]) -> float:
""" """
Коэффициент Жаккара для двух fingerprint'ов. Коэффициент Жаккара для двух fingerprint'ов.

View File

@@ -0,0 +1,132 @@
"""Пакетная запись документов корпуса — путь для массовых источников.
Обычный `index.add_document` пишет по одному документу через ORM: удобно для
парсеров, отдающих сотни статей, но на миллионах это тупик. Замер на проде:
построчная вставка отпечатков — 6.7 тыс. строк/с, COPY — 21 тыс. строк/с, а
миллион статей это ~2 млрд отпечатков. Разница между сутками и неделями.
Поэтому массовые источники (дампы Википедии, бакет PMC) идут сюда: документы
вставляются одной командой с ON CONFLICT, отпечатки — через COPY.
Elasticsearch и эмбеддинги здесь намеренно не трогаются: L1 и L2 начинают
работать сразу, а векторы для L3 досчитываются отдельно
(`scripts/ops/reembed_missing.py`) — иначе заливка упирается в скорость
сервера эмбеддингов и тормозит на порядок.
"""
import contextlib
import io
import logging
import time
from typing import Any
import psycopg2
from app.algorithms.winnowing import sample_evenly, winnow_ordered
from app.config import settings
logger = logging.getLogger(__name__)
# Сколько попыток пережить обрыв соединения с базой. Массовая заливка идёт
# часами, и разрыв сети за это время — норма, а не исключение.
DB_RETRIES = 5
def connect():
"""Отдельное подключение psycopg2: COPY недоступен через ORM-сессию."""
return psycopg2.connect(
host=settings.POSTGRES_HOST,
port=settings.POSTGRES_PORT,
dbname=settings.POSTGRES_DB,
user=settings.POSTGRES_USER,
password=settings.POSTGRES_PASSWORD,
)
def write_batch(conn, docs: list[dict[str, Any]], fp_limit: int) -> tuple[Any, dict[str, int]]:
"""Записать пачку документов с отпечатками, пережив обрыв соединения.
Args:
conn: активное подключение psycopg2 (может быть заменено при обрыве)
docs: документы в унифицированном формате парсеров; нужны ext_id,
source, title и text — остальное необязательно
fp_limit: максимум отпечатков на документ
Returns:
(соединение, счётчики) — соединение может оказаться новым, если
пришлось переподключаться; счётчики: added / duplicates / fingerprints
"""
if not docs:
return conn, {"added": 0, "duplicates": 0, "fingerprints": 0}
for attempt in range(1, DB_RETRIES + 1):
try:
return conn, _write_once(conn, docs, fp_limit)
except psycopg2.OperationalError as e:
wait = min(60, 5 * attempt)
logger.warning(
"bulk_writer: база недоступна (%s), повтор через %sс [%s/%s]",
str(e).strip()[:80], wait, attempt, DB_RETRIES,
)
time.sleep(wait)
with contextlib.suppress(Exception):
conn.close()
try:
conn = connect()
except Exception:
continue
raise RuntimeError(f"не удалось записать пачку после {DB_RETRIES} попыток")
def _write_once(conn, docs: list[dict[str, Any]], fp_limit: int) -> dict[str, int]:
"""Одна попытка записи; счётчики возвращаются только при успехе."""
added = duplicates = fingerprints = 0
with conn.cursor() as cur:
cur.executemany(
"""INSERT INTO documents (ext_id, source, title, doi, year, lang, journal,
abstract, url, authors, indexed_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,'[]'::json, now())
ON CONFLICT (ext_id) DO NOTHING""",
[
(
d["ext_id"], d["source"], (d.get("title") or "")[:1000], d.get("doi"),
d.get("year"), d.get("lang"), (d.get("journal") or None),
(d.get("text") or "")[:2000], d.get("url"),
)
for d in docs
],
)
cur.execute(
"SELECT id, ext_id FROM documents WHERE ext_id = ANY(%s)",
([d["ext_id"] for d in docs],),
)
id_by_ext = {ext: doc_id for doc_id, ext in cur.fetchall()}
buf = io.StringIO()
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():
duplicates += 1
continue
hashes = sample_evenly(winnow_ordered(d.get("text") or ""), fp_limit)
if not hashes:
continue
for pos, h in enumerate(hashes):
buf.write(f"{doc_id}\t{h}\t{pos}\n")
fingerprints += len(hashes)
added += 1
buf.seek(0)
if added:
cur.copy_from(buf, "fingerprints", columns=("doc_id", "hash_value", "position"))
conn.commit()
return {"added": added, "duplicates": duplicates, "fingerprints": fingerprints}

View File

@@ -59,6 +59,12 @@ class Settings(BaseSettings):
FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным
EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов
# Сколько документов массового источника пишется одной транзакцией.
# Больше — меньше накладных расходов, но длиннее транзакция: 500 статей по
# 2000 отпечатков это миллион строк за раз, и на таком объёме соединение
# обрывалось. 150 — компромисс, проверенный на заливке Википедии.
BULK_WRITE_BATCH: int = 150
# Бюджет времени одного прогона index.run_parser. RabbitMQ рвёт канал # Бюджет времени одного прогона index.run_parser. RabbitMQ рвёт канал
# consumer'а, не сделавшего ack за consumer_timeout (по умолчанию 1800с): # consumer'а, не сделавшего ack за consumer_timeout (по умолчанию 1800с):
# воркер падает, сообщение передоставляется, таск начинается заново — # воркер падает, сообщение передоставляется, таск начинается заново —

View File

@@ -78,6 +78,7 @@ class ParseSource(Base):
last_run_at: Mapped[datetime | None] = mapped_column(nullable=True) last_run_at: Mapped[datetime | None] = mapped_column(nullable=True)
docs_added: Mapped[int] = mapped_column(default=0) docs_added: Mapped[int] = mapped_column(default=0)
last_run_id: Mapped[int | None] = mapped_column(nullable=True) last_run_id: Mapped[int | None] = mapped_column(nullable=True)
resume_token: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(server_default=func.now()) created_at: Mapped[datetime] = mapped_column(server_default=func.now())

View File

@@ -1,5 +1,6 @@
"""Celery задачи индексации документов и проверки плагиата (уровни 1-2).""" """Celery задачи индексации документов и проверки плагиата (уровни 1-2)."""
import contextlib
import io import io
from datetime import UTC, datetime from datetime import UTC, datetime
from pathlib import Path from pathlib import Path
@@ -10,7 +11,7 @@ from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError from sqlalchemy.exc import IntegrityError
from app.algorithms.minhash import add_to_lsh, find_similar from app.algorithms.minhash import add_to_lsh, find_similar
from app.algorithms.winnowing import winnow from app.algorithms.winnowing import sample_evenly, winnow, winnow_ordered
from app.celery_app import celery_app from app.celery_app import celery_app
from app.config import settings from app.config import settings
from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status
@@ -284,8 +285,9 @@ def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[
# Вычислить fingerprints # Вычислить fingerprints
text = doc_data.get("full_text") or doc_data.get("abstract", "") or "" text = doc_data.get("full_text") or doc_data.get("abstract", "") or ""
if text: if text:
fp = winnow(text) fingerprints_to_add = sample_evenly(
fingerprints_to_add = list(fp)[: settings.MAX_FINGERPRINTS_PER_DOC] winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC
)
for i, hash_val in enumerate(fingerprints_to_add): for i, hash_val in enumerate(fingerprints_to_add):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i)) session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
@@ -375,7 +377,7 @@ def store_full_text(doc_id: int, text: str) -> dict[str, Any]:
content_type="text/plain; charset=utf-8", content_type="text/plain; charset=utf-8",
) )
hashes = list(winnow(text))[: settings.MAX_FINGERPRINTS_PER_DOC] hashes = sample_evenly(winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC)
with db_session() as session: with db_session() as session:
doc = session.get(Document, doc_id) doc = session.get(Document, doc_id)
@@ -530,6 +532,24 @@ def _stage_work(
logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}") logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}")
def _last_pmc_key() -> str | None:
"""Ключ бакета после последней залитой статьи PMC — старт листинга.
Бакет отдаётся в лексикографическом порядке ключей, и ext_id вида
`pmc:PMC10000000` соответствует ключу `PMC10000000.1`. Берём максимальный
залитый и продолжаем после него.
"""
from sqlalchemy import text as sql_text
with db_session() as session:
row = session.execute(sql_text(
"SELECT max(ext_id) FROM documents WHERE ext_id LIKE 'pmc:PMC%'"
)).first()
if not row or not row[0]:
return None
return row[0].split(":", 1)[1] + ".1"
def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, Any]]: def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, Any]]:
"""Инстанс парсера и совместимые с его fetch() аргументы по типу источника.""" """Инстанс парсера и совместимые с его fetch() аргументы по типу источника."""
limit = cfg["limit"] limit = cfg["limit"]
@@ -559,12 +579,200 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A
"year_from": cfg.get("year_from"), "year_from": cfg.get("year_from"),
"year_to": cfg.get("year_to"), "year_to": cfg.get("year_to"),
} }
elif source_type == "wikipedia_ru":
# Массовый источник: позиция — байтовое смещение в multistream-дампе.
# Раньше хранился номер статьи, и каждый прогон перечитывал дамп с
# начала: на 20 тысячах это стоило 6 минут из 25, дальше росло линейно
from wikipedia_ru import WikipediaRuParser as P
kwargs = {
"limit": limit,
"dump_path": query or None, # query = путь к дампу, если задан
"start_offset": int(cfg.get("resume_token") or 0),
}
elif source_type == "core":
# Массовый источник: позиция — "год:смещение". Выборка режется по
# годам, потому что offset у CORE упирается в 100 000 (см. core.py)
from core import COREParser as P
kwargs = {
"limit": limit,
"query": query,
"year_from": cfg.get("year_from"),
"year_to": cfg.get("year_to"),
"resume_token": cfg.get("resume_token"),
}
elif source_type == "pmc_bulk":
from pmc_bulk import PMCBulkParser as P
kwargs = {
"limit": limit,
"resume_token": cfg.get("resume_token"),
# Токена ещё нет (первый прогон после ручных заливок) — начинаем
# после последней уже залитой статьи, иначе листинг часами
# перемалывает существующие как дубли
"start_after": None if cfg.get("resume_token") else _last_pmc_key(),
}
else: else:
raise ValueError(f"неизвестный тип источника: {source_type}") raise ValueError(f"неизвестный тип источника: {source_type}")
return P(), kwargs return P(), kwargs
def _ingest_one_by_one(parser: Any, raw_docs: Any, prog: Any, run_id: int) -> tuple[bool, bool]:
"""Обычный путь: документ за документом через add_document.
Подходит источникам, отдающим сотни статей: работает ORM-логика с
дедупликацией, индексацией в Elasticsearch и обогащением полным текстом.
Returns:
(отменён, исчерпан бюджет времени)
"""
cancelled = exhausted = False
prog.stage = "index"
prog.count_fetched(len(raw_docs))
prog.log("info", f"получено {len(raw_docs)} документов, индексация")
_write_run(run_id, **prog.snapshot())
# Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu
# кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений,
# а не на каждый документ.
embed_batch: list[int] = []
batch_size = settings.EMBED_BATCH_SIZE
def _flush_embed() -> None:
if embed_batch:
celery_app.send_task(
"gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu"
)
embed_batch.clear()
for raw in raw_docs:
if not cancelled and prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации")
if cancelled or exhausted:
break
try:
doc = parser.transform(raw)
if not (doc and doc.get("title") and doc.get("ext_id")):
prog.count_result("skipped")
continue
result = add_document(doc, dispatch_embed=False)
prog.count_result(result.get("status", "error"))
if result.get("status") == "indexed":
embed_batch.append(result["doc_id"])
if len(embed_batch) >= batch_size:
_flush_embed()
except Exception as e:
prog.count_result("error")
prog.log("warning", f"документ пропущен: {e}")
logger.warning(f"run_parser: ошибка документа: {e}")
if prog.should_flush():
cancelled = _write_run(run_id, **prog.snapshot())
prog.mark_flushed()
if cancelled:
prog.log("warning", "отмена по запросу из админки")
_flush_embed()
return cancelled, exhausted
def _ingest_bulk(
parser: Any, raw_docs: Any, prog: Any, run_id: int, source_id: int
) -> tuple[bool, bool]:
"""Массовый путь: статьи приходят потоком и пишутся пачками через COPY.
Отличия от обычного пути и почему они нужны:
- парсер отдаёт генератор, поэтому счётчик «получено» растёт по ходу дела,
а не известен заранее;
- запись идёт пакетно (bulk_writer), иначе миллионы отпечатков занимают
недели вместо суток;
- позиция продолжения сохраняется в источнике, чтобы следующий прогон
начинал с места остановки, а не перечитывал дамп с начала.
Эмбеддинги здесь не диспатчатся: они считаются заметно медленнее заливки и
делали бы её узким местом. Векторы досчитываются отдельно —
`scripts/ops/reembed_missing.py`.
Returns:
(отменён, исчерпан бюджет времени)
"""
from app.bulk_writer import connect, write_batch
from app.models import ParseSource
cancelled = exhausted = False
prog.stage = "index"
prog.log("info", "массовый источник: запись пачками")
_write_run(run_id, **prog.snapshot())
conn = connect()
batch: list[dict[str, Any]] = []
resume_token: str | None = None
batch_size = settings.BULK_WRITE_BATCH
def save_resume() -> None:
"""Запомнить позицию в источнике — с неё продолжит следующий прогон."""
if resume_token is None:
return
with db_session() as session:
src = session.get(ParseSource, source_id)
if src:
src.resume_token = str(resume_token)
session.commit()
def flush() -> bool:
"""Записать накопленное; True — попросили остановиться."""
nonlocal conn, batch
if not batch:
return False
conn, counts = write_batch(conn, batch, settings.MAX_FINGERPRINTS_PER_DOC)
for _ in range(counts["added"]):
prog.count_result("indexed")
for _ in range(counts["duplicates"]):
prog.count_result("duplicate")
batch = []
save_resume()
stop = _write_run(run_id, **prog.snapshot())
prog.mark_flushed()
return stop
try:
for raw in raw_docs:
if prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан")
break
try:
doc = parser.transform(raw)
except Exception as e:
prog.count_result("error")
logger.warning(f"run_parser bulk: ошибка документа: {e}")
continue
if not (doc and doc.get("ext_id") and doc.get("text")):
prog.count_result("skipped")
continue
# Позиция продолжения: у Википедии — номер статьи в дампе,
# у PMC — токен страницы бакета
resume_token = doc.get("dump_offset") or doc.get("resume_token") or resume_token
batch.append(doc)
prog.count_fetched(prog.fetched + 1)
if len(batch) >= batch_size and flush():
cancelled = True
prog.log("warning", "отмена по запросу из админки")
break
if not cancelled and flush():
cancelled = True
finally:
with contextlib.suppress(Exception):
conn.close()
return cancelled, exhausted
def _write_run(run_id: int, **fields: Any) -> bool: def _write_run(run_id: int, **fields: Any) -> bool:
"""Записать прогресс прогона и вернуть True, если запрошена отмена. """Записать прогресс прогона и вернуть True, если запрошена отмена.
@@ -621,6 +829,8 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any
"year_from": src.year_from, "year_from": src.year_from,
"year_to": src.year_to, "year_to": src.year_to,
"limit": src.limit, "limit": src.limit,
# Массовые источники продолжают с сохранённой позиции
"resume_token": src.resume_token,
} }
src.last_status = "running" src.last_status = "running"
src.last_error = None src.last_error = None
@@ -675,53 +885,13 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any
# fetch+transform без записи в JSONL — работаем in-memory # fetch+transform без записи в JSONL — работаем in-memory
raw_docs = parser.fetch(**fetch_kwargs, progress_cb=on_fetch_progress) raw_docs = parser.fetch(**fetch_kwargs, progress_cb=on_fetch_progress)
prog.stage = "index" if getattr(parser, "bulk", False):
prog.count_fetched(len(raw_docs)) # Массовые источники (дамп Википедии, бакет PMC) отдают генератор:
prog.log("info", f"получено {len(raw_docs)} документов, индексация") # миллионы статей нельзя ни держать в памяти, ни писать по одной
_write_run(run_id, **prog.snapshot()) # через ORM — для них отдельный путь с пакетной записью
cancelled, exhausted = _ingest_bulk(parser, raw_docs, prog, run_id, source_id)
# Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu else:
# кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений, cancelled, exhausted = _ingest_one_by_one(parser, raw_docs, prog, run_id)
# а не на каждый документ.
embed_batch: list[int] = []
batch_size = settings.EMBED_BATCH_SIZE
def _flush_embed() -> None:
if embed_batch:
celery_app.send_task(
"gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu"
)
embed_batch.clear()
for raw in raw_docs:
if not cancelled and prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации")
if cancelled or exhausted:
break
try:
doc = parser.transform(raw)
if not (doc and doc.get("title") and doc.get("ext_id")):
prog.count_result("skipped")
continue
result = add_document(doc, dispatch_embed=False)
prog.count_result(result.get("status", "error"))
if result.get("status") == "indexed":
embed_batch.append(result["doc_id"])
if len(embed_batch) >= batch_size:
_flush_embed()
except Exception as e:
prog.count_result("error")
prog.log("warning", f"документ пропущен: {e}")
logger.warning(f"run_parser: ошибка документа: {e}")
if prog.should_flush():
cancelled = _write_run(run_id, **prog.snapshot())
prog.mark_flushed()
if cancelled:
prog.log("warning", "отмена по запросу из админки")
_flush_embed()
except Exception as e: except Exception as e:
error_msg = str(e)[:500] error_msg = str(e)[:500]
prog.log("error", f"прогон упал: {error_msg}") prog.log("error", f"прогон упал: {error_msg}")

View File

@@ -12,3 +12,4 @@ langdetect==1.0.9
beautifulsoup4==4.12.3 # парсер CyberLeninka: детали статьи (fetch_article_details) beautifulsoup4==4.12.3 # парсер CyberLeninka: детали статьи (fetch_article_details)
pydantic-settings==2.2.1 pydantic-settings==2.2.1
httpx==0.27.0 httpx==0.27.0
socksio==1.0.0 # SOCKS5-прокси для httpx (CORE блокирует запросы из РФ, см. scripts/parsers/core.py)

View File

@@ -5,7 +5,9 @@ from app.algorithms.winnowing import (
get_ngrams, get_ngrams,
hash_ngram, hash_ngram,
jaccard_similarity, jaccard_similarity,
sample_evenly,
winnow, winnow,
winnow_ordered,
) )
# Достаточно длинный текст, чтобы окно Winnowing реально отработало # Достаточно длинный текст, чтобы окно Winnowing реально отработало
@@ -77,3 +79,57 @@ def test_compute_similarity_partial_overlap_is_between():
modified = LONG + " добавлен ещё один совершенно новый хвост предложения здесь" modified = LONG + " добавлен ещё один совершенно новый хвост предложения здесь"
sim = compute_similarity(LONG, modified) sim = compute_similarity(LONG, modified)
assert 0.0 < sim < 1.0 assert 0.0 < sim < 1.0
# ─── Обрезка отпечатков по лимиту ────────────────────────────────────────────
# Регрессия: winnow() возвращает set, и list(fp)[:limit] брал произвольное
# подмножество — у длинного документа целые куски оставались без покрытия,
# и списывание из них не находилось.
def test_winnow_ordered_matches_winnow_by_content():
"""Тот же набор отпечатков, что и у winnow, только с порядком."""
assert set(winnow_ordered(LONG)) == winnow(LONG)
def test_winnow_ordered_has_no_duplicates():
ordered = winnow_ordered(LONG)
assert len(ordered) == len(set(ordered))
def test_winnow_ordered_follows_text_order():
"""Отпечатки начала текста идут раньше отпечатков продолжения."""
tail = " совершенно другой хвост про выпечку хлеба и закваску в тёплой печи"
ordered = winnow_ordered(LONG + tail)
head_prints = set(winnow_ordered(LONG))
positions = [i for i, h in enumerate(ordered) if h in head_prints]
# Отпечатки первой половины сосредоточены в начале списка, а не разбросаны
assert max(positions) < len(ordered)
assert positions[0] == 0
def test_sample_evenly_keeps_everything_under_limit():
items = [1, 2, 3]
assert sample_evenly(items, 10) == items
assert sample_evenly(items, 0) == items # 0 = без ограничения
def test_sample_evenly_respects_limit():
items = list(range(1000))
assert len(sample_evenly(items, 100)) == 100
def test_sample_evenly_covers_whole_document():
"""Главное свойство: выборка растянута по всей длине, а не обрезана с начала."""
items = list(range(1000))
sampled = sample_evenly(items, 10)
assert sampled[0] == 0
assert sampled[-1] >= 900 # хвост документа тоже покрыт
assert sampled == sorted(sampled) # порядок сохранён
def test_sample_evenly_spreads_uniformly():
items = list(range(100))
sampled = sample_evenly(items, 10)
gaps = [b - a for a, b in zip(sampled, sampled[1:], strict=False)]
assert max(gaps) - min(gaps) <= 1 # шаг ровный