Files
yingqing/core/backend/apps/ai/tasks.py
T

147 lines
6.2 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from airshelf.celery import app
QUICK_QUEUE = "airshelf.quick"
@app.task(bind=True, max_retries=3)
def submit_ai_task(self, task_id: str) -> str:
return task_id
@app.task(bind=True, max_retries=5)
def poll_ai_task(self, task_id: str) -> str:
return task_id
@app.task(bind=True, max_retries=0)
def generate_standalone_image_task(self, task_id: str) -> str:
"""单张独立图的慢活(ARK 出图 ~30s)在 worker 内跑,Web 层不被占住。
幂等且失败自退费(见 run_standalone_image_task),故 max_retries=0,不向上抛重试。"""
from apps.ai.services import run_standalone_image_task
run_standalone_image_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0, queue=QUICK_QUEUE)
def extract_entities_task(self, task_id: str) -> str:
"""实体提取的慢活(豆包思考模型流式,可达数十秒)在 worker 内跑,Web 层秒回不被占住、不再 502。
幂等且失败自退费(见 run_extract_entities_task),故 max_retries=0,不向上抛重试。"""
from apps.ai.services import run_extract_entities_task
run_extract_entities_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0)
def run_video_digest_task(self, task_id: str) -> str:
"""视频提炼页的慢活(抽帧 + Gemini,可达 12 分钟)在 worker 内跑,离开页面也不中断。
幂等且失败自退费(见 run_team_digest_task),故 max_retries=0,不向上抛重试。"""
from apps.ai.video_digest import run_team_digest_task
run_team_digest_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0)
def run_video_replace_digest_task(self, task_id: str) -> str:
"""Worker:视频复刻第一道工序(参考视频 → 分镜稿,Gemini 半分钟起)在 worker 内跑。
商品和角色共用。失败自己把任务收成 FAILED(见 run_replace_digest),故 max_retries=0,不向上抛重试。"""
from apps.ai.models import AITask
from apps.ai.video_replace import run_replace_digest
task = AITask.objects.select_related("team", "created_by", "model_config").filter(id=task_id).first()
if task is not None:
run_replace_digest(task)
return task_id
@app.task(bind=True, max_retries=0)
def generate_base_asset_task(self, task_id: str) -> str:
"""基础资产(商品/人物/场景立绘)的慢出图在 worker 内跑,Web 层不被占住。
幂等且失败自退费(见 run_base_asset_task),故 max_retries=0,不向上抛重试。"""
from apps.ai.services import run_base_asset_task
run_base_asset_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0)
def generate_triview_task(self, task_id: str) -> str:
"""三视图(image_edit 以立绘为参考,慢)在 worker 内跑,Web 层不被占住。
幂等且失败自退费(见 run_triview_task),故 max_retries=0,不向上抛重试。"""
from apps.ai.services import run_triview_task
run_triview_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0)
def generate_model_triview_task(self, task_id: str) -> str:
"""团队模特库三视图:不绑定项目,成功更新 Model.triview_asset,失败释放预留积分。"""
from apps.ai.services import run_model_triview_task
run_model_triview_task(task_id=task_id)
return task_id
@app.task(bind=True, max_retries=0, queue=QUICK_QUEUE)
def poll_free_video_task(self, task_id: str, attempt: int = 0) -> str:
"""自由创作视频·worker 兜底轮询:每 30s 一次自重排(不依赖 celery beat),
最多 60 次(约 30 分钟);用尽后只停止 Worker 兜底,不把业务任务判失败。
finalize 幂等(POSTPROCESSING 认领),
与前端主动 poll 并存不双扣。轮询本身出错不重试(max_retries=0),下一次自重排继续。"""
from apps.ai.free_video import finalize_free_video
from apps.ai.models import AITask
from apps.ai.video_replace import advance_video_replace, is_video_replace_task
task = AITask.objects.select_related("model_config", "model_config__provider", "team").filter(id=task_id).first()
if task is None:
return task_id
try:
if is_video_replace_task(task):
task = advance_video_replace(task)
else:
task = finalize_free_video(task=task)
except Exception: # noqa: BLE001 — 单次轮询失败(网络抖动等)不终结任务,等下一轮
import logging
logging.getLogger(__name__).warning("poll_free_video_task %s attempt %s failed", task_id, attempt, exc_info=True)
# eager(本地联调/单测)下 apply_async 会内联立即执行,自重排=同步死循环 → 跳过,收尾交给前端主动 poll
from django.conf import settings as dj_settings
if getattr(dj_settings, "CELERY_TASK_ALWAYS_EAGER", False):
return task_id
# 复刻逐镜生成可能跑 8 段 × 数分钟,60 次(约 30 分钟)不够,拉到约 90 分钟。
poll_limit = 180 if is_video_replace_task(task) else 60
if task.status == AITask.Status.CREATED and attempt < poll_limit:
poll_free_video_task.apply_async(args=[task_id, attempt + 1], countdown=8 if attempt < 12 else 30)
elif task.status in (AITask.Status.SUBMITTED, AITask.Status.POLLING) and attempt < poll_limit:
poll_free_video_task.apply_async(args=[task_id, attempt + 1], countdown=30)
return task_id
@app.task(bind=True, max_retries=0, queue=QUICK_QUEUE)
def drain_asset_reviews_task(self) -> str:
"""全平台人脸审核兜底:未送审自动送火山、审核中自动拉绿/红盾。
自重排(无 celery beat);eager 下不自转,避免单测死循环。"""
from django.conf import settings as dj_settings
from django.core.cache import cache
from apps.assets.review import drain_platform_reviews
if cache.add("asset-review-drain-lock", "1", timeout=50):
try:
drain_platform_reviews()
except Exception: # noqa: BLE001 — 单轮失败不终结循环
import logging
logging.getLogger(__name__).warning("drain_asset_reviews_task failed", exc_info=True)
finally:
cache.delete("asset-review-drain-lock")
if getattr(dj_settings, "CELERY_TASK_ALWAYS_EAGER", False):
return "ok"
drain_asset_reviews_task.apply_async(countdown=45)
return "ok"