fix: гонка dispatch-before-commit (search/plagiarism) + detached user в notify
All checks were successful
Deploy / deploy (push) Successful in 23s
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>
This commit is contained in:
@@ -130,8 +130,13 @@ async def upload_for_plagiarism_check(
|
|||||||
"file_size_bytes": len(file_data),
|
"file_size_bytes": len(file_data),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
task.queue_position = 1
|
||||||
|
task.eta_seconds = 120
|
||||||
db.add(task)
|
db.add(task)
|
||||||
await db.flush()
|
# Коммитим ДО диспатча — иначе воркер может прочитать задачу раньше коммита
|
||||||
|
# (гонка dispatch-before-commit).
|
||||||
|
await db.commit()
|
||||||
|
await db.refresh(task)
|
||||||
|
|
||||||
celery_result = celery_app.send_task(
|
celery_result = celery_app.send_task(
|
||||||
"index.extract_and_check",
|
"index.extract_and_check",
|
||||||
@@ -140,11 +145,7 @@ async def upload_for_plagiarism_check(
|
|||||||
)
|
)
|
||||||
|
|
||||||
task.celery_task_id = celery_result.id
|
task.celery_task_id = celery_result.id
|
||||||
task.queue_position = 1
|
|
||||||
task.eta_seconds = 120
|
|
||||||
|
|
||||||
await db.commit()
|
await db.commit()
|
||||||
await db.refresh(task)
|
|
||||||
|
|
||||||
logger.info("Задача плагиата %s создана для пользователя %d", task.public_id, current_user.id)
|
logger.info("Задача плагиата %s создана для пользователя %d", task.public_id, current_user.id)
|
||||||
return TaskResponse.model_validate(task)
|
return TaskResponse.model_validate(task)
|
||||||
|
|||||||
@@ -66,8 +66,13 @@ async def create_search_task(
|
|||||||
"category": data.category,
|
"category": data.category,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
task.queue_position = 1
|
||||||
|
task.eta_seconds = ETA_PER_POSITION_SECONDS
|
||||||
db.add(task)
|
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(
|
celery_result = celery_app.send_task(
|
||||||
"gpu.search_semantic",
|
"gpu.search_semantic",
|
||||||
@@ -77,11 +82,7 @@ async def create_search_task(
|
|||||||
)
|
)
|
||||||
|
|
||||||
task.celery_task_id = celery_result.id
|
task.celery_task_id = celery_result.id
|
||||||
task.queue_position = 1
|
|
||||||
task.eta_seconds = ETA_PER_POSITION_SECONDS
|
|
||||||
|
|
||||||
await db.commit()
|
await db.commit()
|
||||||
await db.refresh(task)
|
|
||||||
|
|
||||||
logger.info("Задача поиска %s создана для пользователя %d", task.public_id, current_user.id)
|
logger.info("Задача поиска %s создана для пользователя %d", task.public_id, current_user.id)
|
||||||
return TaskResponse.model_validate(task)
|
return TaskResponse.model_validate(task)
|
||||||
|
|||||||
@@ -46,19 +46,24 @@ def send_task_done(self, task_id: str) -> dict:
|
|||||||
|
|
||||||
# Формируем краткое описание результата
|
# Формируем краткое описание результата
|
||||||
summary = _build_summary(task)
|
summary = _build_summary(task)
|
||||||
|
# Достаём значения ДО закрытия сессии — иначе объекты становятся
|
||||||
|
# detached и обращение к атрибутам падает DetachedInstanceError.
|
||||||
|
user_email = user.email
|
||||||
|
user_name = user.name
|
||||||
|
task_type = task.type
|
||||||
|
|
||||||
sender = EmailSender()
|
sender = EmailSender()
|
||||||
sender.send_task_done(
|
sender.send_task_done(
|
||||||
user_email=user.email,
|
user_email=user_email,
|
||||||
user_name=user.name,
|
user_name=user_name,
|
||||||
task_id=task_id,
|
task_id=task_id,
|
||||||
task_type=task.type,
|
task_type=task_type,
|
||||||
summary=summary,
|
summary=summary,
|
||||||
app_url=settings.APP_URL,
|
app_url=settings.APP_URL,
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info(f"Уведомление отправлено пользователю {user.email!r} для задачи {task_id!r}")
|
logger.info(f"Уведомление отправлено пользователю {user_email!r} для задачи {task_id!r}")
|
||||||
return {"status": "sent", "email": user.email}
|
return {"status": "sent", "email": user_email}
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error(f"Ошибка отправки уведомления для задачи {task_id!r}: {exc}", exc_info=True)
|
logger.error(f"Ошибка отправки уведомления для задачи {task_id!r}: {exc}", exc_info=True)
|
||||||
|
|||||||
Reference in New Issue
Block a user