Compare commits
5 Commits
9894fa9320
...
25e55a3f7d
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
25e55a3f7d | ||
|
|
ef2723e1d3 | ||
|
|
38b000703a | ||
|
|
9ca2bd1742 | ||
|
|
c1cf1ddd2f |
11
.env.example
11
.env.example
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 по токену из письма."""
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
126
services/frontend/src/components/EmailVerificationModal.tsx
Normal file
126
services/frontend/src/components/EmailVerificationModal.tsx
Normal 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;
|
||||
}
|
||||
@@ -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>
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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>
|
||||
);
|
||||
|
||||
83
services/frontend/src/pages/VerifyEmail.tsx
Normal file
83
services/frontend/src/pages/VerifyEmail.tsx
Normal 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>
|
||||
);
|
||||
}
|
||||
@@ -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
|
||||
try:
|
||||
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"Не удалось загрузить FAISS индекс ({e}) — создаём новый")
|
||||
cls._index = cls._new_index()
|
||||
cls._reverse_map = {}
|
||||
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
|
||||
logger.info("Создание нового FAISS индекса IndexIDMap2(IndexFlatIP)...")
|
||||
cls._index = cls._new_index()
|
||||
cls._reverse_map = {}
|
||||
|
||||
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
|
||||
|
||||
# Попытка перенести на GPU
|
||||
try:
|
||||
res = faiss.StandardGpuResources()
|
||||
cls._index = faiss.index_cpu_to_gpu(res, 0, cpu_index)
|
||||
cls._use_gpu = True
|
||||
logger.info("FAISS индекс размещён на GPU")
|
||||
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"GPU недоступна: {e}. Используем CPU FAISS.")
|
||||
cls._index = cpu_index
|
||||
cls._use_gpu = False
|
||||
|
||||
if cls._is_trained:
|
||||
cls._index.nprobe = settings.FAISS_NPROBE
|
||||
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()
|
||||
if len(doc_ids) == 0:
|
||||
return
|
||||
|
||||
vectors = vectors.astype(np.float32)
|
||||
vectors = np.asarray(vectors, dtype=np.float32)
|
||||
ids = np.asarray(doc_ids, dtype=np.int64)
|
||||
|
||||
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
|
||||
# Удалить существующие id, чтобы повторный эмбеддинг не создавал дубли
|
||||
try:
|
||||
cls._index.remove_ids(ids)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
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)
|
||||
cls._index.add_with_ids(vectors, ids)
|
||||
for doc_id in doc_ids:
|
||||
cls._reverse_map[int(doc_id)] = int(doc_id)
|
||||
|
||||
# Добавить в 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
|
||||
|
||||
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
|
||||
|
||||
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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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"
|
||||
|
||||
76
services/worker-indexer/app/fulltext.py
Normal file
76
services/worker-indexer/app/fulltext.py
Normal 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
|
||||
@@ -139,50 +139,64 @@ def extract_and_check(
|
||||
)
|
||||
logger.info(f"Фрагментов создано: {len(fragments)}")
|
||||
|
||||
# ──── Уровень 1: Winnowing fingerprints ────────────────────────────────
|
||||
# ──── Уровень 1: Winnowing fingerprints (пофрагментно, с локализацией) ───
|
||||
# Winnow'им КАЖДЫЙ фрагмент отдельно и ищем, с каким источником и на
|
||||
# какой позиции он совпадает. Так в отчёте видно не «документ похож на X»,
|
||||
# а «фрагмент на символах A–B скопирован из источника X».
|
||||
level1_matches: list[dict] = []
|
||||
doc_fingerprint = winnow(text)
|
||||
from app.models import Document, Fingerprint
|
||||
|
||||
if doc_fingerprint:
|
||||
from app.models import Document, Fingerprint
|
||||
with db_session() as session:
|
||||
doc_cache: dict[int, Document] = {}
|
||||
|
||||
with db_session() as session:
|
||||
hashes = list(doc_fingerprint)[: settings.MAX_FINGERPRINTS_PER_DOC]
|
||||
for frag in fragments:
|
||||
frag_fp = winnow(frag["text"])
|
||||
if not frag_fp:
|
||||
continue
|
||||
frag_hashes = list(frag_fp)
|
||||
|
||||
# Найти совпадения в базе fingerprints
|
||||
matching_docs = session.execute(
|
||||
# Источник, разделяющий больше всего отпечатков с этим фрагментом
|
||||
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:
|
||||
continue
|
||||
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),
|
||||
"method": "exact",
|
||||
"source_title": doc.title,
|
||||
"source_url": doc.url,
|
||||
"source_db": doc.source,
|
||||
})
|
||||
level1_matches.append({
|
||||
"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 эмбеддингов
|
||||
celery_app.send_task(
|
||||
"gpu.embed_documents",
|
||||
args=[[doc_id]],
|
||||
queue="queue.gpu",
|
||||
)
|
||||
# Диспатч 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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
import logging
|
||||
import smtplib
|
||||
import ssl
|
||||
from email.mime.multipart import MIMEMultipart
|
||||
from email.mime.text import MIMEText
|
||||
|
||||
@@ -37,9 +38,22 @@ class EmailSender:
|
||||
msg.attach(MIMEText(html_body, "html", "utf-8"))
|
||||
|
||||
try:
|
||||
with smtplib.SMTP_SSL(settings.SMTP_HOST, settings.SMTP_PORT) as smtp:
|
||||
smtp.login(settings.SMTP_USER, settings.SMTP_PASSWORD)
|
||||
smtp.send_message(msg)
|
||||
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}")
|
||||
except smtplib.SMTPException as e:
|
||||
logger.error(f"SMTP ошибка при отправке письма на {to_email!r}: {e}")
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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 "Низкий"
|
||||
|
||||
Reference in New Issue
Block a user