367 lines
15 KiB
Python
367 lines
15 KiB
Python
"""Worker 启动时的初始化任务 — 孤儿任务清理等."""
|
||
|
||
import logging
|
||
from datetime import UTC
|
||
|
||
from celery.signals import worker_ready
|
||
from worker_app.db import SessionLocal
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
def cleanup_stale_running_with_session(repo, timeout_minutes: int) -> int:
|
||
"""清理超时未更新的 running GenerationTask(可注入 repo 的纯核心,便于单测)。
|
||
|
||
Returns:
|
||
清理的任务数量
|
||
"""
|
||
return len(cleanup_stale_running_with_session_ids(repo, timeout_minutes))
|
||
|
||
|
||
def cleanup_stale_running_with_session_ids(repo, timeout_minutes: int) -> list[tuple[str, str]]:
|
||
"""同 cleanup_stale_running_with_session,返回 [(task_id, celery_task_id), ...]。"""
|
||
fn = getattr(repo, "cleanup_stale_running_with_ids", None)
|
||
if fn is not None:
|
||
return fn(timeout_minutes)
|
||
# 旧仓储无 _with_ids 方法:降级为计数,无法撤销消息(执行前状态守卫兜底)
|
||
count = repo.cleanup_stale_running(timeout_minutes)
|
||
return [("", "") for _ in range(count)]
|
||
|
||
|
||
def cleanup_stale_pending_with_session(repo, timeout_minutes: int) -> int:
|
||
"""清理超时 pending GenerationTask(可注入 repo 的纯核心,便于单测)。
|
||
|
||
Returns:
|
||
清理的任务数量
|
||
"""
|
||
return len(cleanup_stale_pending_with_session_ids(repo, timeout_minutes))
|
||
|
||
|
||
def cleanup_stale_pending_with_session_ids(repo, timeout_minutes: int) -> list[tuple[str, str]]:
|
||
"""同 cleanup_stale_pending_with_session,返回 [(task_id, celery_task_id), ...]。"""
|
||
fn = getattr(repo, "cleanup_stale_pending_with_ids", None)
|
||
if fn is not None:
|
||
return fn(timeout_minutes)
|
||
count = repo.cleanup_stale_pending(timeout_minutes)
|
||
return [("", "") for _ in range(count)]
|
||
|
||
|
||
def _revoke_and_purge_stale_messages(items: list[tuple[str, str]]) -> int:
|
||
"""把清理掉的任务对应的 Celery 消息撤销并从 Redis 队列清除(#1714)。
|
||
|
||
防止「DB 已标 failed,但队列消息还在 → 重投执行 → 非法状态转换 → 半成品」。
|
||
失败不阻断清理流程(执行前状态守卫是第二道防线)。
|
||
"""
|
||
biz_ids = [tid for tid, _ in items if tid]
|
||
celery_ids = [cid for _, cid in items if cid]
|
||
if not biz_ids and not celery_ids:
|
||
return 0
|
||
try:
|
||
from worker_app.celery_app import celery_app as app
|
||
from worker_app.core.config import get_settings
|
||
|
||
from packages.shared.celery_orphan_guard import revoke_and_purge
|
||
|
||
broker_url = get_settings().broker_url
|
||
return revoke_and_purge(
|
||
app,
|
||
broker_url,
|
||
business_task_ids=biz_ids,
|
||
celery_task_ids=celery_ids,
|
||
)
|
||
except Exception as e: # noqa: BLE001
|
||
logger.error("撤销作废任务队列消息失败(执行前守卫仍会兜底): %s", e, exc_info=True)
|
||
return 0
|
||
|
||
|
||
# 孤儿任务超时阈值:running 任务超过此时间无进度更新则视为卡死。
|
||
# 依据:worker.generate_video 硬超时 time_limit=11 分钟,正常任务不可能超过;
|
||
# 20 分钟阈值覆盖硬超时 + 重试 + 余量,绝不误杀正常任务。
|
||
ORPHAN_TASK_TIMEOUT_MINUTES = 20
|
||
|
||
# Pending 任务超时阈值:任务创建后超过此时间仍未开始执行则判死。
|
||
# 注意区分 running 孤儿阈值(20 分钟):pending 是「排队等待」时间,
|
||
# 队列积压(如 20+ 转码任务)时视频生成可能正常排队较久,阈值必须放宽,
|
||
# 避免正常排队任务被误杀。队列隔离(#1714)后 generation 队列独占 worker,
|
||
# 理论上排队极短;保留 45 分钟作为兜底,覆盖 worker 短暂停止消费的场景。
|
||
PENDING_TASK_TIMEOUT_MINUTES = 45
|
||
|
||
|
||
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()
|
||
try:
|
||
repo = SQLAlchemyGenerationTaskRepository(session)
|
||
items = cleanup_stale_running_with_session_ids(repo, timeout_minutes)
|
||
finally:
|
||
session.close()
|
||
count = len(items)
|
||
if count > 0:
|
||
logger.warning("清理了 %d 个超时的孤儿 GenerationTask(超过 %d 分钟未更新)", count, timeout_minutes)
|
||
purged = _revoke_and_purge_stale_messages(items)
|
||
logger.info("孤儿任务对应队列消息撤销/清除完成: %d 条", purged)
|
||
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
|
||
|
||
from packages.adapters.sqlalchemy_impl.models import JobModel
|
||
from packages.domain.job import JobStatus
|
||
|
||
try:
|
||
session = SessionLocal()
|
||
cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes)
|
||
stale_jobs = (
|
||
session.query(JobModel)
|
||
.filter(
|
||
JobModel.status == JobStatus.RUNNING.value,
|
||
JobModel.updated_at < cutoff,
|
||
)
|
||
.all()
|
||
)
|
||
count = 0
|
||
stale_items: list[tuple[str, str]] = []
|
||
for model in stale_jobs:
|
||
stale_items.append((model.id, getattr(model, "celery_task_id", "") or ""))
|
||
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()
|
||
if count > 0:
|
||
_revoke_and_purge_generation(stale_items)
|
||
return count
|
||
except Exception as e:
|
||
logger.error("清理孤儿 Job 失败: %s", e, exc_info=True)
|
||
return 0
|
||
|
||
|
||
def cleanup_stale_pending_tasks(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover
|
||
"""清理数据库中卡在 pending 状态超时的 GenerationTask。
|
||
|
||
全局任务队列有 pending 数量上限,长期卡在 pending 的任务会占满队列,
|
||
导致新用户无法创建任务。通过 created_at 超时判断并标记为 failed。
|
||
|
||
Args:
|
||
timeout_minutes: 超时时间(分钟),默认 PENDING_TASK_TIMEOUT_MINUTES
|
||
|
||
Returns:
|
||
清理的任务数量
|
||
"""
|
||
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
|
||
SQLAlchemyGenerationTaskRepository,
|
||
)
|
||
|
||
session = SessionLocal()
|
||
try:
|
||
repo = SQLAlchemyGenerationTaskRepository(session)
|
||
items = cleanup_stale_pending_with_session_ids(repo, timeout_minutes)
|
||
count = len(items)
|
||
if count > 0:
|
||
logger.warning("清理了 %d 个超时的 pending GenerationTask(超过 %d 分钟未处理)", count, timeout_minutes)
|
||
purged = _revoke_and_purge_stale_messages(items)
|
||
logger.info("超时 pending 任务对应队列消息撤销/清除完成: %d 条", purged)
|
||
else:
|
||
logger.info("无超时 pending GenerationTask 需要清理")
|
||
return count
|
||
except Exception as e:
|
||
logger.error("清理超时 pending GenerationTask 失败: %s", e, exc_info=True)
|
||
return 0
|
||
finally:
|
||
session.close()
|
||
|
||
|
||
def _revoke_and_purge_generation(items: list[tuple[str, str]]) -> int:
|
||
"""撤销 Job 表孤儿任务(TTS/配音等)的队列消息,队列覆盖全部已知队列。"""
|
||
biz_ids = [tid for tid, _ in items if tid]
|
||
celery_ids = [cid for _, cid in items if cid]
|
||
if not biz_ids and not celery_ids:
|
||
return 0
|
||
try:
|
||
from worker_app.celery_app import celery_app as app
|
||
from worker_app.core.config import get_settings
|
||
|
||
from packages.shared.celery_orphan_guard import revoke_and_purge
|
||
|
||
broker_url = get_settings().broker_url
|
||
return revoke_and_purge(
|
||
app,
|
||
broker_url,
|
||
business_task_ids=biz_ids,
|
||
celery_task_ids=celery_ids,
|
||
queue_names=("generation", "transcode", "celery"),
|
||
)
|
||
except Exception as e: # noqa: BLE001
|
||
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)
|
||
pending_count = cleanup_stale_pending_tasks(PENDING_TASK_TIMEOUT_MINUTES)
|
||
total = gen_count + job_count + pending_count
|
||
if total > 0:
|
||
logger.warning(
|
||
"任务清理完成: 孤儿 GenerationTask=%d, 孤儿 Job=%d, 超时 pending=%d, 总计=%d",
|
||
gen_count,
|
||
job_count,
|
||
pending_count,
|
||
total,
|
||
)
|
||
return {"generation_tasks": gen_count, "jobs": job_count, "pending": pending_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)
|
||
|
||
|
||
@worker_ready.connect
|
||
def _recover_stuck_ingest_jobs_on_ready(sender, **kwargs): # pragma: no cover
|
||
"""Worker 启动完成后恢复卡死在 processing 的 ingest_job(#1714)。
|
||
|
||
容器重启/进程 OOM 导致 transcode 队列 unacked 消息未重投时,processing
|
||
ingest_job 会永久卡死。启动时扫描 processing 超 10 分钟的 job,CAS 重置
|
||
pending 并重新派单;Redis 锁保证同容器 generation/transcode 双 worker
|
||
只有一个执行恢复。旧消息若后来重投,ingest_asset 执行前守卫会丢弃。
|
||
"""
|
||
try:
|
||
from packages.application.ingest_orphan_cleanup import (
|
||
make_redis_recovery_lock,
|
||
recover_stuck_ingest_jobs_on_startup,
|
||
)
|
||
|
||
session = SessionLocal()
|
||
try:
|
||
recovered = recover_stuck_ingest_jobs_on_startup(
|
||
session,
|
||
lock_acquire=make_redis_recovery_lock(),
|
||
stuck_minutes=10,
|
||
)
|
||
finally:
|
||
session.close()
|
||
logger.info("Worker 启动 ingest 恢复完成,共重新派单 %d 个卡死任务", recovered)
|
||
except Exception as e: # noqa: BLE001 — 启动恢复失败不能阻断 worker 起服
|
||
logger.error("启动 ingest 恢复扫描失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True)
|
||
|
||
|
||
def recover_stale_voice_clones_on_startup(timeout_minutes: int = 10) -> int:
|
||
"""Worker 启动时恢复卡死在 processing 的音色克隆任务。
|
||
|
||
容器重启/进程 OOM 时 worker 中正在轮询的克隆任务会丢失,
|
||
voice_clone_profiles 永久卡在 processing 无兜底。启动时扫描
|
||
updated_at 超过 timeout_minutes 的 processing 记录,直接标记
|
||
为 failed(错误信息指引用户重试)。选择标 failed 而非重新派单,
|
||
因为 CosyVoice 侧的 voice_id 无法在无上下文下恢复轮询,重试需
|
||
用户确认后显式触发。
|
||
|
||
Args:
|
||
timeout_minutes: 判定卡死的阈值,默认 10 分钟
|
||
|
||
Returns:
|
||
恢复的记录数
|
||
"""
|
||
from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import (
|
||
SQLAlchemyVoiceCloneProfileRepository,
|
||
)
|
||
|
||
try:
|
||
session = SessionLocal()
|
||
try:
|
||
repo = SQLAlchemyVoiceCloneProfileRepository(session)
|
||
count = repo.cleanup_stale_processing(timeout_minutes)
|
||
finally:
|
||
session.close()
|
||
if count > 0:
|
||
logger.warning("启动时恢复了 %d 个卡死在 processing 的音色克隆(超时 %d 分钟)", count, timeout_minutes)
|
||
else:
|
||
logger.info("无卡死 processing 音色克隆需要恢复")
|
||
return count
|
||
except Exception as e:
|
||
logger.error("启动时音色克隆恢复扫描失败(beat 巡检仍会兜底): %s", e, exc_info=True)
|
||
return 0
|
||
|
||
|
||
@worker_ready.connect
|
||
def _recover_stuck_voice_clones_on_ready(sender, **kwargs):
|
||
"""Worker 启动完成后恢复卡死的音色克隆任务。"""
|
||
try:
|
||
recovered = recover_stale_voice_clones_on_startup()
|
||
logger.info("Worker 启动音色克隆恢复完成,共标记 %d 个卡死任务为 failed", recovered)
|
||
except Exception as e:
|
||
logger.error("启动音色克隆恢复失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True)
|
||
|
||
|
||
@worker_ready.connect
|
||
def _probe_gpu_encoder_on_ready(sender, **kwargs):
|
||
"""Worker 启动完成后探测 P4000 GPU NVENC 节点状态,打日志。"""
|
||
try:
|
||
from packages.shared.gpu_encoder import get_gpu_encoder
|
||
|
||
client = get_gpu_encoder()
|
||
if client is None:
|
||
logger.info(
|
||
"[gpu-encoder] disabled (ENABLE_GPU_ENCODE=false or endpoint not configured), using CPU libx264"
|
||
)
|
||
return
|
||
health = client.check_health()
|
||
if health.ready:
|
||
logger.info(
|
||
"[gpu-encoder] NVENC enabled: endpoint=%s gpu=%s worker=%s",
|
||
client.endpoint,
|
||
health.gpu_name,
|
||
health.worker,
|
||
)
|
||
else:
|
||
logger.warning(
|
||
"[gpu-encoder] configured but NOT ready: %s (endpoint=%s) — falling back to CPU",
|
||
health.error,
|
||
client.endpoint,
|
||
)
|
||
except Exception as e: # noqa: BLE001
|
||
logger.warning("[gpu-encoder] startup probe error (will retry on first job, CPU fallback): %s", e)
|