"""Celery worker 在线探测 —— 生成类任务的提交前置闸。 设计决策(2026-06-10):无 worker 时**禁止提交**图片/视频生成,而非静默降级到前端轮询。 旧行为的风险:提交到 ARK 后浏览器一关就无人轮询,生成结果悬在云端无人取回入库、 额度预扣长期冻结——等于数据丢失。额度预检挡"钱不够",这道闸挡"环境不齐"。 探测用 control.ping 广播,结果带缓存:在线缓存 30s(别让每次提交都付 ping 开销), 离线只缓存 5s(worker 一启动尽快解封)。CELERY_TASK_ALWAYS_EAGER(测试/同步模式) 下任务就地执行不需要 worker,直接放行。 """ import time from django.conf import settings from rest_framework.exceptions import APIException class WorkerUnavailable(APIException): status_code = 503 default_detail = ( "后台任务服务(Celery worker)未运行:生成结果将无人取回入库,可能造成数据丢失。" "已暂停生成功能,请启动 worker 后重试(celery -A airshelf worker)。" ) default_code = "worker_unavailable" _OK_TTL = 30.0 _FAIL_TTL = 5.0 _cache = {"ok": False, "expires": 0.0} def _ping_workers() -> bool: from airshelf.celery import app as celery_app # limit=1:收到第一个 worker 回包立即返回,不傻等满 timeout return bool(celery_app.control.ping(timeout=1.0, limit=1)) def celery_worker_available() -> bool: if getattr(settings, "CELERY_TASK_ALWAYS_EAGER", False): return True now = time.monotonic() if now < _cache["expires"]: return _cache["ok"] try: ok = _ping_workers() except Exception: # noqa: BLE001 — broker 不可达等同于无 worker,一样要拦 ok = False _cache["ok"] = ok _cache["expires"] = now + (_OK_TTL if ok else _FAIL_TTL) return ok def require_worker() -> None: """生成类提交入口的前置检查:无 worker 直接 503,不让任务出门。""" if not celery_worker_available(): raise WorkerUnavailable()