Раньше эмбеддинг-модель (уровень 3) гоняла на CPU внутри worker-gpu —
GPU CT108 использовался только под LLM-парафраз (уровень 4). Теперь
эмбеддинги идут через Ollama /api/embed на отдельной VM с RX 580
(Vulkan-бэкенд, без возни с ROCm/HIP для этой карты). EMBED_BACKEND
переключаемый ("ollama" | "sentence_transformers"), дефолт — ollama.
Модель сменилась на bge-m3 (1024-мерный вектор вместо 768 у
paraphrase-multilingual-mpnet-base-v2) — несовместимо с уже посчитанным
FAISS-индексом, нужна полная переиндексация корпуса после деплоя.
Заодно докстринг delete_task в tasks.py — снятое раньше ограничение
"нельзя удалить processing" оставляло враньё в докстринге.
124 lines
4.7 KiB
Python
124 lines
4.7 KiB
Python
"""Управление задачами пользователя.
|
||
|
||
В URL используется public_id (не внутренний UUID), чтобы не раскрывать
|
||
структуру внутренних идентификаторов.
|
||
"""
|
||
|
||
import logging
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, status
|
||
from sqlalchemy import select
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.core.minio_client import get_minio
|
||
from app.core.security import get_current_user
|
||
from app.database import get_db
|
||
from app.models.admin import StagedWork
|
||
from app.models.task import Task
|
||
from app.models.user import User
|
||
from app.schemas.tasks import TaskResponse
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/tasks", tags=["tasks"])
|
||
|
||
|
||
@router.get("/", response_model=list[TaskResponse])
|
||
async def list_tasks(
|
||
limit: int = 20,
|
||
offset: int = 0,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> list[TaskResponse]:
|
||
"""Список задач текущего пользователя (новые первыми)."""
|
||
# Жёсткий потолок limit — не даём вытащить всю таблицу одним запросом
|
||
limit = min(limit, 100)
|
||
|
||
result = await db.execute(
|
||
select(Task)
|
||
.where(Task.user_id == current_user.id)
|
||
.order_by(Task.created_at.desc())
|
||
.limit(limit)
|
||
.offset(offset)
|
||
)
|
||
return [TaskResponse.model_validate(t) for t in result.scalars().all()]
|
||
|
||
|
||
@router.get("/{public_id}", response_model=TaskResponse)
|
||
async def get_task(
|
||
public_id: str,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> TaskResponse:
|
||
"""Получить задачу по публичному ID. Возвращает 404 если задача чужая — не раскрываем факт существования."""
|
||
result = await db.execute(
|
||
select(Task).where(
|
||
Task.public_id == public_id,
|
||
Task.user_id == current_user.id, # ownership проверяется в одном запросе
|
||
)
|
||
)
|
||
task = result.scalar_one_or_none()
|
||
|
||
if task is None:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Задача не найдена")
|
||
|
||
return TaskResponse.model_validate(task)
|
||
|
||
|
||
@router.get("/{public_id}/text")
|
||
async def get_task_text(
|
||
public_id: str,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> dict:
|
||
"""Полный текст проверенного документа — для подсветки совпадений в интерфейсе.
|
||
|
||
Тот же текст, по которому worker-gpu считал position_start/position_end в
|
||
отчёте о плагиате (см. StagedWork.text_key — пишется в _stage_work при
|
||
извлечении, worker-indexer/app/tasks/index.py). Запись туда — best-effort
|
||
(сбой не валит саму проверку), поэтому текст доступен не для всех задач.
|
||
"""
|
||
result = await db.execute(
|
||
select(Task).where(Task.public_id == public_id, Task.user_id == current_user.id)
|
||
)
|
||
task = result.scalar_one_or_none()
|
||
if task is None:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Задача не найдена")
|
||
|
||
staged = (
|
||
await db.execute(select(StagedWork).where(StagedWork.task_id == task.id))
|
||
).scalar_one_or_none()
|
||
if staged is None or not staged.text_key:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Текст работы недоступен")
|
||
|
||
try:
|
||
obj = get_minio().get_object("staging", staged.text_key)
|
||
text = obj.read().decode("utf-8", errors="replace")
|
||
except Exception as e:
|
||
logger.warning(f"Не удалось прочитать текст задачи {public_id!r}: {e}")
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Текст работы недоступен") from e
|
||
|
||
return {"text": text}
|
||
|
||
|
||
@router.delete("/{public_id}", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def delete_task(
|
||
public_id: str,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> None:
|
||
"""Удалить задачу."""
|
||
result = await db.execute(
|
||
select(Task).where(
|
||
Task.public_id == public_id,
|
||
Task.user_id == current_user.id,
|
||
)
|
||
)
|
||
task = result.scalar_one_or_none()
|
||
|
||
if task is None:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Задача не найдена")
|
||
|
||
await db.delete(task)
|
||
await db.commit()
|