Files
anti-plagiarism/services/api/app/api/search.py
jze9 61b2275972
All checks were successful
Deploy / deploy (push) Successful in 23s
fix: гонка dispatch-before-commit (search/plagiarism) + detached user в notify
Три бага, найденные при аудите прода:
- search и documents диспатчили Celery-задачу ДО commit() — быстрый воркер
  читал задачу раньше, чем транзакция закоммичена, и падал «задача не найдена»
  (поиск не работал вовсе; плагиат спасала латентность скачивания из MinIO).
  Теперь коммитим до диспатча.
- notify.send_task_done обращался к user.email/name и task.type ПОСЛЕ закрытия
  сессии → DetachedInstanceError, письма о завершении уходили в ретраи.
  Значения достаются внутри сессии.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-30 13:04:49 +05:00

89 lines
3.5 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Роутер семантического поиска источников."""
import logging
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.celery_app import celery_app
from app.core.rate_limiter import acquire_concurrent_slot, check_and_increment_limit
from app.core.security import get_current_verified_user
from app.database import get_db
from app.models.task import Task
from app.models.user import User
from app.schemas.tasks import SearchRequest, TaskResponse
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/search", tags=["search"])
ETA_PER_POSITION_SECONDS = 15
@router.post("/", response_model=TaskResponse, status_code=status.HTTP_202_ACCEPTED)
async def create_search_task(
data: SearchRequest,
# Поиск требует подтверждённого email — защита от массового abuse
current_user: User = Depends(get_current_verified_user),
db: AsyncSession = Depends(get_db),
) -> TaskResponse:
"""
Создать задачу поиска источников.
Возвращает сразу с task public_id и позицией в очереди.
Результат — через GET /tasks/{public_id} или WS /ws/tasks/{public_id}?token=JWT.
"""
# Проверяем лимит (атомарно, Lua-скрипт)
limit_result = await check_and_increment_limit(current_user.id, "search", current_user.plan)
if not limit_result["allowed"]:
raise HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail=(
f"Превышен дневной лимит поиска (тариф '{current_user.plan}'): "
f"{limit_result['current']}/{limit_result['limit']}. "
f"Сбросится {limit_result['reset_at']}."
),
headers={"Retry-After": "86400"},
)
# Проверяем лимит одновременных задач (атомарно)
slot_acquired = await acquire_concurrent_slot(current_user.id, current_user.plan)
if not slot_acquired:
raise HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail="Превышен лимит одновременных задач. Дождитесь завершения текущих.",
)
task = Task(
user_id=current_user.id,
type="search",
status="queued",
input_data={
"query": data.query,
"lang": data.lang,
"year_from": data.year_from,
"year_to": data.year_to,
"category": data.category,
},
)
task.queue_position = 1
task.eta_seconds = ETA_PER_POSITION_SECONDS
db.add(task)
# Коммитим ДО диспатча: иначе быстрый воркер прочитает задачу раньше, чем
# транзакция закоммичена, и не найдёт её в БД (гонка dispatch-before-commit).
await db.commit()
await db.refresh(task)
celery_result = celery_app.send_task(
"gpu.search_semantic",
args=[task.id, data.query],
kwargs={"lang": data.lang, "year_from": data.year_from, "year_to": data.year_to},
queue="queue.gpu",
)
task.celery_task_id = celery_result.id
await db.commit()
logger.info("Задача поиска %s создана для пользователя %d", task.public_id, current_user.id)
return TaskResponse.model_validate(task)