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

View File

@@ -192,9 +192,17 @@ class OpenAlexParser(BaseParser):
source = primary_location.get("source") or {} source = primary_location.get("source") or {}
journal = source.get("display_name") journal = source.get("display_name")
# URL на полный текст # URL на полный текст: предпочитаем прямую ссылку на PDF (её реально можно
# скачать и извлечь текст), иначе oa_url / лендинг журнала.
best_oa = raw.get("best_oa_location") or {}
oa = raw.get("open_access") 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 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) 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) @router.post("/verify-email/{token}", status_code=status.HTTP_200_OK)
async def verify_email(token: str, db: AsyncSession = Depends(get_db)) -> dict: async def verify_email(token: str, db: AsyncSession = Depends(get_db)) -> dict:
"""Подтвердить email по токену из письма.""" """Подтвердить email по токену из письма."""

View File

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

View File

@@ -40,6 +40,9 @@ export const authApi = {
verifyEmail: (token: string) => verifyEmail: (token: string) =>
api.post(`/auth/verify-email/${token}`), api.post(`/auth/verify-email/${token}`),
resendVerification: () =>
api.post('/auth/resend-verification'),
}; };
export const tasksApi = { 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 { GraduationCap, Search, BookOpen, Upload, LayoutDashboard, LogOut, User, BookMarked } from 'lucide-react';
import { clsx } from 'clsx'; import { clsx } from 'clsx';
import { useAuthStore } from '../store/auth'; import { useAuthStore } from '../store/auth';
import { EmailVerificationModal } from './EmailVerificationModal';
interface LayoutProps { interface LayoutProps {
children: React.ReactNode; 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"> <main className="max-w-6xl mx-auto px-4 sm:px-6 lg:px-8 py-8">
{children} {children}
</main> </main>
<EmailVerificationModal />
</div> </div>
); );
} }

View File

