import logging from typing import Any from app.core.celery_app import celery_app logger = logging.getLogger(__name__) def safe_enqueue_generation_task( task: Any, generation_task_repository: Any, *, log_prefix: str = "[任务队列]", log_task_status: bool = False, ) -> bool: """安全入队:send_task 失败时自动把任务标记为 failed,避免留下 pending 僵尸任务。 Args: task: 生成任务对象,需有 id 属性和 mark_failed 方法 generation_task_repository: 任务仓储,用于更新状态 log_prefix: 日志前缀,便于区分调用来源 log_task_status: 成功日志中是否额外打印任务状态 Returns: True 表示入队成功,False 表示入队失败(已标记为 failed) """ try: celery_app.send_task("worker.generate_video", args=[task.id]) if log_task_status: logger.info( "%s 入队成功: task_id=%s, status=%s", log_prefix, task.id, task.status, ) else: logger.info("%s 入队成功: task_id=%s", log_prefix, task.id) return True except Exception as e: logger.error( "%s 入队失败,标记为失败: task_id=%s error=%s", log_prefix, task.id, e, exc_info=True, ) try: task.mark_failed(f"任务入队失败: {e}") generation_task_repository.update(task) except Exception as update_err: logger.error( "%s 入队失败后更新状态也失败: task_id=%s error=%s", log_prefix, task.id, update_err, exc_info=True, ) return False