Заливку Википедии и PMC я сделал отдельными скриптами мимо существующей инфраструктуры: запуск руками через ssh, состояние в файле в /tmp, никакой видимости. Результат предсказуем — за неделю обе умерли молча (обрыв базы на 20 008 статьях из 2 млн и таймаут сети на 85 тыс. из 100 тыс.), прогресс потерялся, а узнали мы об этом через неделю. При том что рядом лежит готовый механизм: parse_sources, прогоны со шкалой, журнал, кнопки, ретраи Celery. Теперь это обычные типы источника — wikipedia_ru и pmc_bulk: - заводятся и запускаются из админки, как OpenAlex или КиберЛенинка; - показывают ту же шкалу, журнал и кнопку остановки; - падение воркера больше не теряет прогресс: позиция продолжения хранится в parse_sources.resume_token (номер статьи в дампе / токен страницы бакета), повторный запуск берёт следующую порцию; - укладываются в бюджет времени таска — заливка идёт порциями, а не одним многосуточным процессом. Чего не хватало конвейеру для миллионов и что добавлено: - парсеры отдают генератор, а не список: 2 млн статей в память не влезают; - app/bulk_writer.py — запись пачками через COPY (21 тыс. строк/с против 6.7 тыс. построчно) с переподключением к базе при обрыве; - эмбеддинги при массовой заливке не диспатчатся: они на порядок медленнее и стали бы узким местом, вектора досчитываются отдельно (reembed_missing.py). scripts/ops/bulk_ingest_*.py удалены — их работу делает конвейер. Проверено на проде: оба парсера отдают документы, прогон через run_parser завершается штатно, позиция продолжения сдвигается (300 → 600), повторный запуск продолжает с неё. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1100 lines
48 KiB
Python
1100 lines
48 KiB
Python
"""Админ-панель: дашборд, клиенты, работы, документы, хранилище, источники, отстойник.
|
||
|
||
Все эндпоинты требуют is_admin (get_admin_user). Эндпоинты под /admin/panel/{code}/*
|
||
дополнительно проверяют секретный код сессии (verify_admin_code).
|
||
"""
|
||
|
||
import contextlib
|
||
import io
|
||
import logging
|
||
import os
|
||
import uuid
|
||
from datetime import datetime, timedelta
|
||
from pathlib import Path
|
||
|
||
import httpx
|
||
from fastapi import APIRouter, Depends, HTTPException, Query, UploadFile, status
|
||
from fastapi.concurrency import run_in_threadpool
|
||
from sqlalchemy import delete, func, select, text
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.config import settings
|
||
from app.core.admin import create_admin_session, get_admin_user
|
||
from app.core.celery_app import celery_app
|
||
from app.core.minio_client import get_minio
|
||
from app.core.redis_client import get_redis
|
||
from app.database import get_db
|
||
from app.models.admin import ParseRun, ParseSource, StagedWork
|
||
from app.models.document import Document, Fingerprint
|
||
from app.models.task import Task
|
||
from app.models.user import User
|
||
from app.schemas.admin import (
|
||
AdminDocumentResponse,
|
||
AdminSessionResponse,
|
||
AdminStats,
|
||
AdminTaskResponse,
|
||
AdminUserResponse,
|
||
AdminUserUpdate,
|
||
ParseRunDetail,
|
||
ParseRunResponse,
|
||
ParseSourceBulkCreate,
|
||
ParseSourceCreate,
|
||
ParseSourceResponse,
|
||
ParseSourceUpdate,
|
||
ServiceHealth,
|
||
StagedWorkDetail,
|
||
StagedWorkResponse,
|
||
)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/admin", tags=["admin"], dependencies=[Depends(get_admin_user)])
|
||
|
||
# Типы источников, для которых в worker-indexer есть парсер (index.run_parser).
|
||
# wikipedia_ru и pmc_bulk — массовые: читают дамп/бакет потоком и пишут пачками,
|
||
# продолжая с сохранённой позиции (parse_sources.resume_token)
|
||
SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc", "wikipedia_ru", "pmc_bulk")
|
||
# Прогон в этих статусах ещё может двигаться — повторно запускать источник нельзя
|
||
RUN_ACTIVE_STATUSES = ("queued", "running")
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# СЕССИЯ ДОСТУПА
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.post("/session", response_model=AdminSessionResponse)
|
||
async def open_admin_session(
|
||
admin: User = Depends(get_admin_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> AdminSessionResponse:
|
||
"""Сгенерировать секретный код доступа и ссылку /admin/<code>."""
|
||
code = await create_admin_session(admin, db)
|
||
from app.models.admin import AdminSession
|
||
|
||
result = await db.execute(select(AdminSession).where(AdminSession.code == code))
|
||
sess = result.scalar_one()
|
||
return AdminSessionResponse(
|
||
code=code,
|
||
url=f"/admin/{code}",
|
||
expires_at=sess.expires_at,
|
||
)
|
||
|
||
|
||
@router.get("/session/verify/{code}")
|
||
async def verify_session(
|
||
code: str,
|
||
admin: User = Depends(get_admin_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> dict:
|
||
"""Проверить валидность кода (для фронта при загрузке /admin/<code>)."""
|
||
from app.models.admin import AdminSession
|
||
|
||
result = await db.execute(select(AdminSession).where(AdminSession.code == code))
|
||
sess = result.scalar_one_or_none()
|
||
if sess is None or sess.user_id != admin.id:
|
||
return {"valid": False}
|
||
expires = sess.expires_at
|
||
if expires.tzinfo is not None:
|
||
expires = expires.replace(tzinfo=None)
|
||
return {"valid": expires >= datetime.utcnow()}
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ДАШБОРД / СЕРВИСЫ
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
async def _rabbitmq_queues() -> dict:
|
||
"""Очереди queue.* из management API RabbitMQ: {имя: сырой словарь очереди}.
|
||
|
||
При недоступности брокера/плагина management возвращает {"error": ...} —
|
||
это не повод валить весь дашборд.
|
||
"""
|
||
try:
|
||
# amqp://user:pass@host:port/
|
||
creds, _, hostpart = settings.RABBITMQ_URL.split("//", 1)[1].rpartition("@")
|
||
rmq_user, _, rmq_pass = creds.partition(":")
|
||
rmq_host = hostpart.split(":")[0].split("/")[0]
|
||
async with httpx.AsyncClient(timeout=5) as c:
|
||
resp = await c.get(f"http://{rmq_host}:15672/api/queues", auth=(rmq_user, rmq_pass))
|
||
if resp.status_code != 200:
|
||
return {"error": f"management API вернул {resp.status_code}"}
|
||
return {
|
||
q["name"]: q for q in resp.json() if str(q.get("name", "")).startswith("queue.")
|
||
}
|
||
except Exception as e:
|
||
return {"error": str(e)[:120]}
|
||
|
||
|
||
@router.get("/health", response_model=list[ServiceHealth])
|
||
async def health(db: AsyncSession = Depends(get_db)) -> list[ServiceHealth]:
|
||
"""Статус всех сервисов инфраструктуры."""
|
||
checks: list[ServiceHealth] = []
|
||
|
||
# PostgreSQL
|
||
try:
|
||
await db.execute(select(func.now()))
|
||
checks.append(ServiceHealth(name="PostgreSQL", ok=True))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="PostgreSQL", ok=False, detail=str(e)[:200]))
|
||
|
||
# Redis
|
||
try:
|
||
r = get_redis()
|
||
await r.ping()
|
||
checks.append(ServiceHealth(name="Redis", ok=True))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="Redis", ok=False, detail=str(e)[:200]))
|
||
|
||
# MinIO
|
||
try:
|
||
get_minio().bucket_exists(settings.MINIO_BUCKET_DOCS)
|
||
checks.append(ServiceHealth(name="MinIO", ok=True))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="MinIO", ok=False, detail=str(e)[:200]))
|
||
|
||
# Elasticsearch
|
||
try:
|
||
async with httpx.AsyncClient(timeout=5) as c:
|
||
resp = await c.get(f"{settings.ELASTICSEARCH_URL}/_cluster/health")
|
||
ok = resp.status_code == 200 and resp.json().get("status") != "red"
|
||
checks.append(ServiceHealth(name="Elasticsearch", ok=ok,
|
||
detail=resp.json().get("status")))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="Elasticsearch", ok=False, detail=str(e)[:200]))
|
||
|
||
# Ollama
|
||
try:
|
||
async with httpx.AsyncClient(timeout=5) as c:
|
||
resp = await c.get(f"{settings.OLLAMA_URL}/api/version")
|
||
checks.append(ServiceHealth(name="Ollama", ok=resp.status_code == 200,
|
||
detail=resp.json().get("version")))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="Ollama", ok=False, detail=str(e)[:200]))
|
||
|
||
# RabbitMQ + воркеры через Celery ping
|
||
try:
|
||
insp = celery_app.control.inspect(timeout=3)
|
||
pong = insp.ping() or {}
|
||
workers = list(pong.keys())
|
||
checks.append(ServiceHealth(name="RabbitMQ/Celery", ok=bool(workers),
|
||
detail=f"воркеров онлайн: {len(workers)}"))
|
||
except Exception as e:
|
||
checks.append(ServiceHealth(name="RabbitMQ/Celery", ok=False, detail=str(e)[:200]))
|
||
|
||
return checks
|
||
|
||
|
||
@router.get("/stats", response_model=AdminStats)
|
||
async def stats(db: AsyncSession = Depends(get_db)) -> AdminStats:
|
||
"""Сводные счётчики для дашборда."""
|
||
users_total = (await db.execute(select(func.count()).select_from(User))).scalar_one()
|
||
|
||
rows = (await db.execute(select(Task.status, func.count()).group_by(Task.status))).all()
|
||
tasks_by_status = dict(rows)
|
||
|
||
docs_total = (await db.execute(select(func.count()).select_from(Document))).scalar_one()
|
||
rows = (await db.execute(select(Document.source, func.count()).group_by(Document.source))).all()
|
||
docs_by_source = dict(rows)
|
||
|
||
staging_pending = (
|
||
await db.execute(
|
||
select(func.count()).select_from(StagedWork).where(StagedWork.status == "pending")
|
||
)
|
||
).scalar_one()
|
||
|
||
# Хранилище MinIO
|
||
storage: dict = {}
|
||
try:
|
||
client = get_minio()
|
||
for bucket in (settings.MINIO_BUCKET_DOCS, settings.MINIO_BUCKET_BACKUPS, "staging"):
|
||
try:
|
||
if not client.bucket_exists(bucket):
|
||
continue
|
||
total = 0
|
||
count = 0
|
||
for obj in client.list_objects(bucket, recursive=True):
|
||
total += obj.size or 0
|
||
count += 1
|
||
storage[bucket] = {"objects": count, "bytes": total}
|
||
except Exception:
|
||
continue
|
||
except Exception as e:
|
||
storage = {"error": str(e)[:200]}
|
||
|
||
# Длины очередей через RabbitMQ management API (если доступно)
|
||
rmq = await _rabbitmq_queues()
|
||
queues: dict = (
|
||
{"error": rmq["error"]}
|
||
if "error" in rmq
|
||
else {name: q.get("messages", 0) for name, q in rmq.items()}
|
||
)
|
||
|
||
return AdminStats(
|
||
users_total=users_total,
|
||
tasks_by_status=tasks_by_status,
|
||
documents_total=docs_total,
|
||
documents_by_source=docs_by_source,
|
||
staging_pending=staging_pending,
|
||
storage=storage,
|
||
queues=queues,
|
||
)
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# КЛИЕНТЫ
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.get("/users", response_model=list[AdminUserResponse])
|
||
async def list_users(
|
||
q: str | None = None,
|
||
limit: int = Query(50, le=200),
|
||
offset: int = 0,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[AdminUserResponse]:
|
||
stmt = select(User)
|
||
if q:
|
||
like = f"%{q.lower()}%"
|
||
stmt = stmt.where(func.lower(User.email).like(like) | func.lower(User.name).like(like))
|
||
stmt = stmt.order_by(User.created_at.desc()).limit(limit).offset(offset)
|
||
users = (await db.execute(stmt)).scalars().all()
|
||
|
||
result = []
|
||
for u in users:
|
||
cnt = (
|
||
await db.execute(select(func.count()).select_from(Task).where(Task.user_id == u.id))
|
||
).scalar_one()
|
||
item = AdminUserResponse.model_validate(u)
|
||
item.tasks_count = cnt
|
||
result.append(item)
|
||
return result
|
||
|
||
|
||
@router.patch("/users/{user_id}", response_model=AdminUserResponse)
|
||
async def update_user(
|
||
user_id: int,
|
||
data: AdminUserUpdate,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> AdminUserResponse:
|
||
user = (await db.execute(select(User).where(User.id == user_id))).scalar_one_or_none()
|
||
if user is None:
|
||
raise HTTPException(status_code=404, detail="Пользователь не найден")
|
||
if data.plan is not None:
|
||
user.plan = data.plan
|
||
if data.is_verified is not None:
|
||
user.is_verified = data.is_verified
|
||
if data.is_admin is not None:
|
||
user.is_admin = data.is_admin
|
||
await db.commit()
|
||
await db.refresh(user)
|
||
# Сбросить кэш пользователя
|
||
with contextlib.suppress(Exception):
|
||
await get_redis().delete(f"user:cache:{user_id}")
|
||
return AdminUserResponse.model_validate(user)
|
||
|
||
|
||
@router.delete("/users/{user_id}", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def delete_user(user_id: int, db: AsyncSession = Depends(get_db)) -> None:
|
||
user = (await db.execute(select(User).where(User.id == user_id))).scalar_one_or_none()
|
||
if user is None:
|
||
raise HTTPException(status_code=404, detail="Пользователь не найден")
|
||
await db.execute(delete(Task).where(Task.user_id == user_id))
|
||
await db.delete(user)
|
||
await db.commit()
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# РАБОТЫ (задачи)
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.get("/tasks", response_model=list[AdminTaskResponse])
|
||
async def list_tasks(
|
||
status_filter: str | None = Query(None, alias="status"),
|
||
type_filter: str | None = Query(None, alias="type"),
|
||
user_id: int | None = None,
|
||
limit: int = Query(50, le=200),
|
||
offset: int = 0,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[AdminTaskResponse]:
|
||
stmt = select(Task)
|
||
if status_filter:
|
||
stmt = stmt.where(Task.status == status_filter)
|
||
if type_filter:
|
||
stmt = stmt.where(Task.type == type_filter)
|
||
if user_id:
|
||
stmt = stmt.where(Task.user_id == user_id)
|
||
stmt = stmt.order_by(Task.created_at.desc()).limit(limit).offset(offset)
|
||
tasks = (await db.execute(stmt)).scalars().all()
|
||
return [AdminTaskResponse.model_validate(t) for t in tasks]
|
||
|
||
|
||
@router.get("/tasks/{public_id}")
|
||
async def get_task(public_id: str, db: AsyncSession = Depends(get_db)) -> dict:
|
||
task = (
|
||
await db.execute(select(Task).where(Task.public_id == public_id))
|
||
).scalar_one_or_none()
|
||
if task is None:
|
||
raise HTTPException(status_code=404, detail="Задача не найдена")
|
||
return {
|
||
"public_id": task.public_id,
|
||
"user_id": task.user_id,
|
||
"type": task.type,
|
||
"status": task.status,
|
||
"input_data": task.input_data,
|
||
"result": task.result,
|
||
"error": task.error,
|
||
"created_at": task.created_at,
|
||
}
|
||
|
||
|
||
@router.post("/tasks/{public_id}/requeue")
|
||
async def requeue_task(public_id: str, db: AsyncSession = Depends(get_db)) -> dict:
|
||
task = (
|
||
await db.execute(select(Task).where(Task.public_id == public_id))
|
||
).scalar_one_or_none()
|
||
if task is None:
|
||
raise HTTPException(status_code=404, detail="Задача не найдена")
|
||
if task.type != "plagiarism":
|
||
raise HTTPException(status_code=400, detail="Перезапуск поддержан только для проверок плагиата")
|
||
minio_key = (task.input_data or {}).get("minio_key")
|
||
filename = (task.input_data or {}).get("filename")
|
||
if not minio_key:
|
||
raise HTTPException(status_code=400, detail="Нет minio_key для перезапуска")
|
||
task.status = "queued"
|
||
task.error = None
|
||
task.result = None
|
||
await db.commit()
|
||
celery_app.send_task(
|
||
"index.extract_and_check",
|
||
args=[task.id, minio_key, filename],
|
||
queue="queue.index",
|
||
)
|
||
return {"status": "requeued", "public_id": public_id}
|
||
|
||
|
||
@router.delete("/tasks/{public_id}", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def delete_task(public_id: str, db: AsyncSession = Depends(get_db)) -> None:
|
||
task = (
|
||
await db.execute(select(Task).where(Task.public_id == public_id))
|
||
).scalar_one_or_none()
|
||
if task is None:
|
||
raise HTTPException(status_code=404, detail="Задача не найдена")
|
||
await db.delete(task)
|
||
await db.commit()
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# БАЗА ДОКУМЕНТОВ
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.get("/documents", response_model=list[AdminDocumentResponse])
|
||
async def list_documents(
|
||
q: str | None = None,
|
||
source: str | None = None,
|
||
year: int | None = None,
|
||
limit: int = Query(50, le=200),
|
||
offset: int = 0,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[AdminDocumentResponse]:
|
||
stmt = select(Document)
|
||
if q:
|
||
stmt = stmt.where(func.lower(Document.title).like(f"%{q.lower()}%"))
|
||
if source:
|
||
stmt = stmt.where(Document.source == source)
|
||
if year:
|
||
stmt = stmt.where(Document.year == year)
|
||
stmt = stmt.order_by(Document.indexed_at.desc()).limit(limit).offset(offset)
|
||
docs = (await db.execute(stmt)).scalars().all()
|
||
return [AdminDocumentResponse.model_validate(d) for d in docs]
|
||
|
||
|
||
@router.delete("/documents/{doc_id}", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def delete_document(doc_id: int, db: AsyncSession = Depends(get_db)) -> None:
|
||
doc = (await db.execute(select(Document).where(Document.id == doc_id))).scalar_one_or_none()
|
||
if doc is None:
|
||
raise HTTPException(status_code=404, detail="Документ не найден")
|
||
await db.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id))
|
||
await db.delete(doc)
|
||
await db.commit()
|
||
# Удалить из Elasticsearch (best-effort)
|
||
try:
|
||
async with httpx.AsyncClient(timeout=5) as c:
|
||
await c.delete(f"{settings.ELASTICSEARCH_URL}/documents/_doc/{doc_id}")
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ХРАНИЛИЩЕ
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.get("/storage")
|
||
async def storage_info() -> dict:
|
||
client = get_minio()
|
||
result: dict = {"buckets": []}
|
||
for bucket in (settings.MINIO_BUCKET_DOCS, settings.MINIO_BUCKET_BACKUPS, "staging"):
|
||
try:
|
||
if not client.bucket_exists(bucket):
|
||
continue
|
||
total = 0
|
||
objs = []
|
||
for obj in client.list_objects(bucket, recursive=True):
|
||
total += obj.size or 0
|
||
if len(objs) < 100:
|
||
objs.append({"name": obj.object_name, "size": obj.size})
|
||
result["buckets"].append({
|
||
"name": bucket,
|
||
"objects_total": len(objs),
|
||
"bytes": total,
|
||
"sample": objs,
|
||
})
|
||
except Exception as e:
|
||
result["buckets"].append({"name": bucket, "error": str(e)[:120]})
|
||
return result
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ИСТОЧНИКИ ПАРСИНГА
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
async def _attach_runs(
|
||
sources: list[ParseSource], db: AsyncSession
|
||
) -> list[ParseSourceResponse]:
|
||
"""Подтянуть последние прогоны одним запросом (иначе N+1 на 170+ источников)."""
|
||
run_ids = [s.last_run_id for s in sources if s.last_run_id]
|
||
runs: dict[int, ParseRun] = {}
|
||
if run_ids:
|
||
rows = (await db.execute(select(ParseRun).where(ParseRun.id.in_(run_ids)))).scalars().all()
|
||
runs = {r.id: r for r in rows}
|
||
|
||
result = []
|
||
for s in sources:
|
||
item = ParseSourceResponse.model_validate(s)
|
||
run = runs.get(s.last_run_id) if s.last_run_id else None
|
||
item.last_run = ParseRunResponse.model_validate(run) if run else None
|
||
result.append(item)
|
||
return result
|
||
|
||
|
||
async def _start_run(src: ParseSource, db: AsyncSession) -> ParseRun:
|
||
"""Создать строку прогона и отправить таск парсера.
|
||
|
||
Строку создаёт API, а не воркер: при массовом запуске 170+ источников
|
||
очередь разбирается минутами, и без неё в админке не было бы видно, что
|
||
источник уже поставлен в очередь (шкала висела бы на прошлом прогоне).
|
||
"""
|
||
run = ParseRun(source_id=src.id, status="queued", stage="queued", target=src.limit, log=[])
|
||
db.add(run)
|
||
await db.flush()
|
||
|
||
src.last_status = "running"
|
||
src.last_error = None
|
||
src.last_run_at = datetime.utcnow()
|
||
src.last_run_id = run.id
|
||
await db.commit()
|
||
await db.refresh(run)
|
||
|
||
celery_result = celery_app.send_task(
|
||
"index.run_parser", args=[src.id, run.id], queue="queue.index"
|
||
)
|
||
run.celery_task_id = celery_result.id
|
||
await db.commit()
|
||
await db.refresh(run)
|
||
return run
|
||
|
||
|
||
@router.get("/sources", response_model=list[ParseSourceResponse])
|
||
async def list_sources(db: AsyncSession = Depends(get_db)) -> list[ParseSourceResponse]:
|
||
rows = (
|
||
await db.execute(select(ParseSource).order_by(ParseSource.created_at.desc()))
|
||
).scalars().all()
|
||
return await _attach_runs(list(rows), db)
|
||
|
||
|
||
@router.post("/sources", response_model=ParseSourceResponse, status_code=201)
|
||
async def create_source(data: ParseSourceCreate, db: AsyncSession = Depends(get_db)) -> ParseSourceResponse:
|
||
if data.source_type not in SOURCE_TYPES:
|
||
raise HTTPException(status_code=400, detail="Недопустимый тип источника")
|
||
src = ParseSource(**data.model_dump())
|
||
db.add(src)
|
||
await db.commit()
|
||
await db.refresh(src)
|
||
return ParseSourceResponse.model_validate(src)
|
||
|
||
|
||
@router.post("/sources/bulk", status_code=201)
|
||
async def create_sources_bulk(
|
||
data: ParseSourceBulkCreate, db: AsyncSession = Depends(get_db)
|
||
) -> dict:
|
||
"""Добавить пачку источников одного типа — по одному на каждый запрос-тему."""
|
||
if data.source_type not in SOURCE_TYPES:
|
||
raise HTTPException(status_code=400, detail="Недопустимый тип источника")
|
||
|
||
queries = [q.strip() for q in data.queries if q.strip()]
|
||
if not queries:
|
||
raise HTTPException(status_code=400, detail="Пустой список запросов")
|
||
|
||
prefix = data.name_prefix or data.source_type
|
||
created: list[ParseSource] = []
|
||
for q in queries:
|
||
src = ParseSource(
|
||
source_type=data.source_type,
|
||
name=f"{prefix}:{q}"[:255],
|
||
query=q,
|
||
lang=data.lang,
|
||
year_from=data.year_from,
|
||
year_to=data.year_to,
|
||
limit=data.limit,
|
||
)
|
||
db.add(src)
|
||
created.append(src)
|
||
await db.commit()
|
||
|
||
started = 0
|
||
if data.run_now:
|
||
for src in created:
|
||
await db.refresh(src)
|
||
await _start_run(src, db)
|
||
started += 1
|
||
|
||
return {"created": len(created), "started": started}
|
||
|
||
|
||
@router.patch("/sources/{source_id}", response_model=ParseSourceResponse)
|
||
async def update_source(
|
||
source_id: int, data: ParseSourceUpdate, db: AsyncSession = Depends(get_db)
|
||
) -> ParseSourceResponse:
|
||
src = (await db.execute(select(ParseSource).where(ParseSource.id == source_id))).scalar_one_or_none()
|
||
if src is None:
|
||
raise HTTPException(status_code=404, detail="Источник не найден")
|
||
for k, v in data.model_dump(exclude_unset=True).items():
|
||
setattr(src, k, v)
|
||
await db.commit()
|
||
await db.refresh(src)
|
||
return ParseSourceResponse.model_validate(src)
|
||
|
||
|
||
@router.delete("/sources/{source_id}", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def delete_source(source_id: int, db: AsyncSession = Depends(get_db)) -> None:
|
||
src = (await db.execute(select(ParseSource).where(ParseSource.id == source_id))).scalar_one_or_none()
|
||
if src is None:
|
||
raise HTTPException(status_code=404, detail="Источник не найден")
|
||
await db.delete(src)
|
||
await db.commit()
|
||
|
||
|
||
@router.post("/sources/run-all")
|
||
async def run_all_sources(
|
||
source_type: str | None = None,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> dict:
|
||
"""Запустить заливку по всем включённым источникам (кроме уже идущих)."""
|
||
stmt = select(ParseSource).where(ParseSource.enabled.is_(True))
|
||
if source_type:
|
||
stmt = stmt.where(ParseSource.source_type == source_type)
|
||
sources = (await db.execute(stmt.order_by(ParseSource.id))).scalars().all()
|
||
|
||
active_ids = set(
|
||
(
|
||
await db.execute(
|
||
select(ParseRun.source_id).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES))
|
||
)
|
||
).scalars().all()
|
||
)
|
||
|
||
started, skipped = 0, 0
|
||
for src in sources:
|
||
if src.id in active_ids:
|
||
skipped += 1
|
||
continue
|
||
await _start_run(src, db)
|
||
started += 1
|
||
|
||
logger.info("Массовый запуск источников: старт %d, пропущено %d", started, skipped)
|
||
return {"started": started, "skipped_active": skipped, "total": len(sources)}
|
||
|
||
|
||
@router.post("/sources/stop-all")
|
||
async def stop_all_sources(db: AsyncSession = Depends(get_db)) -> dict:
|
||
"""Отменить все активные прогоны (кооперативно — воркер увидит на своём тике)."""
|
||
runs = (
|
||
await db.execute(select(ParseRun).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES)))
|
||
).scalars().all()
|
||
for run in runs:
|
||
await _cancel_run(run, db)
|
||
await db.commit()
|
||
return {"cancelled": len(runs)}
|
||
|
||
|
||
@router.post("/sources/{source_id}/run", response_model=ParseRunResponse)
|
||
async def run_source(source_id: int, db: AsyncSession = Depends(get_db)) -> ParseRunResponse:
|
||
src = (await db.execute(select(ParseSource).where(ParseSource.id == source_id))).scalar_one_or_none()
|
||
if src is None:
|
||
raise HTTPException(status_code=404, detail="Источник не найден")
|
||
|
||
active = (
|
||
await db.execute(
|
||
select(ParseRun).where(
|
||
ParseRun.source_id == source_id, ParseRun.status.in_(RUN_ACTIVE_STATUSES)
|
||
)
|
||
)
|
||
).scalars().first()
|
||
if active is not None:
|
||
raise HTTPException(status_code=409, detail="Источник уже парсится")
|
||
|
||
run = await _start_run(src, db)
|
||
return ParseRunResponse.model_validate(run)
|
||
|
||
|
||
async def _cancel_run(run: ParseRun, db: AsyncSession) -> None:
|
||
"""Отменить прогон: флаг для воркера + revoke, если таск ещё не начали.
|
||
|
||
Пока таск в очереди — revoke снимает его, и никто уже не переведёт прогон
|
||
из queued, поэтому закрываем строку прямо здесь. Начатый прогон трогать
|
||
жёстко нельзя (оборвётся посреди записи в базу) — ставим флаг, воркер
|
||
остановится сам на ближайшем тике прогресса.
|
||
"""
|
||
run.cancel_requested = True
|
||
if run.status == "queued":
|
||
if run.celery_task_id:
|
||
with contextlib.suppress(Exception):
|
||
celery_app.control.revoke(run.celery_task_id)
|
||
run.status = "cancelled"
|
||
run.stage = "finished"
|
||
run.finished_at = datetime.utcnow()
|
||
src = await db.get(ParseSource, run.source_id)
|
||
if src and src.last_run_id == run.id:
|
||
src.last_status = "cancelled"
|
||
|
||
|
||
@router.post("/sources/{source_id}/cancel")
|
||
async def cancel_source_run(source_id: int, db: AsyncSession = Depends(get_db)) -> dict:
|
||
runs = (
|
||
await db.execute(
|
||
select(ParseRun).where(
|
||
ParseRun.source_id == source_id, ParseRun.status.in_(RUN_ACTIVE_STATUSES)
|
||
)
|
||
)
|
||
).scalars().all()
|
||
if not runs:
|
||
raise HTTPException(status_code=404, detail="Активных прогонов нет")
|
||
for run in runs:
|
||
await _cancel_run(run, db)
|
||
await db.commit()
|
||
return {"cancelled": len(runs), "source_id": source_id}
|
||
|
||
|
||
@router.get("/sources/{source_id}/runs", response_model=list[ParseRunResponse])
|
||
async def list_source_runs(
|
||
source_id: int,
|
||
limit: int = Query(20, le=100),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[ParseRunResponse]:
|
||
rows = (
|
||
await db.execute(
|
||
select(ParseRun)
|
||
.where(ParseRun.source_id == source_id)
|
||
.order_by(ParseRun.started_at.desc())
|
||
.limit(limit)
|
||
)
|
||
).scalars().all()
|
||
return [ParseRunResponse.model_validate(r) for r in rows]
|
||
|
||
|
||
@router.get("/runs/active", response_model=list[ParseRunResponse])
|
||
async def list_active_runs(db: AsyncSession = Depends(get_db)) -> list[ParseRunResponse]:
|
||
"""Все прогоны в работе — для панели отладки и общей шкалы заливки."""
|
||
rows = (
|
||
await db.execute(
|
||
select(ParseRun)
|
||
.where(ParseRun.status.in_(RUN_ACTIVE_STATUSES))
|
||
.order_by(ParseRun.started_at)
|
||
)
|
||
).scalars().all()
|
||
return [ParseRunResponse.model_validate(r) for r in rows]
|
||
|
||
|
||
@router.get("/runs/{run_id}", response_model=ParseRunDetail)
|
||
async def get_run(run_id: int, db: AsyncSession = Depends(get_db)) -> ParseRunDetail:
|
||
run = await db.get(ParseRun, run_id)
|
||
if run is None:
|
||
raise HTTPException(status_code=404, detail="Прогон не найден")
|
||
detail = ParseRunDetail.model_validate(run)
|
||
detail.log = run.log or []
|
||
return detail
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ЗАГРУЗКА РАБОТ В КОРПУС
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
UPLOAD_EXTENSIONS = {".pdf", ".docx", ".txt"}
|
||
UPLOAD_MAX_BYTES = 100 * 1024 * 1024
|
||
UPLOAD_CONTENT_TYPES = {
|
||
".pdf": "application/pdf",
|
||
".docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
|
||
".txt": "text/plain",
|
||
}
|
||
|
||
|
||
@router.post("/documents/upload", status_code=202)
|
||
async def upload_documents(files: list[UploadFile]) -> dict:
|
||
"""Залить готовые работы прямо в базу сравнения (минуя проверку и отстойник).
|
||
|
||
Отличие от пользовательской загрузки (/documents/check): там файл проверяют
|
||
на плагиат, здесь он сам становится источником, с которым будут сравнивать.
|
||
"""
|
||
if not files:
|
||
raise HTTPException(status_code=400, detail="Файлы не переданы")
|
||
|
||
accepted: list[dict] = []
|
||
rejected: list[dict] = []
|
||
minio = get_minio()
|
||
|
||
for file in files:
|
||
name = file.filename or "без имени"
|
||
ext = Path(name).suffix.lower()
|
||
if ext not in UPLOAD_EXTENSIONS:
|
||
rejected.append({"filename": name, "reason": f"формат {ext or '—'} не поддержан"})
|
||
continue
|
||
|
||
data = await file.read()
|
||
if not data:
|
||
rejected.append({"filename": name, "reason": "пустой файл"})
|
||
continue
|
||
if len(data) > UPLOAD_MAX_BYTES:
|
||
rejected.append({"filename": name, "reason": "больше 100 МБ"})
|
||
continue
|
||
|
||
minio_key = f"corpus-upload/{uuid.uuid4()}{ext}"
|
||
try:
|
||
minio.put_object(
|
||
bucket_name=settings.MINIO_BUCKET_DOCS,
|
||
object_name=minio_key,
|
||
data=io.BytesIO(data),
|
||
length=len(data),
|
||
content_type=UPLOAD_CONTENT_TYPES.get(ext, "application/octet-stream"),
|
||
)
|
||
except Exception as e:
|
||
logger.error("Загрузка %s в MinIO не удалась: %s", name, e)
|
||
rejected.append({"filename": name, "reason": "ошибка сохранения в хранилище"})
|
||
continue
|
||
|
||
celery_app.send_task(
|
||
"index.ingest_upload",
|
||
args=[minio_key, name],
|
||
kwargs={"meta": {"title": Path(name).stem}},
|
||
queue="queue.index",
|
||
)
|
||
accepted.append({"filename": name, "minio_key": minio_key, "bytes": len(data)})
|
||
|
||
logger.info("Админ залил в корпус: принято %d, отклонено %d", len(accepted), len(rejected))
|
||
return {"accepted": accepted, "rejected": rejected}
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ОТЛАДКА
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
def _vector_index_stats() -> dict:
|
||
"""Реальное наполнение векторного индекса — спросить у worker-gpu.
|
||
|
||
Колонка `documents.faiss_id` для этого не годится: пометка остаётся и когда
|
||
вектор в индекс не попал, и когда индекс пересоздали после смены модели
|
||
эмбеддингов. Отладка, показывающая покрытие L3 по ней, завышает его в разы.
|
||
"""
|
||
try:
|
||
result = celery_app.send_task("gpu.index_stats", queue="queue.gpu")
|
||
return result.get(timeout=10) or {}
|
||
except Exception as e:
|
||
return {"error": str(e)[:200] or type(e).__name__}
|
||
|
||
|
||
def _celery_snapshot() -> dict:
|
||
"""Живое состояние воркеров через Celery inspect (блокирующий вызов)."""
|
||
try:
|
||
insp = celery_app.control.inspect(timeout=4)
|
||
ping = insp.ping() or {}
|
||
active = insp.active() or {}
|
||
reserved = insp.reserved() or {}
|
||
stats = insp.stats() or {}
|
||
except Exception as e:
|
||
return {"error": str(e)[:200], "workers": []}
|
||
|
||
workers = []
|
||
for name in sorted(set(ping) | set(stats)):
|
||
wstats = stats.get(name, {})
|
||
workers.append({
|
||
"name": name,
|
||
"concurrency": (wstats.get("pool") or {}).get("max-concurrency"),
|
||
"reserved": len(reserved.get(name, [])),
|
||
"active": [
|
||
{
|
||
"name": t.get("name"),
|
||
"id": t.get("id"),
|
||
# args целиком не отдаём: в них бывает текст работы целиком
|
||
"args": str(t.get("args"))[:120],
|
||
"started_ago_s": (
|
||
round(datetime.utcnow().timestamp() - t["time_start"])
|
||
if t.get("time_start") else None
|
||
),
|
||
}
|
||
for t in active.get(name, [])
|
||
],
|
||
})
|
||
return {"workers": workers}
|
||
|
||
|
||
@router.get("/debug")
|
||
async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
|
||
"""Полный срез состояния системы для отладки заливки и проверок.
|
||
|
||
Один запрос вместо похода по пяти вкладкам: кто из воркеров жив и что
|
||
именно сейчас крутит, что копится в очередях, как наполняется корпус,
|
||
какие прогоны идут и на чём падали последние.
|
||
"""
|
||
# Воркеры и статистика индекса — блокирующие вызовы Celery, в пул потоков,
|
||
# чтобы не вешать событийный цикл
|
||
celery_state = await run_in_threadpool(_celery_snapshot)
|
||
vector_stats = await run_in_threadpool(_vector_index_stats)
|
||
|
||
rmq = await _rabbitmq_queues()
|
||
queues = (
|
||
{"error": rmq["error"]}
|
||
if "error" in rmq
|
||
else {
|
||
name: {
|
||
"messages": q.get("messages", 0),
|
||
"unacked": q.get("messages_unacknowledged", 0),
|
||
"consumers": q.get("consumers", 0),
|
||
}
|
||
for name, q in sorted(rmq.items())
|
||
}
|
||
)
|
||
|
||
# ── Корпус ────────────────────────────────────────────────────────────────
|
||
docs_total = (await db.execute(select(func.count()).select_from(Document))).scalar_one()
|
||
docs_embedded = (
|
||
await db.execute(
|
||
select(func.count()).select_from(Document).where(Document.faiss_id.isnot(None))
|
||
)
|
||
).scalar_one()
|
||
# count(*) по fingerprints — это десятки миллионов строк и секунды ожидания;
|
||
# для панели отладки достаточно оценки планировщика
|
||
fingerprints_est = (
|
||
await db.execute(text("SELECT reltuples::bigint FROM pg_class WHERE relname='fingerprints'"))
|
||
).scalar() or 0
|
||
|
||
es_docs = None
|
||
try:
|
||
async with httpx.AsyncClient(timeout=5) as c:
|
||
resp = await c.get(f"{settings.ELASTICSEARCH_URL}/documents/_count")
|
||
if resp.status_code == 200:
|
||
es_docs = resp.json().get("count")
|
||
except Exception:
|
||
es_docs = None
|
||
|
||
# ── Источники и прогоны ───────────────────────────────────────────────────
|
||
src_rows = (
|
||
await db.execute(select(ParseSource.last_status, func.count()).group_by(ParseSource.last_status))
|
||
).all()
|
||
active_runs = (
|
||
await db.execute(
|
||
select(ParseRun).where(ParseRun.status.in_(RUN_ACTIVE_STATUSES)).order_by(ParseRun.id)
|
||
)
|
||
).scalars().all()
|
||
recent_failed_runs = (
|
||
await db.execute(
|
||
select(ParseRun)
|
||
.where(ParseRun.status.in_(("error", "partial")))
|
||
.order_by(ParseRun.started_at.desc())
|
||
.limit(10)
|
||
)
|
||
).scalars().all()
|
||
|
||
# Зависший прогон: числится в работе, но воркер давно не отчитывался
|
||
stale_cutoff = datetime.utcnow() - timedelta(minutes=10)
|
||
stale_runs = [
|
||
r.id for r in active_runs
|
||
if r.status == "running" and (r.heartbeat_at is None or r.heartbeat_at < stale_cutoff)
|
||
]
|
||
|
||
# ── Задачи пользователей ──────────────────────────────────────────────────
|
||
day_ago = datetime.utcnow() - timedelta(hours=24)
|
||
tasks_24h = dict(
|
||
(
|
||
await db.execute(
|
||
select(Task.status, func.count()).where(Task.created_at >= day_ago).group_by(Task.status)
|
||
)
|
||
).all()
|
||
)
|
||
failed_tasks = (
|
||
await db.execute(
|
||
select(Task)
|
||
.where(Task.status == "failed")
|
||
.order_by(Task.created_at.desc())
|
||
.limit(10)
|
||
)
|
||
).scalars().all()
|
||
|
||
# Зависшая проверка: числится в работе, но воркер давно её не трогал.
|
||
# Пользователь видит вечное «обрабатывается» — молча и без следов в логах,
|
||
# поэтому такие задачи должны быть видны в отладке (нашлись экземпляры,
|
||
# висевшие с мая).
|
||
#
|
||
# Возраст считает сама база: колонки в PostgreSQL — timestamptz, а модель
|
||
# объявляет их без зоны, и передача сюда Python-времени ломается на
|
||
# несовпадении (asyncpg отвергает aware-параметр, naive даёт неверный сдвиг).
|
||
stale_tasks = (
|
||
await db.execute(
|
||
text("""
|
||
SELECT public_id, type, created_at,
|
||
round(extract(epoch FROM now() - coalesce(updated_at, created_at))
|
||
/ 3600.0, 1) AS stuck_hours
|
||
FROM tasks
|
||
WHERE status = 'processing'
|
||
AND coalesce(updated_at, created_at) < now() - interval '2 hours'
|
||
ORDER BY created_at DESC
|
||
LIMIT 20
|
||
""")
|
||
)
|
||
).all()
|
||
|
||
return {
|
||
"generated_at": datetime.utcnow().isoformat(timespec="seconds"),
|
||
"celery": celery_state,
|
||
"queues": queues,
|
||
"corpus": {
|
||
"documents": docs_total,
|
||
# Пометка в БД и реальное содержимое индекса расходятся — показываем
|
||
# обе величины, иначе покрытие L3 выглядит лучше, чем оно есть
|
||
"documents_marked_embedded": docs_embedded,
|
||
"vector_index": vector_stats,
|
||
"fingerprints_estimate": int(fingerprints_est),
|
||
"elasticsearch_documents": es_docs,
|
||
},
|
||
"sources": {
|
||
"by_status": dict(src_rows),
|
||
"active_runs": [ParseRunResponse.model_validate(r).model_dump() for r in active_runs],
|
||
"stale_run_ids": stale_runs,
|
||
"recent_problem_runs": [
|
||
ParseRunResponse.model_validate(r).model_dump() for r in recent_failed_runs
|
||
],
|
||
},
|
||
"tasks_24h": tasks_24h,
|
||
"recent_failed_tasks": [
|
||
{
|
||
"public_id": t.public_id,
|
||
"type": t.type,
|
||
"error": (t.error or "")[:200],
|
||
"created_at": t.created_at,
|
||
}
|
||
for t in failed_tasks
|
||
],
|
||
"stale_tasks": [
|
||
{
|
||
"public_id": t.public_id,
|
||
"type": t.type,
|
||
"created_at": t.created_at,
|
||
"stuck_hours": float(t.stuck_hours),
|
||
}
|
||
for t in stale_tasks
|
||
],
|
||
# Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на
|
||
# поведение воркеров — при разборе «почему не так считает» нужны первыми
|
||
"config": {
|
||
"environment": settings.ENVIRONMENT,
|
||
"embed_backend": os.environ.get("EMBED_BACKEND", "ollama"),
|
||
"llm_backend": os.environ.get("LLM_BACKEND", "ollama"),
|
||
"vector_backend": os.environ.get("VECTOR_BACKEND", "faiss"),
|
||
"embed_model": os.environ.get("EMBED_MODEL", "—"),
|
||
"fetch_full_text": os.environ.get("FETCH_FULL_TEXT", "false"),
|
||
"auto_approve_submissions": os.environ.get("AUTO_APPROVE_SUBMISSIONS", "false"),
|
||
"parser_time_budget_s": os.environ.get("PARSER_TIME_BUDGET_S", "1500"),
|
||
"ollama_url": settings.OLLAMA_URL,
|
||
"elasticsearch_url": settings.ELASTICSEARCH_URL,
|
||
},
|
||
}
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
# ОТСТОЙНИК
|
||
# ═══════════════════════════════════════════════════════════════════════════════
|
||
@router.get("/staging", response_model=list[StagedWorkResponse])
|
||
async def list_staging(
|
||
status_filter: str = Query("pending", alias="status"),
|
||
limit: int = Query(50, le=200),
|
||
offset: int = 0,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[StagedWorkResponse]:
|
||
stmt = select(StagedWork)
|
||
if status_filter != "all":
|
||
stmt = stmt.where(StagedWork.status == status_filter)
|
||
stmt = stmt.order_by(StagedWork.created_at.desc()).limit(limit).offset(offset)
|
||
rows = (await db.execute(stmt)).scalars().all()
|
||
return [StagedWorkResponse.model_validate(s) for s in rows]
|
||
|
||
|
||
@router.get("/staging/{staged_id}", response_model=StagedWorkDetail)
|
||
async def get_staging(staged_id: int, db: AsyncSession = Depends(get_db)) -> StagedWorkDetail:
|
||
sw = (await db.execute(select(StagedWork).where(StagedWork.id == staged_id))).scalar_one_or_none()
|
||
if sw is None:
|
||
raise HTTPException(status_code=404, detail="Работа не найдена")
|
||
detail = StagedWorkDetail.model_validate(sw)
|
||
# Превью текста из MinIO
|
||
if sw.text_key:
|
||
try:
|
||
obj = get_minio().get_object("staging", sw.text_key)
|
||
text = obj.read().decode("utf-8", errors="replace")
|
||
detail.text_preview = text[:5000]
|
||
except Exception as e:
|
||
detail.text_preview = f"(не удалось прочитать текст: {e})"
|
||
return detail
|
||
|
||
|
||
@router.post("/staging/{staged_id}/approve")
|
||
async def approve_staging(
|
||
staged_id: int,
|
||
admin: User = Depends(get_admin_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> dict:
|
||
sw = (await db.execute(select(StagedWork).where(StagedWork.id == staged_id))).scalar_one_or_none()
|
||
if sw is None:
|
||
raise HTTPException(status_code=404, detail="Работа не найдена")
|
||
if sw.status != "pending":
|
||
raise HTTPException(status_code=409, detail=f"Работа уже {sw.status}")
|
||
|
||
# Прочитать извлечённый текст
|
||
full_text = ""
|
||
if sw.text_key:
|
||
try:
|
||
obj = get_minio().get_object("staging", sw.text_key)
|
||
full_text = obj.read().decode("utf-8", errors="replace")
|
||
except Exception as e:
|
||
raise HTTPException(status_code=500, detail=f"Не удалось прочитать текст: {e}") from e
|
||
|
||
# Добавить в базу документов через существующую задачу индексатора
|
||
doc_data = {
|
||
"source": "user_submission",
|
||
"ext_id": f"staged:{sw.id}",
|
||
"title": sw.title or sw.filename or f"Работа #{sw.id}",
|
||
"authors": sw.authors or [],
|
||
"year": sw.year,
|
||
"lang": sw.lang,
|
||
"abstract": full_text[:2000],
|
||
"full_text": full_text,
|
||
}
|
||
celery_app.send_task("index.add_document", args=[doc_data], queue="queue.index")
|
||
|
||
sw.status = "approved"
|
||
sw.reviewed_by = admin.id
|
||
sw.reviewed_at = datetime.utcnow()
|
||
await db.commit()
|
||
return {"status": "approved", "id": staged_id}
|
||
|
||
|
||
@router.post("/staging/{staged_id}/reject")
|
||
async def reject_staging(
|
||
staged_id: int,
|
||
admin: User = Depends(get_admin_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> dict:
|
||
sw = (await db.execute(select(StagedWork).where(StagedWork.id == staged_id))).scalar_one_or_none()
|
||
if sw is None:
|
||
raise HTTPException(status_code=404, detail="Работа не найдена")
|
||
sw.status = "rejected"
|
||
sw.reviewed_by = admin.id
|
||
sw.reviewed_at = datetime.utcnow()
|
||
await db.commit()
|
||
return {"status": "rejected", "id": staged_id}
|