From 61b22759720374feab406b0a8c8cae75fe330b24 Mon Sep 17 00:00:00 2001 From: jze9 Date: Thu, 30 Jul 2026 13:04:49 +0500 Subject: [PATCH] =?UTF-8?q?fix:=20=D0=B3=D0=BE=D0=BD=D0=BA=D0=B0=20dispatc?= =?UTF-8?q?h-before-commit=20(search/plagiarism)=20+=20detached=20user=20?= =?UTF-8?q?=D0=B2=20notify?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Три бага, найденные при аудите прода: - search и documents диспатчили Celery-задачу ДО commit() — быстрый воркер читал задачу раньше, чем транзакция закоммичена, и падал «задача не найдена» (поиск не работал вовсе; плагиат спасала латентность скачивания из MinIO). Теперь коммитим до диспатча. - notify.send_task_done обращался к user.email/name и task.type ПОСЛЕ закрытия сессии → DetachedInstanceError, письма о завершении уходили в ретраи. Значения достаются внутри сессии. Co-Authored-By: Claude Opus 4.8 --- services/api/app/api/documents.py | 11 ++++++----- services/api/app/api/search.py | 11 ++++++----- services/worker-notifier/app/tasks/notify.py | 15 ++++++++++----- 3 files changed, 22 insertions(+), 15 deletions(-) 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)