diff --git a/services/api/app/api/documents.py b/services/api/app/api/documents.py index 99b6cf1..a0e1ec1 100644 --- a/services/api/app/api/documents.py +++ b/services/api/app/api/documents.py @@ -130,8 +130,13 @@ async def upload_for_plagiarism_check( "file_size_bytes": len(file_data), }, ) + task.queue_position = 1 + task.eta_seconds = 120 db.add(task) - await db.flush() + # Коммитим ДО диспатча — иначе воркер может прочитать задачу раньше коммита + # (гонка dispatch-before-commit). + await db.commit() + await db.refresh(task) celery_result = celery_app.send_task( "index.extract_and_check", @@ -140,11 +145,7 @@ async def upload_for_plagiarism_check( ) task.celery_task_id = celery_result.id - task.queue_position = 1 - task.eta_seconds = 120 - await db.commit() - await db.refresh(task) logger.info("Задача плагиата %s создана для пользователя %d", task.public_id, current_user.id) return TaskResponse.model_validate(task) diff --git a/services/api/app/api/search.py b/services/api/app/api/search.py index 5f3abb3..9d33c79 100644 --- a/services/api/app/api/search.py +++ b/services/api/app/api/search.py @@ -66,8 +66,13 @@ async def create_search_task( "category": data.category, }, ) + task.queue_position = 1 + task.eta_seconds = ETA_PER_POSITION_SECONDS db.add(task) - await db.flush() # получаем id и public_id + # Коммитим ДО диспатча: иначе быстрый воркер прочитает задачу раньше, чем + # транзакция закоммичена, и не найдёт её в БД (гонка dispatch-before-commit). + await db.commit() + await db.refresh(task) celery_result = celery_app.send_task( "gpu.search_semantic", @@ -77,11 +82,7 @@ async def create_search_task( ) task.celery_task_id = celery_result.id - task.queue_position = 1 - task.eta_seconds = ETA_PER_POSITION_SECONDS - await db.commit() - await db.refresh(task) logger.info("Задача поиска %s создана для пользователя %d", task.public_id, current_user.id) return TaskResponse.model_validate(task) diff --git a/services/worker-notifier/app/tasks/notify.py b/services/worker-notifier/app/tasks/notify.py index c664426..9cb8b2b 100644 --- a/services/worker-notifier/app/tasks/notify.py +++ b/services/worker-notifier/app/tasks/notify.py @@ -46,19 +46,24 @@ def send_task_done(self, task_id: str) -> dict: # Формируем краткое описание результата summary = _build_summary(task) + # Достаём значения ДО закрытия сессии — иначе объекты становятся + # detached и обращение к атрибутам падает DetachedInstanceError. + user_email = user.email + user_name = user.name + task_type = task.type sender = EmailSender() sender.send_task_done( - user_email=user.email, - user_name=user.name, + user_email=user_email, + user_name=user_name, task_id=task_id, - task_type=task.type, + task_type=task_type, summary=summary, app_url=settings.APP_URL, ) - logger.info(f"Уведомление отправлено пользователю {user.email!r} для задачи {task_id!r}") - return {"status": "sent", "email": user.email} + logger.info(f"Уведомление отправлено пользователю {user_email!r} для задачи {task_id!r}") + return {"status": "sent", "email": user_email} except Exception as exc: logger.error(f"Ошибка отправки уведомления для задачи {task_id!r}: {exc}", exc_info=True)