"""Worker 启动时的初始化任务 — 孤儿任务清理等.""" import logging from celery.signals import worker_ready from worker_app.db import SessionLocal logger = logging.getLogger(__name__) # 孤儿任务超时阈值:渲染任务超过此时间未更新则视为卡死 ORPHAN_TASK_TIMEOUT_MINUTES = 10 def cleanup_orphan_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover """清理数据库中超时未更新的 running GenerationTask(孤儿任务)。 worker 重启或崩溃后,之前处于 running 状态的任务会变成孤儿任务, 一直卡在 running 不动。通过 updated_at 超时判断并标记为 failed。 Args: timeout_minutes: 超时时间(分钟),默认 10 分钟 Returns: 清理的任务数量 """ from packages.adapters.sqlalchemy_impl.generation_task_repository import ( SQLAlchemyGenerationTaskRepository, ) try: session = SessionLocal() repo = SQLAlchemyGenerationTaskRepository(session) count = repo.cleanup_stale_running(timeout_minutes) session.close() if count > 0: logger.warning("清理了 %d 个超时的孤儿 GenerationTask(超过 %d 分钟未更新)", count, timeout_minutes) else: logger.info("无孤儿 GenerationTask 需要清理") return count except Exception as e: logger.error("清理孤儿 GenerationTask 失败: %s", e, exc_info=True) return 0 def cleanup_stale_jobs(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover """清理数据库中超时未更新的 running Job(孤儿任务)。 与 cleanup_orphan_tasks 配合,同时清理 Job 表和 GenerationTask 表。 Returns: 清理的任务数量 """ from datetime import datetime, timedelta, timezone from packages.adapters.sqlalchemy_impl.models import JobModel from packages.domain.job import JobStatus try: session = SessionLocal() cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes) stale_jobs = ( session.query(JobModel) .filter( JobModel.status == JobStatus.RUNNING.value, JobModel.updated_at < cutoff, ) .all() ) count = 0 for model in stale_jobs: model.status = JobStatus.FAILED.value model.error_message = f"任务执行中断(超过 {timeout_minutes} 分钟未更新)" count += 1 if count > 0: session.commit() logger.warning("清理了 %d 个超时的孤儿 Job(超过 %d 分钟未更新)", count, timeout_minutes) else: logger.info("无孤儿 Job 需要清理") session.close() return count except Exception as e: logger.error("清理孤儿 Job 失败: %s", e, exc_info=True) return 0 def cleanup_all_stale_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> dict: # pragma: no cover """统一清理所有超时的孤儿任务。 同时清理 GenerationTask 和 Job 两类表。 Returns: {"generation_tasks": int, "jobs": int} """ gen_count = cleanup_orphan_tasks(timeout_minutes) job_count = cleanup_stale_jobs(timeout_minutes) total = gen_count + job_count if total > 0: logger.warning( "孤儿任务清理完成: GenerationTask=%d, Job=%d, 总计=%d", gen_count, job_count, total, ) return {"generation_tasks": gen_count, "jobs": job_count} @worker_ready.connect def _on_worker_ready(sender, **kwargs): # pragma: no cover """Worker 启动完成后执行 — 清理孤儿任务。""" logger.info("Worker 启动完成,开始清理孤儿 running 任务...") result = cleanup_all_stale_tasks() total = result["generation_tasks"] + result["jobs"] logger.info("Worker 启动清理完成,共清理 %d 个孤儿任务", total)