Compare commits

..

5 Commits

Author SHA1 Message Date
jze9
25e55a3f7d fix(detection): пофрагментная локализация уровня 1 + полное хранение отпечатков
Две связанные проблемы, ломавшие качество детекции:

1. Уровень 1 (Winnowing) репортил весь документ одним совпадением на позиции
   0..длина — в отчёте нельзя было понять, ГДЕ плагиат. Теперь winnow'им
   каждый фрагмент и находим источник + позицию для каждого, так отчёт
   показывает "символы A-B скопированы из источника X".

2. MAX_FINGERPRINTS_PER_DOC=500 обрезал отпечатки статьи (у типичной статьи
   их ~3600), причём произвольную выборку. Из-за этого скопированный фрагмент
   почти не разделял отпечатки с источником (проверено: хранилось ~14%,
   фрагмент находил 17% своих хэшей → не срабатывало). Winnowing рассчитан на
   ПОЛНОЕ хранение; поднял лимит до 20000.

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

Проверено end-to-end: документ с дословной вставкой из статьи корпуса →
вставка локализована (chars 1027-2385 = 100%, источник атрибутирован);
чисто оригинальный текст → 0% (нет ложных срабатываний).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-23 19:57:05 +05:00
jze9
ef2723e1d3 feat(indexer): скачивание полного текста статей + батчинг эмбеддингов
Раньше корпус состоял только из метаданных: парсеры отдают full_text=None,
fingerprints/эмбеддинги считались из аннотаций, MinIO статьями не наполнялся
вовсе. Проверка плагиата шла против абстрактов, а не тел статей.

Добавлено:
- app/fulltext.py: best-effort скачивание PDF по url источника (стрим с
  лимитом размера, детект PDF по content-type/magic), извлечение текста
  через PyMuPDF.
- index.enrich_full_text: новая задача — качает полный текст, кладёт в MinIO
  (documents/corpus/{id}.txt), пересчитывает fingerprints по полному тексту,
  обновляет MinHash. Диспатчится из add_document при FETCH_FULL_TEXT=true.
- openalex: предпочитаем прямую ссылку на PDF (best_oa_location.pdf_url)
  вместо лендинга — покрытие full-text выросло с 17% до 33% на выборке.
- run_parser батчит эмбеддинги (EMBED_BATCH_SIZE) вместо диспатча по одному
  документу: worker-gpu кодирует пачку разом и реже переписывает FAISS-индекс.

Покрытие ~33% (прямые OA-PDF: arxiv/usenix/springer/techscience и т.п.);
для остального остаётся фолбэк на аннотацию. Управляется FETCH_FULL_TEXT.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-23 19:42:02 +05:00
jze9
38b000703a fix(gpu): рабочий FAISS-индекс вместо необучаемого IVFFlat
Семантический поиск (уровень 3) не работал вообще: индекс IndexIVFFlat с
nlist=1024 требует ~40000 векторов для обучения, а до обучения векторы
складывались во временный in-memory _flat_index, который:
  - не сохранялся на диск (save() писал только пустой _index) → терялся при
    рестарте;
  - не участвовал в поиске (search() искал только в необученном _index и
    сразу возвращал []).
Итог: в FAISS всегда было 0 векторов, семантика возвращала пусто.

Заменено на IndexIDMap2(IndexFlatIP): без обучения, работает с первого
вектора, doc_id хранится внутри индекса, корректно персистится. На
нормализованных векторах inner product = cosine, порог 0.75 сохраняет смысл.
Добавлена идемпотентность add_vectors (remove_ids перед add).

Проверено: 66 документов → ntotal=66, поиск возвращает релевантные
результаты со score 0.72-0.78, round-trip save/load сохраняет векторы.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-23 19:28:34 +05:00
jze9
9ca2bd1742 fix(notifier): маппинг input_data в модели Task для уведомлений о завершении
send_task_done падал с AttributeError: минимальная модель Task в notifier
не мапила колонку input_data (она есть в БД и в полной модели API), а
_build_summary читал task.input_data для search/plagiarism. Задача уходила
в бесконечные ретраи. Добавлен маппинг колонки + защита от NULL.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-23 19:16:13 +05:00
jze9
c1cf1ddd2f feat(auth): подтверждение email + актуализация SMTP/Ollama конфигов
- Резенд письма верификации (/auth/resend-verification), модалка на
  фронте с поллингом статуса, страница /verify-email/:token
- SMTP переведён на собственный Postfix (mail.jze9mail.ru, STARTTLS,
  SMTP_TLS_VERIFY) вместо Yandex-заглушки в дефолтах и .env.example
- OLLAMA_URL и модель в worker-gpu синхронизированы с новым GPU-хостом
  (llama3:8b -> qwen2.5:7b, которой раньше не было на сервере)

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-23 18:37:22 +05:00
19 changed files with 642 additions and 187 deletions

View File

@@ -28,12 +28,13 @@ OLLAMA_URL=http://ollama:11434
SECRET_KEY=change-me-in-production-use-openssl-rand-hex-32
ACCESS_TOKEN_EXPIRE_MINUTES=10080
# SMTP (Yandex)
SMTP_HOST=smtp.yandex.ru
SMTP_PORT=465
SMTP_USER=noreply@jze9.ru
# SMTP (собственный Postfix+Dovecot, mail.jze9mail.ru, STARTTLS)
SMTP_HOST=mail.jze9mail.ru
SMTP_PORT=587
SMTP_USER=noreply
SMTP_PASSWORD=changeme
SMTP_FROM=noreply@jze9.ru
SMTP_FROM=noreply@jze9mail.ru
SMTP_TLS_VERIFY=true
# App
APP_URL=https://academic.jze9.ru

