Files
xiaoxia-saas/apps/worker/worker_app/tasks/_startup.py
T
xiaoxia 0ef4e1d633
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / Validate - Migration (alembic) (push) Successful in 51s
CI/CD Pipeline / Validate - Type Check (mypy) (push) Successful in 59s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 2m49s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 3m26s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 4m55s
CI/CD Pipeline / Validate - Code Quality (push) Successful in 5m47s
CI/CD Pipeline / Integration Tests (push) Successful in 1m56s
CI/CD Pipeline / Unit Tests (push) Successful in 9m42s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Build Staging API Image (push) Successful in 12m15s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 34s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 37s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 2m27s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 4m39s
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
feat: 渲染失败检测+任务超时机制 (#1219)
- 新增 video_validation 模块(moov atom 检测 + ffprobe 验证 + 退出码映射)
- render_adapter 渲染后自动校验输出再上传 OSS
- 三个渲染任务加 10 分钟 soft_time_limit
- Worker 启动时清理 GenerationTask + Job 两张表的孤儿任务
- 21 个新单元测试,全量 13713 测试通过
2026-08-02 15:34:22 +08:00

114 lines
4.0 KiB
Python

"""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)