Compare commits

..

8 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
8 changed files with 560 additions and 9 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

@@ -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

@@ -15,11 +15,36 @@
Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи, Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи,
прогон закрывается с нулём добавленных, и скрипт больше его не трогает прогон закрывается с нулём добавленных, и скрипт больше его не трогает
(--stop-when-empty, по умолчанию включено). (--stop-when-empty, по умолчанию включено).
Перед подсчётом «занято ли» снимает зависшие прогоны: если worker-indexer упал
или перезапустился, прогон навсегда остаётся в статусе running с замолчавшим
heartbeat — без этой чистки сторож видит «занято» бесконечно и не запускает
вообще ничего, для всех источников сразу (порог тот же, что и в панели
отладки — services/api/app/api/admin.py, stale_cutoff).
Молчащий heartbeat сам по себе ещё не значит «мёртв»: пачка COPY на большой
базе идёт десятки минут без единого тика. Поэтому прогон снимается, только
если его таска нет среди выполняющихся на воркерах queue.index; не ответил
хоть один такой воркер — не снимается ничего. Иначе сторож снимал живой
медленный прогон и тут же запускал его дубль по тем же статьям.
""" """
import argparse import argparse
BULK_TYPES = ("wikipedia_ru", "pmc_bulk") 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: def main() -> None:
@@ -38,7 +63,44 @@ def main() -> None:
types = tuple(t.strip() for t in args.types.split(",") if t.strip()) types = tuple(t.strip() for t in args.types.split(",") if t.strip())
live = live_task_ids(celery_app)
with db_session() as s: 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( active = s.execute(text(
"SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')" "SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')"
)).scalar_one() )).scalar_one()
@@ -53,12 +115,25 @@ def main() -> None:
started = [] started = []
for source_id, source_type, last_run_id in sources: 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: if last_run_id and not args.keep_going_empty:
prev = s.get(ParseRun, last_run_id) last_two = s.execute(text("""
if prev and prev.status == "done" and prev.added == 0 and prev.fetched == 0: SELECT status, added, fetched FROM parse_runs
print(f"{source_type}: источник исчерпан, пропускаем") 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 continue
src = s.get(ParseSource, source_id) src = s.get(ParseSource, source_id)

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

View File

@@ -14,7 +14,7 @@
import logging import logging
import re import re
from collections.abc import Iterator from collections.abc import Iterator
from concurrent.futures import ThreadPoolExecutor from concurrent.futures import ThreadPoolExecutor, as_completed
from typing import Any from typing import Any
from urllib.parse import quote from urllib.parse import quote
@@ -39,8 +39,11 @@ class PMCBulkParser(BaseParser):
def __init__(self, workers: int = 12) -> None: def __init__(self, workers: int = 12) -> None:
super().__init__() super().__init__()
self.workers = workers self.workers = workers
# 60с на файл при бюджете задачи в 1500с — один залипший запрос съедал
# почти весь бюджет. Обычный GET сюда укладывается в доли секунды,
# 12с — с большим запасом на джиттер, но без риска съесть весь прогон
self.client = httpx.Client( self.client = httpx.Client(
timeout=60, timeout=12,
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"}, headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"},
) )
# Позиция листинга: бакет отдаётся страницами, и продолжать прогон # Позиция листинга: бакет отдаётся страницами, и продолжать прогон
@@ -82,8 +85,13 @@ class PMCBulkParser(BaseParser):
logger.info("PMC bulk: бакет закончился") logger.info("PMC bulk: бакет закончился")
return return
# as_completed вместо map(): map() отдаёт результаты строго по
# порядку отправки, поэтому один залипший запрос блокирует все
# уже готовые — даже если остальные 11 потоков давно отработали
with ThreadPoolExecutor(max_workers=self.workers) as pool: with ThreadPoolExecutor(max_workers=self.workers) as pool:
for raw in pool.map(self._fetch_article, ids): 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: if raw is None:
continue continue
given += 1 given += 1

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

@@ -589,6 +589,17 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A
"dump_path": query or None, # query = путь к дампу, если задан "dump_path": query or None, # query = путь к дампу, если задан
"start_offset": int(cfg.get("resume_token") or 0), "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": elif source_type == "pmc_bulk":
from pmc_bulk import PMCBulkParser as P from pmc_bulk import PMCBulkParser as P
kwargs = { kwargs = {

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)