@@ -14,6 +14,7 @@ import { Task } from './pages/Task';
import { Pricing } from './pages/Pricing'; import { Pricing } from './pages/Pricing';
import { Login } from './pages/Login'; import { Login } from './pages/Login';
import { Register } from './pages/Register'; import { Register } from './pages/Register';
import { VerifyEmail } from './pages/VerifyEmail';
import { AdminLayout } from './pages/admin/AdminLayout'; import { AdminLayout } from './pages/admin/AdminLayout';
import { Dashboard } from './pages/admin/Dashboard'; import { Dashboard } from './pages/admin/Dashboard';
import { Users as AdminUsers } from './pages/admin/Users'; import { Users as AdminUsers } from './pages/admin/Users';
@@ -48,6 +49,7 @@ function PublicApp() {
<Route path="/pricing" element={<Pricing />} /> <Route path="/pricing" element={<Pricing />} />
<Route path="/login" element={<Login />} /> <Route path="/login" element={<Login />} />
<Route path="/register" element={<Register />} /> <Route path="/register" element={<Register />} />
<Route path="/verify-email/:token" element={<VerifyEmail />} />
</Routes> </Routes>
</Layout> </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). Используется плоский индекс IndexFlatIP, обёрнутый в IndexIDMap2, что даёт:
Поддерживает graceful degradation на CPU если GPU недоступна. - отсутствие этапа обучения (в отличие от 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 logging
import os import os
from pathlib import Path from pathlib import Path
from typing import Optional
import numpy as np import numpy as np
@@ -18,70 +24,83 @@ logger = logging.getLogger(__name__)
class FAISSManager: class FAISSManager:
"""Singleton для управления FAISS индексом на GPU.""" """Singleton для управления FAISS индексом (IndexIDMap2 поверх IndexFlatIP)."""
_index = None _index = None
_id_map: dict[int, int] = {} # faiss_internal_id -> doc_id (PostgreSQL) # doc_id -> faiss_id. Для IDMap2 faiss_id == doc_id, но маппинг сохраняем
_reverse_map: dict[int, int] = {} # doc_id -> faiss_internal_id # для совместимости с вызывающим кодом (plagiarism.embed_documents).
_reverse_map: dict[int, int] = {}
_use_gpu: bool = False _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 @classmethod
def load_or_create(cls) -> None: def load_or_create(cls) -> None:
""" """Загрузить индекс с диска или создать новый.
Загрузить индекс с диска или создать новый.
Пытается перенести индекс на GPU, при ошибке остаётся на CPU. Старый несовместимый индекс (например, IVFFlat от предыдущей версии,
который не хранит id_map) безопасно пересоздаётся — полезных векторов в
нём всё равно не было.
""" """
import faiss import faiss
index_path = settings.FAISS_INDEX_PATH 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): 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: try:
res = faiss.StandardGpuResources() loaded = faiss.read_index(index_path)
cls._index = faiss.index_cpu_to_gpu(res, 0, cpu_index) if hasattr(loaded, "id_map"):
cls._use_gpu = True cls._index = loaded
logger.info("FAISS индекс размещён на GPU") 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: except Exception as e:
logger.warning(f"GPU недоступна: {e}. Используем CPU FAISS.") logger.warning(f"Не удалось загрузить FAISS индекс ({e}) — создаём новый")
cls._index = cpu_index cls._index = cls._new_index()
cls._use_gpu = False cls._reverse_map = {}
else:
logger.info("Создание нового FAISS индекса IndexIDMap2(IndexFlatIP)...")
cls._index = cls._new_index()
cls._reverse_map = {}
if cls._is_trained: cls._use_gpu = False # FlatIP на CPU достаточно быстр для целевого масштаба
cls._index.nprobe = settings.FAISS_NPROBE
@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 @classmethod
def search(cls, query_vector: np.ndarray, k: int = 20) -> list[tuple[int, float]]: def search(cls, query_vector: np.ndarray, k: int = 20) -> list[tuple[int, float]]:
""" """Поиск k ближайших векторов.
Поиск k ближайших векторов.
Args: Args:
query_vector: Нормализованный вектор запроса, форма (768,) query_vector: Нормализованный вектор запроса, форма (768,)
@@ -90,21 +109,21 @@ class FAISSManager:
Returns: Returns:
Список кортежей (doc_id, cosine_score), отсортированных по убыванию score Список кортежей (doc_id, cosine_score), отсортированных по убыванию score
""" """
if cls._index is None or not cls._is_trained: cls._ensure()
logger.warning("FAISS индекс не инициализирован или не обучен, пропускаем поиск")
if cls._index.ntotal == 0:
return [] return []
try: try:
query = query_vector.reshape(1, -1).astype(np.float32) 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 = [] results = []
for idx, dist in zip(indices[0], distances[0]): for idx, dist in zip(ids[0], distances[0]):
if idx == -1: if idx == -1:
continue continue
doc_id = cls._id_map.get(int(idx)) # Для IDMap2 idx — это уже doc_id из PostgreSQL
if doc_id is not None: results.append((int(idx), float(dist)))
results.append((doc_id, float(dist)))
return results return results
@@ -114,95 +133,49 @@ class FAISSManager:
@classmethod @classmethod
def add_vectors(cls, vectors: np.ndarray, doc_ids: list[int]) -> None: def add_vectors(cls, vectors: np.ndarray, doc_ids: list[int]) -> None:
""" """Добавить (или обновить) векторы в индекс.
Добавить векторы в индекс.
Если индекс не обучен и накопилось достаточно векторов — обучить его. Идемпотентно по doc_id: при повторном эмбеддинге старый вектор документа
удаляется перед добавлением нового, чтобы не плодить дубли.
Args: Args:
vectors: numpy массив формы (N, 768) vectors: numpy массив формы (N, 768)
doc_ids: Список doc_id из PostgreSQL doc_ids: Список doc_id из PostgreSQL
""" """
import faiss cls._ensure()
if cls._index is None: if len(doc_ids) == 0:
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()
return return
if cls._is_trained: vectors = np.asarray(vectors, dtype=np.float32)
start_id = cls._index.ntotal ids = np.asarray(doc_ids, dtype=np.int64)
cls._index.add(vectors)
for i, doc_id in enumerate(doc_ids): # Удалить существующие id, чтобы повторный эмбеддинг не создавал дубли
internal_id = start_id + i try:
cls._id_map[internal_id] = doc_id cls._index.remove_ids(ids)
cls._reverse_map[doc_id] = internal_id 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}") logger.info(f"Добавлено {len(doc_ids)} векторов в FAISS. Всего: {cls._index.ntotal}")
@classmethod @classmethod
def save(cls) -> None: def save(cls) -> None:
"""Сохранить индекс на диск (CPU версия).""" """Сохранить индекс на диск."""
if cls._index is None: cls._ensure()
return
import faiss import faiss
index_path = Path(settings.FAISS_INDEX_PATH) index_path = Path(settings.FAISS_INDEX_PATH)
index_path.parent.mkdir(parents=True, exist_ok=True) index_path.parent.mkdir(parents=True, exist_ok=True)
# Перенести на CPU перед сохранением faiss.write_index(cls._index, str(index_path))
if cls._use_gpu: logger.info(f"FAISS индекс сохранён: {index_path} ({cls._index.ntotal} векторов)")
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} векторов)")
@classmethod @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: def total_vectors(cls) -> int:
"""Количество векторов в индексе.""" """Количество векторов в индексе."""
if cls._index is None: if cls._index is None:

