53fb25efcf
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Check push changed paths (push) Successful in 19s
CI/CD Pipeline / Build Staging API Image (push) Successful in 41s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 48s
CI/CD Pipeline / Integration Tests (push) Successful in 3m10s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 3m17s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m30s
CI/CD Pipeline / Validate - Style (push) Successful in 4m17s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 59s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m33s
CI/CD Pipeline / Validate - Security (push) Successful in 7m12s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m38s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m13s
CI/CD Pipeline / Unit Tests (push) Successful in 10m11s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Staging E2E Tests (push) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (push) Failing after 26h14m3s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 26h24m21s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 26h19m47s
CI/CD Pipeline / PR Build Web Image (push) Failing after 26h23m44s
CI/CD Pipeline / PR Build API Image (push) Failing after 26h23m44s
CI/CD Pipeline / Deploy Production (push) Failing after 26h13m23s
CI/CD Pipeline / Build Production Web Image (push) Failing after 26h13m26s
CI/CD Pipeline / CI Gate (push) Failing after 26h13m25s
CI/CD Pipeline / Build Production API Image (push) Failing after 26h13m26s
CI/CD Pipeline / Canary Release to Production (push) Failing after 26h13m23s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Failing after 26h19m46s
CI/CD Pipeline / Frontend Lint (push) Failing after 26h23m37s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 26h23m45s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 26h19m46s
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
312 lines
12 KiB
Python
312 lines
12 KiB
Python
"""上传/转码链路(IngestJob + Asset)孤儿清理核心逻辑。
|
||
|
||
#1714:generation 链路有 cleanup_stale_running/pending 兜底,但上传链路
|
||
(ingest_jobs + assets)没有。worker 容器重启/进程 OOM 时,已 prefetch 的
|
||
celery 消息会丢失(transcode 队列 worker_prefetch_multiplier=1,消息预取后
|
||
宕机即丢失,Redis 队列里也不再存在),导致:
|
||
|
||
- ingest_jobs.status 永久卡 pending/processing
|
||
- assets.status 永久卡 processing/uploading(complete 阶段预建的占位)
|
||
|
||
本模块提供纯核心(session 注入,便于单测):超时阈值内无更新的记录
|
||
批量标终态(job→failed、asset→error),并返回 (job_id, celery_task_id)
|
||
列表供调用方 revoke + purge 残留队列消息。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from collections.abc import Callable
|
||
from datetime import UTC, datetime, timedelta
|
||
from typing import Any
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# ingest_job PROCESSING 超时阈值:ingest 任务包含下载 + ffprobe + HEVC 转码
|
||
# (1GB 视频约 10-20 分钟)+ 回传 OSS,正常任务可能跑 20-30 分钟;
|
||
# 60 分钟阈值覆盖大文件转码 + 抖动,绝不误杀正常任务。
|
||
INGEST_PROCESSING_TIMEOUT_MINUTES = 60
|
||
|
||
# ingest_job PENDING 超时阈值:transcode 队列 concurrency=1,队列积压时
|
||
# 正常排队可能较久;90 分钟覆盖 worker 短暂停消费 + 排队。
|
||
INGEST_PENDING_TIMEOUT_MINUTES = 90
|
||
|
||
# Asset 占位超时阈值:无关联 ingest_job 的孤儿占位(complete 预建后派单失败等),
|
||
# 阈值放宽到 120 分钟,避免与 ingest_job 生命周期错杀。
|
||
ASSET_ORPHAN_TIMEOUT_MINUTES = 120
|
||
|
||
_TERMINAL_JOB_STATUSES = ("failed", "completed")
|
||
_TERMINAL_ASSET_STATUSES = ("ready", "error", "deleted")
|
||
|
||
|
||
def _now() -> datetime:
|
||
return datetime.now(UTC)
|
||
|
||
|
||
def cleanup_stale_ingest_jobs(
|
||
session: Any,
|
||
*,
|
||
processing_timeout_minutes: int = INGEST_PROCESSING_TIMEOUT_MINUTES,
|
||
pending_timeout_minutes: int = INGEST_PENDING_TIMEOUT_MINUTES,
|
||
commit: bool = True,
|
||
) -> tuple[list[tuple[str, str]], list[str]]:
|
||
"""清理超时卡 pending/processing 的 ingest_jobs,并联动关联 asset。
|
||
|
||
Args:
|
||
session: SQLAlchemy session(或提供 query/commit 的鸭子类型)
|
||
processing_timeout_minutes: processing 状态超时阈值
|
||
pending_timeout_minutes: pending 状态超时阈值
|
||
commit: 是否提交事务
|
||
|
||
Returns:
|
||
(job_items, asset_ids)
|
||
- job_items: [(job_id, celery_task_id), ...] 供 revoke/purge
|
||
- asset_ids: 被联动标记为 error 的 asset id 列表
|
||
"""
|
||
from packages.adapters.sqlalchemy_impl.models import AssetModel, IngestJobModel
|
||
|
||
now = _now()
|
||
processing_cutoff = now - timedelta(minutes=processing_timeout_minutes)
|
||
pending_cutoff = now - timedelta(minutes=pending_timeout_minutes)
|
||
|
||
stale_jobs = (
|
||
session.query(IngestJobModel)
|
||
.filter(
|
||
IngestJobModel.status.in_(["pending", "processing"]),
|
||
(
|
||
(IngestJobModel.status == "processing") & (IngestJobModel.updated_at < processing_cutoff)
|
||
| (IngestJobModel.status == "pending") & (IngestJobModel.created_at < pending_cutoff)
|
||
),
|
||
)
|
||
.all()
|
||
)
|
||
|
||
job_items: list[tuple[str, str]] = []
|
||
asset_ids: list[str] = []
|
||
stale_asset_models: list[Any] = []
|
||
for job_model in stale_jobs:
|
||
ref_time = job_model.updated_at or job_model.created_at
|
||
if ref_time.tzinfo is None: # SQLite 读回 naive datetime 的防御
|
||
ref_time = ref_time.replace(tzinfo=UTC)
|
||
stale_minutes = int((now - ref_time).total_seconds() // 60)
|
||
job_model.status = "failed"
|
||
job_model.error_message = (
|
||
f"转码任务执行中断(超过超时阈值未更新,疑似 worker 重启/进程退出,已卡死 {stale_minutes} 分钟)"
|
||
)
|
||
job_model.updated_at = now
|
||
job_items.append((job_model.id, getattr(job_model, "celery_task_id", "") or ""))
|
||
if job_model.asset_id:
|
||
asset_ids.append(job_model.asset_id)
|
||
|
||
if asset_ids:
|
||
stale_asset_models = (
|
||
session.query(AssetModel)
|
||
.filter(
|
||
AssetModel.id.in_(asset_ids),
|
||
AssetModel.status.in_(["processing", "uploading"]),
|
||
)
|
||
.all()
|
||
)
|
||
for asset_model in stale_asset_models:
|
||
asset_model.status = "error"
|
||
asset_model.updated_at = now
|
||
|
||
if commit and (job_items or stale_asset_models):
|
||
session.commit()
|
||
|
||
if job_items:
|
||
logger.warning(
|
||
"[ingest-cleanup] 清理 %d 个超时 ingest_job(processing>%dm / pending>%dm),联动 %d 个 asset 标 error",
|
||
len(job_items),
|
||
processing_timeout_minutes,
|
||
pending_timeout_minutes,
|
||
len(stale_asset_models),
|
||
)
|
||
return job_items, [a.id for a in stale_asset_models]
|
||
|
||
|
||
def cleanup_orphan_processing_assets(
|
||
session: Any,
|
||
*,
|
||
timeout_minutes: int = ASSET_ORPHAN_TIMEOUT_MINUTES,
|
||
commit: bool = True,
|
||
) -> list[str]:
|
||
"""清理无 ingest_job 关联、超时卡 processing/uploading 的孤儿 asset 占位。
|
||
|
||
complete 阶段预建 asset 后若派单失败(或 direct 上传 complete 后
|
||
未触发 ingest),占位会永久卡住。这类 asset 没有对应 ingest_job,
|
||
只能按 created_at 超时兜底标 error。
|
||
"""
|
||
from packages.adapters.sqlalchemy_impl.models import AssetModel, IngestJobModel
|
||
|
||
cutoff = _now() - timedelta(minutes=timeout_minutes)
|
||
orphan_assets = (
|
||
session.query(AssetModel)
|
||
.outerjoin(IngestJobModel, IngestJobModel.asset_id == AssetModel.id)
|
||
.filter(
|
||
AssetModel.status.in_(["processing", "uploading"]),
|
||
AssetModel.created_at < cutoff,
|
||
IngestJobModel.id.is_(None),
|
||
)
|
||
.all()
|
||
)
|
||
for asset_model in orphan_assets:
|
||
asset_model.status = "error"
|
||
asset_model.updated_at = _now()
|
||
if commit and orphan_assets:
|
||
session.commit()
|
||
logger.warning("[ingest-cleanup] 清理 %d 个无 job 关联的超时孤儿 asset 占位", len(orphan_assets))
|
||
return [a.id for a in orphan_assets]
|
||
|
||
|
||
def revoke_stale_ingest_messages(
|
||
job_items: list[tuple[str, str]],
|
||
*,
|
||
celery_app_factory: Callable[[], Any] | None = None,
|
||
broker_url_factory: Callable[[], str] | None = None,
|
||
) -> int:
|
||
"""revoke + 物理清理 ingest 作废消息(transcode/celery 队列)。
|
||
|
||
消息可能已在 worker 宕机时丢失(队列里查不到),那也无害;
|
||
若消息还在(极端重复投递),物理清除防止重投执行。
|
||
失败不阻断清理(ingest_asset 的执行前状态守卫是第二道防线)。
|
||
"""
|
||
biz_ids = [jid for jid, _ in job_items if jid]
|
||
celery_ids = [cid for _, cid in job_items if cid]
|
||
if not biz_ids and not celery_ids:
|
||
return 0
|
||
try:
|
||
from packages.shared.celery_orphan_guard import revoke_and_purge
|
||
|
||
app = celery_app_factory() if celery_app_factory else None
|
||
broker_url = broker_url_factory() if broker_url_factory else ""
|
||
if app is None or not broker_url:
|
||
from worker_app.celery_app import celery_app as _app
|
||
from worker_app.core.config import get_settings
|
||
|
||
app = _app
|
||
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=("transcode", "celery"),
|
||
)
|
||
except Exception as e: # noqa: BLE001
|
||
logger.error("撤销作废 ingest 队列消息失败(执行前守卫仍会兜底): %s", e, exc_info=True)
|
||
return 0
|
||
|
||
|
||
# ── worker 启动恢复(#1714)──────────────────────────────────────────────
|
||
#
|
||
# task_acks_late=True 下,worker 崩溃/容器重启时未 ack 的消息理论上会在
|
||
# visibility_timeout 到期后重新投递;但 prefork 进程异常、部署窗口跨
|
||
# visibility 配置边界等场景仍可能留下卡在 processing 的 ingest_job
|
||
# (staging 实证:03:16 派单、03:45 置 processing 后 worker 重启,
|
||
# unacked 消息未重投,任务永久卡死)。启动时做一次显式恢复扫描兜底。
|
||
#
|
||
# 恢复策略:processing 超过 stuck_minutes(默认 10 分钟,部署中跨进程
|
||
# 交接的正常窗口 < 10 分钟,不会误抢别的 worker 正在执行的任务)的 job,
|
||
# CAS 重置为 pending 并重新 send_task;旧消息若后来重投,ingest_asset
|
||
# 的执行前守卫会把状态不匹配的旧 celery 消息丢弃。
|
||
|
||
|
||
def recover_stuck_ingest_jobs_on_startup(
|
||
session: Any,
|
||
*,
|
||
send_task: Callable[..., Any] | None = None,
|
||
update_celery_task_id: Callable[[str, str], None] | None = None,
|
||
lock_acquire: Callable[[], bool] | None = None,
|
||
stuck_minutes: int = 10,
|
||
commit: bool = True,
|
||
) -> int:
|
||
"""worker 启动时把卡在 processing 超时的 ingest_job 重新派单。
|
||
|
||
Args:
|
||
session: SQLAlchemy session
|
||
send_task: celery send_task 可调用(注入便于测试);不传则用 worker celery_app
|
||
update_celery_task_id: 回写新 celery task id 的回调(job_id, new_task_id)
|
||
lock_acquire: 分布式锁获取回调(多 worker 进程同时启动时只允许一个恢复);
|
||
返回 False 表示未抢到锁,本次跳过
|
||
stuck_minutes: processing 超过该分钟数视为卡死
|
||
|
||
Returns:
|
||
重新派单的 job 数
|
||
"""
|
||
if lock_acquire is not None and not lock_acquire():
|
||
logger.info("[ingest-recover] 未抢到恢复锁,跳过(另一进程正在恢复)")
|
||
return 0
|
||
|
||
from packages.adapters.sqlalchemy_impl.models import IngestJobModel
|
||
|
||
cutoff = _now() - timedelta(minutes=stuck_minutes)
|
||
stuck_jobs = (
|
||
session.query(IngestJobModel)
|
||
.filter(IngestJobModel.status == "processing", IngestJobModel.updated_at < cutoff)
|
||
.order_by(IngestJobModel.updated_at.asc())
|
||
.all()
|
||
)
|
||
|
||
if not stuck_jobs:
|
||
logger.info("[ingest-recover] 无卡死 processing ingest_job 需要恢复")
|
||
return 0
|
||
|
||
if send_task is None:
|
||
from worker_app.celery_app import celery_app as _app
|
||
|
||
send_task = _app.send_task
|
||
|
||
recovered = 0
|
||
for job_model in stuck_jobs:
|
||
# CAS:只有仍是 processing 才重置(并发/旧消息已回写终态时不碰)
|
||
updated = (
|
||
session.query(IngestJobModel)
|
||
.filter(IngestJobModel.id == job_model.id, IngestJobModel.status == "processing")
|
||
.update({"status": "pending", "error_message": "", "updated_at": _now()})
|
||
)
|
||
if not updated:
|
||
continue
|
||
try:
|
||
result = send_task("worker.ingest_asset", args=[job_model.id])
|
||
new_task_id = getattr(result, "id", "") or ""
|
||
except Exception as e: # noqa: BLE001
|
||
logger.error("[ingest-recover] 重新派单失败 job_id=%s: %s", job_model.id, e)
|
||
continue
|
||
if new_task_id:
|
||
job_model.celery_task_id = new_task_id
|
||
if update_celery_task_id is not None:
|
||
update_celery_task_id(job_model.id, new_task_id)
|
||
logger.warning(
|
||
"[ingest-recover] 卡死 ingest_job %s 已重置 pending 并重新派单 (new celery task=%s)",
|
||
job_model.id,
|
||
new_task_id,
|
||
)
|
||
recovered += 1
|
||
|
||
if commit and recovered:
|
||
session.commit()
|
||
logger.warning("[ingest-recover] 启动恢复完成,共重新派单 %d 个卡死 ingest_job", recovered)
|
||
return recovered
|
||
|
||
|
||
def make_redis_recovery_lock(lock_key: str = "ingest:recover:startup", ttl_seconds: int = 300):
|
||
"""构造基于 Redis SET NX 的恢复锁工厂(多 worker 进程互斥)。
|
||
|
||
返回一个无参 callable,调用时尝试抢锁:抢到返回 True,未抢到返回 False。
|
||
Redis 不可用时不阻断启动恢复(返回 True,恢复逻辑自身有 CAS 幂等保护)。
|
||
"""
|
||
|
||
def _acquire() -> bool:
|
||
try:
|
||
import redis as redis_lib
|
||
from worker_app.core.config import get_settings
|
||
|
||
client = redis_lib.Redis.from_url(get_settings().broker_url)
|
||
return bool(client.set(lock_key, "1", nx=True, ex=ttl_seconds))
|
||
except Exception as e: # noqa: BLE001
|
||
logger.warning("[ingest-recover] Redis 锁不可用,降级为无锁执行(CAS 兜底): %s", e)
|
||
return True
|
||
|
||
return _acquire
|