View File

@@ -192,9 +192,17 @@ class OpenAlexParser(BaseParser):
source = primary_location.get("source") or {}
journal = source.get("display_name")
# URL на полный текст
# URL на полный текст: предпочитаем прямую ссылку на PDF (её реально можно
# скачать и извлечь текст), иначе oa_url / лендинг журнала.
best_oa = raw.get("best_oa_location") or {}
oa = raw.get("open_access") or {}
url = oa.get("oa_url") or primary_location.get("landing_page_url")
url = (
best_oa.get("pdf_url")
or primary_location.get("pdf_url")
or oa.get("oa_url")
or best_oa.get("landing_page_url")
or primary_location.get("landing_page_url")
)
# Аннотация (восстановить из инвертированного индекса)
abstract = None

View File

@@ -98,6 +98,37 @@ async def get_me(current_user: User = Depends(get_current_user)) -> UserInToken:
return UserInToken.model_validate(current_user)
@router.post("/resend-verification", status_code=status.HTTP_204_NO_CONTENT)
async def resend_verification(
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
) -> None:
"""Повторно отправить письмо с подтверждением email."""
if current_user.is_verified:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Email уже подтверждён",
)
if not current_user.verification_token:
current_user.verification_token = secrets.token_urlsafe(32)
await db.commit()
await db.refresh(current_user)
try:
celery_app.send_task(
"notify.send_verification",
args=[current_user.email, current_user.name, current_user.verification_token],
queue="queue.notify",
)
except Exception as e:
logger.warning(f"Не удалось поставить задачу повторной верификации: {e}")
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Не удалось отправить письмо, попробуйте позже",
)
@router.post("/verify-email/{token}", status_code=status.HTTP_200_OK)
async def verify_email(token: str, db: AsyncSession = Depends(get_db)) -> dict:
"""Подтвердить email по токену из письма."""

View File

@@ -45,11 +45,12 @@ class Settings(BaseSettings):
ACCESS_TOKEN_EXPIRE_MINUTES: int = 10080 # 7 дней
# SMTP
SMTP_HOST: str = "smtp.yandex.ru"
SMTP_PORT: int = 465
SMTP_USER: str = "noreply@jze9.ru"
SMTP_HOST: str = "mail.jze9mail.ru"
SMTP_PORT: int = 587
SMTP_USER: str = "noreply"
SMTP_PASSWORD: str = "changeme"
SMTP_FROM: str = "noreply@jze9.ru"
SMTP_FROM: str = "noreply@jze9mail.ru"
SMTP_TLS_VERIFY: bool = True
# App
APP_URL: str = "https://academic.jze9.ru"

View File