View File

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

View File

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

View File

@@ -39,7 +39,22 @@ class Settings(BaseSettings):
# Настройки обработки текста # Настройки обработки текста
FRAGMENT_WINDOW_WORDS: int = 200 # Размер окна для фрагментов FRAGMENT_WINDOW_WORDS: int = 200 # Размер окна для фрагментов
FRAGMENT_OVERLAP_WORDS: int = 50 # Перекрытие фрагментов 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 # App
ENVIRONMENT: str = "development" 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)}") logger.info(f"Фрагментов создано: {len(fragments)}")
# ──── Уровень 1: Winnowing fingerprints ──────────────────────────────── # ──── Уровень 1: Winnowing fingerprints (пофрагментно, с локализацией) ───
# Winnow'им КАЖДЫЙ фрагмент отдельно и ищем, с каким источником и на
# какой позиции он совпадает. Так в отчёте видно не «документ похож на X»,
# а «фрагмент на символах AB скопирован из источника X».
level1_matches: list[dict] = [] level1_matches: list[dict] = []
doc_fingerprint = winnow(text)
if doc_fingerprint:
from app.models import Document, Fingerprint from app.models import Document, Fingerprint
with db_session() as session: with db_session() as session:
hashes = list(doc_fingerprint)[: settings.MAX_FINGERPRINTS_PER_DOC] doc_cache: dict[int, Document] = {}
# Найти совпадения в базе fingerprints for frag in fragments:
matching_docs = session.execute( frag_fp = winnow(frag["text"])
if not frag_fp:
continue
frag_hashes = list(frag_fp)
# Источник, разделяющий больше всего отпечатков с этим фрагментом
row = session.execute(
select( select(
Fingerprint.doc_id, 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) .group_by(Fingerprint.doc_id)
.having(func.count(Fingerprint.id) > len(hashes) * 0.1)
.order_by(func.count(Fingerprint.id).desc()) .order_by(func.count(Fingerprint.id).desc())
.limit(20) .limit(1)
).all() ).first()
for doc_id, match_count in matching_docs: if row is None:
similarity = match_count / len(hashes) * 100
if similarity < 20:
continue 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) doc = session.get(Document, doc_id)
if not doc: if doc is None:
continue continue
doc_cache[doc_id] = doc
level1_matches.append({ level1_matches.append({
"fragment": text[:200], "fragment": frag["text"][:300],
"position_start": 0, "position_start": frag["start"],
"position_end": len(text), "position_end": frag["end"],
"similarity": round(similarity, 1), "similarity": round(min(similarity, 100.0), 1),
"method": "exact", "method": "exact",
"source_title": doc.title, "source_title": doc.title,
"source_url": doc.url, "source_url": doc.url,
"source_db": doc.source, "source_db": doc.source,
}) })
logger.info(f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений") logger.info(
f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений-фрагментов"
)
# ──── Уровень 2: MinHash LSH ──────────────────────────────────────────── # ──── Уровень 2: MinHash LSH ────────────────────────────────────────────
level2_matches: list[dict] = [] level2_matches: list[dict] = []
@@ -242,20 +256,24 @@ def extract_and_check(
@celery_app.task(name="index.add_document") @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 1. Дедупликация по ext_id
2. Сохранить метаданные в PostgreSQL 2. Сохранить метаданные в PostgreSQL
3. Индексировать в Elasticsearch 3. Вычислить провизорные Winnowing fingerprints (из аннотации)
4. Вычислить Winnowing fingerprints 4. Добавить в MinHash LSH
5. Добавить в MinHash LSH 5. Индексировать в Elasticsearch
6. Диспатч gpu.embed_documents для FAISS эмбеддингов 6. Диспатч gpu.embed_documents для FAISS эмбеддингов (если dispatch_embed)
7. Если включён FETCH_FULL_TEXT и есть url — диспатч index.enrich_full_text,
который скачает полный текст и пересчитает fingerprints по нему.
Args: Args:
doc_data: Словарь с метаданными документа doc_data: Словарь с метаданными документа
dispatch_embed: Диспатчить ли эмбеддинг по одному документу. При массовой
заливке run_parser выключает это и батчит эмбеддинги сам.
Returns: Returns:
dict со статусом операции и doc_id dict со статусом операции и doc_id
@@ -325,16 +343,91 @@ def add_document(doc_data: dict[str, Any]) -> dict[str, Any]:
except Exception as e: except Exception as e:
logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}") logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}")
# Диспатч FAISS эмбеддингов # Диспатч FAISS эмбеддингов (по одному документу; при массовой заливке
# run_parser выключает это и батчит сам)
if dispatch_embed:
celery_app.send_task( celery_app.send_task(
"gpu.embed_documents", "gpu.embed_documents",
args=[[doc_id]], args=[[doc_id]],
queue="queue.gpu", 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} 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( def _stage_work(
task_id: str, task_id: str,
minio_key: str, minio_key: str,
@@ -451,16 +544,35 @@ def run_parser(source_id: int) -> dict[str, Any]:
parser = P() parser = P()
# fetch+transform без записи в JSONL — работаем in-memory # fetch+transform без записи в JSONL — работаем in-memory
raw_docs = parser.fetch(**fetch_kwargs) 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: for raw in raw_docs:
try: try:
doc = parser.transform(raw) doc = parser.transform(raw)
if not (doc and doc.get("title") and doc.get("ext_id")): if not (doc and doc.get("title") and doc.get("ext_id")):
continue continue
result = add_document(doc) result = add_document(doc, dispatch_embed=False)
if result.get("status") == "indexed": if result.get("status") == "indexed":
added += 1 added += 1
embed_batch.append(result["doc_id"])
if len(embed_batch) >= batch_size:
_flush_embed()
except Exception as e: except Exception as e:
logger.warning(f"run_parser: ошибка документа: {e}") logger.warning(f"run_parser: ошибка документа: {e}")
_flush_embed()
except Exception as e: except Exception as e:
error_msg = str(e)[:500] error_msg = str(e)[:500]
logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True) 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/" RABBITMQ_URL: str = "amqp://guest:guest@rabbitmq:5672/"
# SMTP # SMTP
SMTP_HOST: str = "smtp.yandex.ru" SMTP_HOST: str = "mail.jze9mail.ru"
SMTP_PORT: int = 465 SMTP_PORT: int = 587
SMTP_USER: str = "noreply@jze9.ru" SMTP_USER: str = "noreply"
SMTP_PASSWORD: str = "changeme" SMTP_PASSWORD: str = "changeme"
SMTP_FROM: str = "noreply@jze9.ru" SMTP_FROM: str = "noreply@jze9mail.ru"
SMTP_TLS_VERIFY: bool = True
# App # App
APP_URL: str = "https://academic.jze9.ru" APP_URL: str = "https://academic.jze9.ru"

View File

@@ -2,6 +2,7 @@
import logging import logging
import smtplib import smtplib
import ssl
from email.mime.multipart import MIMEMultipart from email.mime.multipart import MIMEMultipart
from email.mime.text import MIMEText from email.mime.text import MIMEText
@@ -37,7 +38,20 @@ class EmailSender:
msg.attach(MIMEText(html_body, "html", "utf-8")) msg.attach(MIMEText(html_body, "html", "utf-8"))
try: 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.login(settings.SMTP_USER, settings.SMTP_PASSWORD)
smtp.send_message(msg) smtp.send_message(msg)
logger.info(f"Email отправлен: {to_email!r}, тема: {subject!r}") 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")) user_id: Mapped[int] = mapped_column(ForeignKey("users.id"))
type: Mapped[str] = mapped_column(String(20)) type: Mapped[str] = mapped_column(String(20))
status: 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) result: Mapped[dict | None] = mapped_column(JSON, nullable=True)
error: Mapped[str | None] = mapped_column(Text, nullable=True) error: 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

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