147 lines
5.6 KiB
Python
147 lines
5.6 KiB
Python
"""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
|