"""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} # True / False / None(inspect 没答上来,不当成「未部署」) _registered_task_cache: dict[str, tuple[bool | None, float]] = {} _queue_cache: dict[str, tuple[bool | None, float]] = {} def _task_name_listed(task_name: str, task_names) -> bool: """Celery inspect.registered() 可能给全名、短名,或带 [rate] 后缀。""" short = task_name.rsplit(".", 1)[-1] for raw in task_names or []: name = str(raw).split("[", 1)[0].strip() if name in {task_name, short} or name.endswith("." + short): return True return False def _inspect_task_registered(task_name: str) -> bool | None: """True 已注册 / False 明确没有 / None 问不到(超时或空应答)。""" try: from airshelf.celery import app as celery_app registered = celery_app.control.inspect(timeout=2.5).registered() except Exception: # noqa: BLE001 — inspect 失败不等于任务没部署 return None if not registered: return None return any(_task_name_listed(task_name, names) for names in registered.values()) def _ping_workers() -> bool: from airshelf.celery import app as celery_app # limit=1:收到第一个 worker 回包立即返回。本机连火山 Redis 往返常超过 1s, # 超时会被误判成 worker 未运行,生成入口整页锁死。 return bool(celery_app.control.ping(timeout=3.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 _inspect_queue_consumed(queue_name: str) -> bool | None: """True 有 worker 在听该队列 / False 在线 worker 都不听 / None 问不到。""" try: from airshelf.celery import app as celery_app active_queues = celery_app.control.inspect(timeout=2.5).active_queues() except Exception: # noqa: BLE001 — inspect 失败不能当成「没人听」 return None if not active_queues: return None found_any_queue = False for queues in active_queues.values(): for item in queues or []: found_any_queue = True name = item.get("name") if isinstance(item, dict) else str(item) if name == queue_name: return True if found_any_queue: return False return None def worker_consumes_queue(queue_name: str) -> bool | None: """线上 worker 默认只听 celery;极速成片发到 airshelf.quick 时用来决定要不要本机兜底。""" if getattr(settings, "CELERY_TASK_ALWAYS_EAGER", False): return True now = time.monotonic() cached = _queue_cache.get(queue_name) if cached and now < cached[1]: return cached[0] available = _inspect_queue_consumed(queue_name) ttl = _OK_TTL if available else _FAIL_TTL _queue_cache[queue_name] = (available, now + ttl) return available def require_worker() -> None: """生成类提交入口的前置检查:无 worker 直接 503,不让任务出门。""" if not celery_worker_available(): raise WorkerUnavailable() def require_worker_task(task_name: str) -> None: """确认在线 Worker 已加载指定任务,避免新功能被旧进程静默丢弃。 Worker 在线不代表它已经加载了最新 tasks.py。Celery 收到未注册任务会直接丢弃消息, 页面只能看到永久「生成中」。提交极速成片等新编排前,额外确认任务名已注册。 inspect 超时/空应答不能当成「未部署」——本机常见 ping 通、inspect 1s 问不到, 会误报「正在升级」。只有明确看到已注册列表且没有该任务时才拦截。 """ if getattr(settings, "CELERY_TASK_ALWAYS_EAGER", False): return require_worker() now = time.monotonic() cached = _registered_task_cache.get(task_name) if cached and now < cached[1]: available = cached[0] else: available = _inspect_task_registered(task_name) ttl = _OK_TTL if available else _FAIL_TTL _registered_task_cache[task_name] = (available, now + ttl) if available is False: exc = WorkerUnavailable() exc.detail = "一键成片后台任务尚未加载,请重启 Celery worker 后再试。" raise exc