@@ -40,6 +40,9 @@ export const authApi = {
verifyEmail: (token: string) =>
api.post(`/auth/verify-email/${token}`),
resendVerification: () =>
api.post('/auth/resend-verification'),
};
export const tasksApi = {

View File

@@ -0,0 +1,126 @@
import { useState } from 'react';
import { useLocation } from 'react-router-dom';
import { useMutation, useQuery } from '@tanstack/react-query';
import { Mail, X, CheckCircle2, AlertCircle } from 'lucide-react';
import toast from 'react-hot-toast';
import { authApi } from '../api/client';
import { useAuthStore } from '../store/auth';
export function EmailVerificationModal() {
const { user, updateUser } = useAuthStore();
const location = useLocation();
const [dismissed, setDismissed] = useState(false);
const [verified, setVerified] = useState(false);
// Не показывать на странице подтверждения — там своя UI
if (location.pathname.startsWith('/verify-email')) return null;
const resend = useMutation({
mutationFn: () => authApi.resendVerification(),
onSuccess: () => toast.success('Письмо отправлено, проверьте почту'),
onError: () => toast.error('Не удалось отправить письмо'),
});
useQuery({
queryKey: ['email-verification-poll'],
queryFn: async () => {
const res = await authApi.me();
if (res.data.is_verified) {
updateUser({ is_verified: true });
setVerified(true);
}
return res.data;
},
enabled: !!(user && !user.is_verified && !dismissed && !verified),
refetchInterval: 5000,
staleTime: 0,
});
if (!user || dismissed) return null;
if (!user.is_verified && !verified) {
return (
<div className="fixed inset-0 z-50 flex items-center justify-center bg-black/40 backdrop-blur-sm">
<div className="relative w-full max-w-md mx-4 bg-white rounded-2xl shadow-2xl p-8">
<button
onClick={() => setDismissed(true)}
className="absolute top-4 right-4 text-gray-400 hover:text-gray-600 transition-colors"
>
<X className="w-5 h-5" />
</button>
<div className="flex justify-center mb-5">
<div className="p-4 bg-amber-50 rounded-full">
<Mail className="w-8 h-8 text-amber-500" />
</div>
</div>
<h2 className="text-xl font-bold text-gray-900 text-center mb-1">
Подтвердите ваш email
</h2>
<p className="text-center text-sm text-gray-400 mb-6">{user.email}</p>
{/* Статус */}
<div className="flex items-center justify-center gap-2 mb-6 px-4 py-3 bg-amber-50 rounded-xl">
<AlertCircle className="w-4 h-4 text-amber-500 shrink-0" />
<span className="text-sm font-medium text-amber-700">Email не подтверждён</span>
</div>
<p className="text-sm text-gray-500 text-center mb-6 leading-relaxed">
Мы отправили письмо на{' '}
<span className="font-medium text-gray-800">{user.email}</span>.
Нажмите кнопку «Подтвердить» в письме для активации аккаунта.
</p>
<button
onClick={() => resend.mutate()}
disabled={resend.isPending}
className="w-full py-3 bg-brand-600 text-white rounded-xl text-sm font-medium hover:bg-brand-700 disabled:opacity-50 transition-colors"
>
{resend.isPending ? 'Отправляем...' : 'Подтвердить'}
</button>
<button
onClick={() => setDismissed(true)}
className="w-full mt-3 py-2.5 text-sm text-gray-400 hover:text-gray-600 transition-colors"
>
Напомнить позже
</button>
</div>
</div>
);
}
if (verified) {
return (
<div className="fixed inset-0 z-50 flex items-center justify-center bg-black/40 backdrop-blur-sm">
<div className="relative w-full max-w-md mx-4 bg-white rounded-2xl shadow-2xl p-8 text-center">
<div className="flex justify-center mb-5">
<div className="p-4 bg-green-50 rounded-full">
<CheckCircle2 className="w-8 h-8 text-green-500" />
</div>
</div>
<h2 className="text-xl font-bold text-gray-900 mb-2">Email подтверждён!</h2>
<div className="flex items-center justify-center gap-2 mb-6 px-4 py-3 bg-green-50 rounded-xl">
<CheckCircle2 className="w-4 h-4 text-green-500 shrink-0" />
<span className="text-sm font-medium text-green-700">Email подтверждён</span>
</div>
<p className="text-sm text-gray-500 mb-6">
Аккаунт активирован. Теперь вы получаете уведомления о завершении задач.
</p>
<button
onClick={() => setDismissed(true)}
className="w-full py-3 bg-brand-600 text-white rounded-xl text-sm font-medium hover:bg-brand-700 transition-colors"
>
Закрыть
</button>
</div>
</div>
);
}
return null;
}

View File

@@ -3,6 +3,7 @@ import { Link, NavLink, useNavigate } from 'react-router-dom';
import { GraduationCap, Search, BookOpen, Upload, LayoutDashboard, LogOut, User, BookMarked } from 'lucide-react';
import { clsx } from 'clsx';
import { useAuthStore } from '../store/auth';
import { EmailVerificationModal } from './EmailVerificationModal';
interface LayoutProps {
children: React.ReactNode;
@@ -144,6 +145,8 @@ export function Layout({ children }: LayoutProps) {
<main className="max-w-6xl mx-auto px-4 sm:px-6 lg:px-8 py-8">
{children}
</main>
<EmailVerificationModal />
</div>
);
}

View File

@@ -14,6 +14,7 @@ import { Task } from './pages/Task';
import { Pricing } from './pages/Pricing';
import { Login } from './pages/Login';
import { Register } from './pages/Register';
import { VerifyEmail } from './pages/VerifyEmail';
import { AdminLayout } from './pages/admin/AdminLayout';
import { Dashboard } from './pages/admin/Dashboard';
import { Users as AdminUsers } from './pages/admin/Users';
@@ -48,6 +49,7 @@ function PublicApp() {
<Route path="/pricing" element={<Pricing />} />
<Route path="/login" element={<Login />} />
<Route path="/register" element={<Register />} />
<Route path="/verify-email/:token" element={<VerifyEmail />} />
</Routes>
</Layout>
);

View File

@@ -0,0 +1,83 @@
import { useEffect } from 'react';
import { useParams, useNavigate, Link } from 'react-router-dom';
import { useMutation } from '@tanstack/react-query';
import { CheckCircle2, XCircle, Loader2 } from 'lucide-react';
import { authApi } from '../api/client';
import { useAuthStore } from '../store/auth';
export function VerifyEmail() {
const { token } = useParams<{ token: string }>();
const navigate = useNavigate();
const { updateUser } = useAuthStore();
const verify = useMutation({
mutationFn: () => authApi.verifyEmail(token!),
onSuccess: () => {
updateUser({ is_verified: true });
setTimeout(() => navigate('/cabinet'), 3000);
},
});
useEffect(() => {
if (token) verify.mutate();
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [token]);
return (
<div className="max-w-md mx-auto pt-16 text-center">
<div className="bg-white rounded-2xl border border-gray-100 p-10">
{verify.isPending && (
<>
<div className="flex justify-center mb-5">
<Loader2 className="w-14 h-14 text-brand-600 animate-spin" />
</div>
<h1 className="text-xl font-bold text-gray-900 mb-2">Подтверждаем email</h1>
<p className="text-sm text-gray-500">Пожалуйста, подождите</p>
</>
)}
{verify.isSuccess && (
<>
<div className="flex justify-center mb-5">
<div className="p-4 bg-green-50 rounded-full">
<CheckCircle2 className="w-14 h-14 text-green-500" />
</div>
</div>
<h1 className="text-2xl font-bold text-gray-900 mb-2">Email подтверждён!</h1>
<div className="inline-flex items-center gap-2 px-4 py-2 bg-green-50 rounded-xl mb-4">
<CheckCircle2 className="w-4 h-4 text-green-500" />
<span className="text-sm font-medium text-green-700">Email подтверждён</span>
</div>
<p className="text-sm text-gray-500">
Аккаунт активирован. Перенаправляем в личный кабинет
</p>
</>
)}
{verify.isError && (
<>
<div className="flex justify-center mb-5">
<div className="p-4 bg-red-50 rounded-full">
<XCircle className="w-14 h-14 text-red-500" />
</div>
</div>
<h1 className="text-2xl font-bold text-gray-900 mb-2">Ошибка подтверждения</h1>
<div className="inline-flex items-center gap-2 px-4 py-2 bg-red-50 rounded-xl mb-4">
<XCircle className="w-4 h-4 text-red-500" />
<span className="text-sm font-medium text-red-700">Email не подтверждён</span>
</div>
<p className="text-sm text-gray-500 mb-6">
Ссылка недействительна или устарела. Запросите новое письмо в личном кабинете.
</p>
<Link
to="/cabinet"
className="inline-block px-6 py-2.5 bg-brand-600 text-white rounded-xl text-sm font-medium hover:bg-brand-700 transition-colors"
>
В личный кабинет
</Link>
</>
)}
</div>
</div>
);
}

View File

@@ -1,14 +1,20 @@
"""Singleton менеджер FAISS GPU индекса.
"""Singleton менеджер FAISS индекса.
Использует IVFFlat (а не HNSW — не поддерживается на GPU).
Поддерживает graceful degradation на CPU если GPU недоступна.
Используется плоский индекс IndexFlatIP, обёрнутый в IndexIDMap2, что даёт:
- отсутствие этапа обучения (в отличие от IVFFlat) — индекс работоспособен сразу,
начиная с первого вектора;
- хранение doc_id прямо внутри индекса (add_with_ids) — не нужен отдельный
маппинг на диске, а поиск сразу возвращает doc_id из PostgreSQL;
- корректную персистентность через faiss.write_index / read_index.
Векторы модели нормализованы (normalize_embeddings=True), поэтому inner product
эквивалентен косинусной близости. Для масштаба проекта (сотни тысяч — единицы
миллионов документов) полный перебор по FlatIP по скорости приемлем.
"""
import json
import logging
import os
from pathlib import Path
from typing import Optional
import numpy as np
@@ -18,70 +24,83 @@ logger = logging.getLogger(__name__)
class FAISSManager:
"""Singleton для управления FAISS индексом на GPU."""
"""Singleton для управления FAISS индексом (IndexIDMap2 поверх IndexFlatIP)."""
_index = None
_id_map: dict[int, int] = {} # faiss_internal_id -> doc_id (PostgreSQL)
_reverse_map: dict[int, int] = {} # doc_id -> faiss_internal_id
# doc_id -> faiss_id. Для IDMap2 faiss_id == doc_id, но маппинг сохраняем
# для совместимости с вызывающим кодом (plagiarism.embed_documents).
_reverse_map: dict[int, int] = {}
_use_gpu: bool = False
_is_trained: bool = False
@classmethod
def _new_index(cls):
"""Создать новый пустой индекс нужного типа."""
import faiss
base = faiss.IndexFlatIP(settings.EMBED_DIM)
return faiss.IndexIDMap2(base)
@classmethod
def load_or_create(cls) -> None:
"""
Загрузить индекс с диска или создать новый.
"""Загрузить индекс с диска или создать новый.
Пытается перенести индекс на GPU, при ошибке остаётся на CPU.
Старый несовместимый индекс (например, IVFFlat от предыдущей версии,
который не хранит id_map) безопасно пересоздаётся — полезных векторов в
нём всё равно не было.
"""
import faiss
index_path = settings.FAISS_INDEX_PATH
id_map_path = settings.FAISS_ID_MAP_PATH
# Загрузить ID маппинг
if os.path.exists(id_map_path):
with open(id_map_path, "r") as f:
raw_map = json.load(f)
cls._id_map = {int(k): int(v) for k, v in raw_map.items()}
cls._reverse_map = {v: k for k, v in cls._id_map.items()}
logger.info(f"ID маппинг загружен: {len(cls._id_map)} записей")
if os.path.exists(index_path):
# Загрузить существующий индекс
logger.info(f"Загрузка FAISS индекса из {index_path}")
cpu_index = faiss.read_index(index_path)
cls._is_trained = cpu_index.is_trained
else:
# Создать новый IVFFlat индекс
logger.info("Создание нового FAISS IVFFlat индекса...")
quantizer = faiss.IndexFlatIP(settings.EMBED_DIM)
cpu_index = faiss.IndexIVFFlat(
quantizer,
settings.EMBED_DIM,
settings.FAISS_NLIST,
faiss.METRIC_INNER_PRODUCT,
)
# IVFFlat требует обучения перед использованием
cls._is_trained = False
# Попытка перенести на GPU
try:
res = faiss.StandardGpuResources()
cls._index = faiss.index_cpu_to_gpu(res, 0, cpu_index)
cls._use_gpu = True
logger.info("FAISS индекс размещён на GPU")
loaded = faiss.read_index(index_path)
if hasattr(loaded, "id_map"):
cls._index = loaded
cls._rebuild_reverse_map()
logger.info(
f"Загрузка FAISS индекса из {index_path} "
f"({cls._index.ntotal} векторов)"
)
else:
logger.warning(
"На диске несовместимый FAISS индекс (без id_map) — "
"пересоздаём как IndexIDMap2(IndexFlatIP)"
)
cls._index = cls._new_index()
cls._reverse_map = {}
except Exception as e:
logger.warning(f"GPU недоступна: {e}. Используем CPU FAISS.")
cls._index = cpu_index
cls._use_gpu = False
logger.warning(f"Не удалось загрузить FAISS индекс ({e}) — создаём новый")
cls._index = cls._new_index()
cls._reverse_map = {}
else:
logger.info("Создание нового FAISS индекса IndexIDMap2(IndexFlatIP)...")
cls._index = cls._new_index()
cls._reverse_map = {}
if cls._is_trained:
cls._index.nprobe = settings.FAISS_NPROBE
cls._use_gpu = False # FlatIP на CPU достаточно быстр для целевого масштаба
@classmethod
def _ensure(cls) -> None:
"""Ленивая инициализация индекса при первом обращении."""
if cls._index is None:
cls.load_or_create()
@classmethod
def _rebuild_reverse_map(cls) -> None:
"""Восстановить _reverse_map из id, хранящихся внутри загруженного индекса."""
import faiss
try:
ids = faiss.vector_to_array(cls._index.id_map)
cls._reverse_map = {int(i): int(i) for i in ids}
except Exception as e:
logger.warning(f"Не удалось восстановить reverse_map из индекса: {e}")
cls._reverse_map = {}
@classmethod
def search(cls, query_vector: np.ndarray, k: int = 20) -> list[tuple[int, float]]:
"""
Поиск k ближайших векторов.
"""Поиск k ближайших векторов.
Args:
query_vector: Нормализованный вектор запроса, форма (768,)
@@ -90,21 +109,21 @@ class FAISSManager:
Returns:
Список кортежей (doc_id, cosine_score), отсортированных по убыванию score
"""
if cls._index is None or not cls._is_trained:
logger.warning("FAISS индекс не инициализирован или не обучен, пропускаем поиск")
cls._ensure()
if cls._index.ntotal == 0:
return []
try:
query = query_vector.reshape(1, -1).astype(np.float32)
distances, indices = cls._index.search(query, k)
distances, ids = cls._index.search(query, min(k, cls._index.ntotal))
results = []
for idx, dist in zip(indices[0], distances[0]):
for idx, dist in zip(ids[0], distances[0]):
if idx == -1:
continue
doc_id = cls._id_map.get(int(idx))
if doc_id is not None:
results.append((doc_id, float(dist)))
# Для IDMap2 idx — это уже doc_id из PostgreSQL
results.append((int(idx), float(dist)))
return results
@@ -114,95 +133,49 @@ class FAISSManager:
@classmethod
def add_vectors(cls, vectors: np.ndarray, doc_ids: list[int]) -> None:
"""
Добавить векторы в индекс.
"""Добавить (или обновить) векторы в индекс.
Если индекс не обучен и накопилось достаточно векторов — обучить его.
Идемпотентно по doc_id: при повторном эмбеддинге старый вектор документа
удаляется перед добавлением нового, чтобы не плодить дубли.
Args:
vectors: numpy массив формы (N, 768)
doc_ids: Список doc_id из PostgreSQL
"""
import faiss
cls._ensure()
if cls._index is None:
cls.load_or_create()
vectors = vectors.astype(np.float32)
if not cls._is_trained:
# Для IVFFlat нужно минимум nlist * 39 обучающих примеров
min_train = settings.FAISS_NLIST * 39
current_n = cls._index.ntotal if cls._index is not None else 0
if current_n + len(vectors) >= min_train:
logger.info(f"Обучение IVFFlat индекса на {current_n + len(vectors)} векторах...")
cls._index.train(vectors)
cls._is_trained = True
cls._index.nprobe = settings.FAISS_NPROBE
logger.info("Обучение завершено")
else:
logger.info(
f"Недостаточно векторов для обучения IVFFlat "
f"({current_n + len(vectors)} < {min_train}). "
"Используйте FlatIP до накопления достаточного количества документов."
)
# Временный flat индекс для малого количества документов
if not hasattr(cls, '_flat_index') or cls._flat_index is None:
cls._flat_index = faiss.IndexFlatIP(settings.EMBED_DIM)
# Добавить в flat индекс
start_id = cls._flat_index.ntotal
cls._flat_index.add(vectors)
for i, doc_id in enumerate(doc_ids):
internal_id = start_id + i
cls._id_map[internal_id] = doc_id
cls._reverse_map[doc_id] = internal_id
cls._save_id_map()
if len(doc_ids) == 0:
return
if cls._is_trained:
start_id = cls._index.ntotal
cls._index.add(vectors)
for i, doc_id in enumerate(doc_ids):
internal_id = start_id + i
cls._id_map[internal_id] = doc_id
cls._reverse_map[doc_id] = internal_id
vectors = np.asarray(vectors, dtype=np.float32)
ids = np.asarray(doc_ids, dtype=np.int64)
# Удалить существующие id, чтобы повторный эмбеддинг не создавал дубли
try:
cls._index.remove_ids(ids)
except Exception:
pass
cls._index.add_with_ids(vectors, ids)
for doc_id in doc_ids:
cls._reverse_map[int(doc_id)] = int(doc_id)
cls._save_id_map()
logger.info(f"Добавлено {len(doc_ids)} векторов в FAISS. Всего: {cls._index.ntotal}")
@classmethod
def save(cls) -> None:
"""Сохранить индекс на диск (CPU версия)."""
if cls._index is None:
return
"""Сохранить индекс на диск."""
cls._ensure()
import faiss
index_path = Path(settings.FAISS_INDEX_PATH)
index_path.parent.mkdir(parents=True, exist_ok=True)
# Перенести на CPU перед сохранением
if cls._use_gpu:
cpu_index = faiss.index_gpu_to_cpu(cls._index)
else:
cpu_index = cls._index
faiss.write_index(cpu_index, str(index_path))
cls._save_id_map()
logger.info(f"FAISS индекс сохранён: {index_path} ({cpu_index.ntotal} векторов)")
faiss.write_index(cls._index, str(index_path))
logger.info(f"FAISS индекс сохранён: {index_path} ({cls._index.ntotal} векторов)")
@classmethod
def _save_id_map(cls) -> None:
"""Сохранить маппинг faiss_id -> doc_id на диск."""
id_map_path = Path(settings.FAISS_ID_MAP_PATH)
id_map_path.parent.mkdir(parents=True, exist_ok=True)
with open(id_map_path, "w") as f:
json.dump({str(k): v for k, v in cls._id_map.items()}, f)
@classmethod
@property
def total_vectors(cls) -> int:
"""Количество векторов в индексе."""
if cls._index is None:

View File

@@ -15,7 +15,7 @@ class OllamaClient:
def __init__(self) -> None:
self.base_url = settings.OLLAMA_URL
self.model = "llama3:8b"
self.model = "qwen2.5:7b"
self.timeout = 60.0 # секунд
def check_paraphrase(self, text_a: str, text_b: str) -> dict:

View File

@@ -152,10 +152,14 @@ def check_plagiarism(
seen_sources.add(key)
unique_matches.append(m)
# Вычислить общий процент схожести
# Вычислить общий процент схожести по ДОЛЕ помеченных фрагментов документа.
# Считаем уникальные позиции (фрагмент, совпавший с несколькими источниками,
# не должен раздувать процент выше 100).
total_frags = len(fragments)
flagged_frags = len(unique_matches)
flagged_positions = {m.get("position_start") for m in unique_matches}
flagged_frags = len(flagged_positions)
overall_similarity = (flagged_frags / total_frags * 100) if total_frags > 0 else 0.0
overall_similarity = min(overall_similarity, 100.0)
result = {
"overall_similarity": round(overall_similarity, 2),

View File

@@ -39,7 +39,22 @@ class Settings(BaseSettings):
# Настройки обработки текста
FRAGMENT_WINDOW_WORDS: int = 200 # Размер окна для фрагментов
FRAGMENT_OVERLAP_WORDS: int = 50 # Перекрытие фрагментов
MAX_FINGERPRINTS_PER_DOC: int = 500 # Максимум хэшей Winnowing на документ
# Максимум хэшей Winnowing на документ. Winnowing рассчитан на ПОЛНОЕ
# хранение отпечатков: обрезка теряет гарантию, что скопированный фрагмент
# разделит отпечатки с источником (при 500 хранилось ~14% отпечатков статьи,
# и локализованный поиск фрагментов почти не срабатывал). Держим высоким;
# цена — размер таблицы fingerprints (~3-4К строк на статью).
MAX_FINGERPRINTS_PER_DOC: int = 20000
# Порог уровня 1: % отпечатков фрагмента, найденных в источнике, чтобы
# пометить фрагмент как скопированный
EXACT_FRAGMENT_THRESHOLD: float = 40.0
# Скачивание полного текста статей (PDF по URL источника)
FETCH_FULL_TEXT: bool = False # Включить обогащение полным текстом при заливке корпуса
FULL_TEXT_TIMEOUT: float = 30.0 # Таймаут скачивания одного документа, сек
FULL_TEXT_MAX_BYTES: int = 30 * 1024 * 1024 # Лимит размера скачиваемого файла (30 МБ)
FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным
EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов
# App
ENVIRONMENT: str = "development"

View File

@@ -0,0 +1,76 @@
"""Скачивание и извлечение полного текста статьи по URL источника.
OpenAlex/arXiv отдают в метаданных ссылку (`url` = oa_url / pdf), которая часто
ведёт на PDF открытого доступа. Здесь мы best-effort скачиваем файл и извлекаем
из него текст. HTML-страницы (лендинги журналов) пропускаем — надёжно доставать
текст статьи из произвольного HTML нельзя.
Всё завёрнуто в широкий except: недоступный/битый источник не должен ронять
заливку корпуса, просто у документа не будет полного текста.
"""
import logging
import httpx
from app.config import settings
from app.extractors.pdf import extract_text_from_pdf
logger = logging.getLogger(__name__)
_HEADERS = {
"User-Agent": "AcademicHelper/1.0 (+https://academic.jze9.ru; mailto:noreply@jze9mail.ru)",
"Accept": "application/pdf,*/*",
}
def fetch_full_text(url: str) -> str | None:
"""Скачать документ по URL и вернуть извлечённый текст, либо None.
Возвращает None, если: url пустой, файл не PDF, скачивание не удалось,
или извлечённого текста слишком мало (< FULL_TEXT_MIN_CHARS).
"""
if not url:
return None
try:
with httpx.Client(
follow_redirects=True,
timeout=settings.FULL_TEXT_TIMEOUT,
headers=_HEADERS,
) as client:
with client.stream("GET", url) as resp:
resp.raise_for_status()
ctype = resp.headers.get("content-type", "").lower()
# Скачиваем с ограничением размера
buf = bytearray()
for chunk in resp.iter_bytes():
buf += chunk
if len(buf) > settings.FULL_TEXT_MAX_BYTES:
logger.info(
f"full-text превысил лимит {settings.FULL_TEXT_MAX_BYTES} байт, "
f"обрезаю: {url}"
)
break
data = bytes(buf)
except Exception as e:
logger.info(f"full-text: скачать не удалось {url!r}: {type(e).__name__}: {str(e)[:120]}")
return None
# Определяем PDF по content-type или magic-байтам
is_pdf = "pdf" in ctype or data[:5] == b"%PDF-"
if not is_pdf:
logger.debug(f"full-text: не PDF (content-type={ctype!r}), пропускаю: {url}")
return None
try:
text = (extract_text_from_pdf(data) or "").strip()
except Exception as e:
logger.info(f"full-text: извлечение PDF не удалось {url!r}: {type(e).__name__}: {str(e)[:120]}")
return None
if len(text) < settings.FULL_TEXT_MIN_CHARS:
return None
return text

View File

@@ -139,50 +139,64 @@ def extract_and_check(
)
logger.info(f"Фрагментов создано: {len(fragments)}")
# ──── Уровень 1: Winnowing fingerprints ────────────────────────────────
# ──── Уровень 1: Winnowing fingerprints (пофрагментно, с локализацией) ───
# Winnow'им КАЖДЫЙ фрагмент отдельно и ищем, с каким источником и на
# какой позиции он совпадает. Так в отчёте видно не «документ похож на X»,
# а «фрагмент на символах AB скопирован из источника X».
level1_matches: list[dict] = []
doc_fingerprint = winnow(text)
if doc_fingerprint:
from app.models import Document, Fingerprint
with db_session() as session:
hashes = list(doc_fingerprint)[: settings.MAX_FINGERPRINTS_PER_DOC]
doc_cache: dict[int, Document] = {}
# Найти совпадения в базе fingerprints
matching_docs = session.execute(
for frag in fragments:
frag_fp = winnow(frag["text"])
if not frag_fp:
continue
frag_hashes = list(frag_fp)
# Источник, разделяющий больше всего отпечатков с этим фрагментом
row = session.execute(
select(
Fingerprint.doc_id,
func.count(Fingerprint.id).label("match_count"),
func.count(Fingerprint.id).label("cnt"),
)
.where(Fingerprint.hash_value.in_(hashes))
.where(Fingerprint.hash_value.in_(frag_hashes))
.group_by(Fingerprint.doc_id)
.having(func.count(Fingerprint.id) > len(hashes) * 0.1)
.order_by(func.count(Fingerprint.id).desc())
.limit(20)
).all()
.limit(1)
).first()
for doc_id, match_count in matching_docs:
similarity = match_count / len(hashes) * 100
if similarity < 20:
if row is None:
continue
doc_id, cnt = row
# Доля отпечатков фрагмента, найденных в источнике
similarity = cnt / len(frag_hashes) * 100
if similarity < settings.EXACT_FRAGMENT_THRESHOLD:
continue
doc = doc_cache.get(doc_id)
if doc is None:
doc = session.get(Document, doc_id)
if not doc:
if doc is None:
continue
doc_cache[doc_id] = doc
level1_matches.append({
"fragment": text[:200],
"position_start": 0,
"position_end": len(text),
"similarity": round(similarity, 1),
"fragment": frag["text"][:300],
"position_start": frag["start"],
"position_end": frag["end"],
"similarity": round(min(similarity, 100.0), 1),
"method": "exact",
"source_title": doc.title,
"source_url": doc.url,
"source_db": doc.source,
})
logger.info(f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений")
logger.info(
f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений-фрагментов"
)
# ──── Уровень 2: MinHash LSH ────────────────────────────────────────────
level2_matches: list[dict] = []
@@ -242,20 +256,24 @@ def extract_and_check(
@celery_app.task(name="index.add_document")
def add_document(doc_data: dict[str, Any]) -> dict[str, Any]:
def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[str, Any]:
"""
Добавить документ из внешнего источника в систему.
Алгоритм:
1. Дедупликация по ext_id
2. Сохранить метаданные в PostgreSQL
3. Индексировать в Elasticsearch
4. Вычислить Winnowing fingerprints
5. Добавить в MinHash LSH
6. Диспатч gpu.embed_documents для FAISS эмбеддингов
3. Вычислить провизорные Winnowing fingerprints (из аннотации)
4. Добавить в MinHash LSH
5. Индексировать в Elasticsearch
6. Диспатч gpu.embed_documents для FAISS эмбеддингов (если dispatch_embed)
7. Если включён FETCH_FULL_TEXT и есть url — диспатч index.enrich_full_text,
который скачает полный текст и пересчитает fingerprints по нему.
Args:
doc_data: Словарь с метаданными документа
dispatch_embed: Диспатчить ли эмбеддинг по одному документу. При массовой
заливке run_parser выключает это и батчит эмбеддинги сам.
Returns:
dict со статусом операции и doc_id
@@ -325,16 +343,91 @@ def add_document(doc_data: dict[str, Any]) -> dict[str, Any]:
except Exception as e:
logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}")
# Диспатч FAISS эмбеддингов
# Диспатч FAISS эмбеддингов (по одному документу; при массовой заливке
# run_parser выключает это и батчит сам)
if dispatch_embed:
celery_app.send_task(
"gpu.embed_documents",
args=[[doc_id]],
queue="queue.gpu",
)
# Обогащение полным текстом: скачать PDF по url и пересчитать fingerprints
if settings.FETCH_FULL_TEXT and doc_data.get("url"):
celery_app.send_task(
"index.enrich_full_text",
args=[doc_id, doc_data["url"]],
queue="queue.index",
)
return {"status": "indexed", "doc_id": doc_id}
@celery_app.task(
name="index.enrich_full_text",
bind=True,
max_retries=2,
default_retry_delay=120,
)
def enrich_full_text(self, doc_id: int, url: str) -> dict[str, Any]:
"""Скачать полный текст статьи и пересчитать по нему fingerprints.
Метаданные и провизорные fingerprints (из аннотации) уже сохранены в
add_document. Здесь мы:
1. Скачиваем PDF по url и извлекаем текст (best-effort).
2. Сохраняем полный текст в MinIO (bucket documents, префикс corpus/).
3. Заменяем fingerprints документа на посчитанные по полному тексту.
4. Обновляем MinHash LSH.
Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext.
"""
from app.fulltext import fetch_full_text
from app.models import Document, Fingerprint
text = fetch_full_text(url)
if not text:
return {"status": "no_fulltext", "doc_id": doc_id}
# Сохранить полный текст в MinIO
try:
minio = get_minio()
key = f"corpus/{doc_id}.txt"
data = text.encode("utf-8")
minio.put_object(
settings.MINIO_BUCKET_DOCS, key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8",
)
except Exception as exc:
logger.error(f"enrich_full_text: не удалось сохранить текст в MinIO для {doc_id}: {exc}")
raise self.retry(exc=exc, countdown=120)
# Пересчитать fingerprints по полному тексту
doc_fp = winnow(text)
hashes = list(doc_fp)[: settings.MAX_FINGERPRINTS_PER_DOC]
from sqlalchemy import delete
with db_session() as session:
doc = session.get(Document, doc_id)
if doc is None:
return {"status": "doc_gone", "doc_id": doc_id}
doc.minio_key = key
# Удалить провизорные fingerprints и записать новые
session.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id))
for i, hash_val in enumerate(hashes):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
session.commit()
# Обновить MinHash LSH по полному тексту
add_to_lsh(f"doc:{doc_id}", text)
logger.info(
f"enrich_full_text: doc={doc_id} полный текст {len(text)} симв., "
f"fingerprints={len(hashes)}"
)
return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)}
def _stage_work(
task_id: str,
minio_key: str,
@@ -451,16 +544,35 @@ def run_parser(source_id: int) -> dict[str, Any]:
parser = P()
# fetch+transform без записи в JSONL — работаем in-memory
raw_docs = parser.fetch(**fetch_kwargs)
# Эмбеддинги диспатчим пачками, а не по одному документу: так 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:
try:
doc = parser.transform(raw)
if not (doc and doc.get("title") and doc.get("ext_id")):
continue
result = add_document(doc)
result = add_document(doc, dispatch_embed=False)
if result.get("status") == "indexed":
added += 1
embed_batch.append(result["doc_id"])
if len(embed_batch) >= batch_size:
_flush_embed()
except Exception as e:
logger.warning(f"run_parser: ошибка документа: {e}")
_flush_embed()
except Exception as e:
error_msg = str(e)[:500]
logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True)

View File

@@ -24,11 +24,12 @@ class Settings(BaseSettings):
RABBITMQ_URL: str = "amqp://guest:guest@rabbitmq:5672/"
# SMTP
SMTP_HOST: str = "smtp.yandex.ru"
SMTP_PORT: int = 465
SMTP_USER: str = "noreply@jze9.ru"
SMTP_HOST: str = "mail.jze9mail.ru"
SMTP_PORT: int = 587
SMTP_USER: str = "noreply"
SMTP_PASSWORD: str = "changeme"
SMTP_FROM: str = "noreply@jze9.ru"
SMTP_FROM: str = "noreply@jze9mail.ru"
SMTP_TLS_VERIFY: bool = True
# App
APP_URL: str = "https://academic.jze9.ru"

View File

@@ -2,6 +2,7 @@
import logging
import smtplib
import ssl
from email.mime.multipart import MIMEMultipart
from email.mime.text import MIMEText
@@ -37,7 +38,20 @@ class EmailSender:
msg.attach(MIMEText(html_body, "html", "utf-8"))
try:
with smtplib.SMTP_SSL(settings.SMTP_HOST, settings.SMTP_PORT) as smtp:
tls_ctx = ssl.create_default_context()
if not settings.SMTP_TLS_VERIFY:
tls_ctx.check_hostname = False
tls_ctx.verify_mode = ssl.CERT_NONE
if settings.SMTP_PORT == 465:
with smtplib.SMTP_SSL(settings.SMTP_HOST, settings.SMTP_PORT, context=tls_ctx) as smtp:
smtp.login(settings.SMTP_USER, settings.SMTP_PASSWORD)
smtp.send_message(msg)
else:
with smtplib.SMTP(settings.SMTP_HOST, settings.SMTP_PORT) as smtp:
smtp.ehlo()
smtp.starttls(context=tls_ctx)
smtp.ehlo()
smtp.login(settings.SMTP_USER, settings.SMTP_PASSWORD)
smtp.send_message(msg)
logger.info(f"Email отправлен: {to_email!r}, тема: {subject!r}")

View File

@@ -25,6 +25,7 @@ class Task(Base):
user_id: Mapped[int] = mapped_column(ForeignKey("users.id"))
type: Mapped[str] = mapped_column(String(20))
status: Mapped[str] = mapped_column(String(20))
input_data: Mapped[dict | None] = mapped_column(JSON, nullable=True)
result: Mapped[dict | None] = mapped_column(JSON, nullable=True)
error: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(server_default=func.now())

View File

@@ -68,10 +68,11 @@ def send_task_done(self, task_id: str) -> dict:
def _build_summary(task) -> str:
"""Сформировать краткое описание результата для email."""
result = task.result or {}
input_data = task.input_data or {}
if task.type == "search":
count = len(result.get("sources", []))
query = task.input_data.get("query", "")
query = input_data.get("query", "")
return (
f'Найдено <strong>{count} источников</strong> по запросу "{query[:80]}".'
f' Откройте результат, чтобы просмотреть список с ГОСТ-цитатами.'
@@ -81,7 +82,7 @@ def _build_summary(task) -> str:
similarity = result.get("overall_similarity", 0)
total = result.get("total_fragments", 0)
flagged = result.get("flagged_fragments", 0)
filename = task.input_data.get("filename", "")
filename = input_data.get("filename", "")
color = "#dc2626" if similarity > 30 else "#d97706" if similarity > 10 else "#16a34a"
label = "Высокий" if similarity > 30 else "Средний" if similarity > 10 else "Низкий"