All checks were successful
Deploy / deploy (push) Successful in 23s
Три бага, найденные при аудите прода: - 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>
89 lines
3.5 KiB
Python
89 lines
3.5 KiB
Python
"""Роутер семантического поиска источников."""
|
||
|
||
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)
|