大量优化全能创作

This commit is contained in:
Azmat@qq.com
2026-09-16 16:07:56 +08:00
parent a95cd159db
commit 7f5307d97c
9 changed files with 1155 additions and 65 deletions
+193
View File
@@ -6,6 +6,7 @@
from __future__ import annotations
from datetime import timedelta
import uuid
from django.core.cache import cache
from django.db import transaction
@@ -123,9 +124,16 @@ def sync_generating_message(message: CreationMessage) -> bool:
if message.kind != CreationMessage.Kind.GENERATING:
return False
payload = message.payload or {}
if payload.get("kind") == "video_segments":
return _sync_segmented_video_message(message)
task = _message_task(message)
if task is None:
return False
if payload.get("kind") == "video_merge" and task.status in (AITask.Status.SUBMITTED, AITask.Status.POLLING):
# Celery 不在线的本地环境同样可完成用户已确认的合并;此前从不在普通出片完成时调用。
run_segmented_video_merge(str(task.id))
task.refresh_from_db()
if task.status == AITask.Status.SUCCEEDED:
assets = _assets_from_task(task)
if not assets:
@@ -138,6 +146,191 @@ def sync_generating_message(message: CreationMessage) -> bool:
return False
def _sync_segmented_video_message(message: CreationMessage) -> bool:
"""同步全能创作的多段出片。
每段仍是普通 FREE_VIDEO 任务,故可沿用既有计费、轮询和资产落库;这里仅把它们聚合成
一张结果卡。所有分段完成后停在结果卡,绝不在此处触发 ffmpeg。
"""
from .free_video import finalize_free_video
from .models import AITask
payload = message.payload or {}
ids = [str(value) for value in payload.get("task_ids") or [] if value]
if not ids:
return False
by_id = {
str(task.id): task
for task in AITask.objects.filter(id__in=ids).select_related("model_config")
}
tasks = [by_id.get(task_id) for task_id in ids]
if any(task is None for task in tasks):
fail_generating_message(message, "分段视频任务不完整,请重新生成。")
return True
# 本地没有 worker 时,用户的会话轮询本身即可推进每个片段;已在 worker 中的任务则幂等返回。
refreshed = []
for task in tasks:
assert task is not None
if task.status in (AITask.Status.SUBMITTED, AITask.Status.POLLING):
try:
task = finalize_free_video(task=task)
except Exception: # noqa: BLE001 - 单段网络抖动不能使整组直接失败
task.refresh_from_db()
else:
task.refresh_from_db()
refreshed.append(task)
failed = next(
(task for task in refreshed if task.status in (AITask.Status.FAILED, AITask.Status.CANCELLED)),
None,
)
if failed is not None:
number = next((index + 1 for index, task in enumerate(refreshed) if task.id == failed.id), 1)
fail_generating_message(message, f"第 {number} 段生成失败:{failed.error_message or '请重试'}")
return True
if not all(task.status == AITask.Status.SUCCEEDED for task in refreshed):
return False
assets: list[dict] = []
for index, task in enumerate(refreshed, start=1):
task_assets = _assets_from_task(task)
if not task_assets:
return False
for asset in task_assets:
assets.append({**asset, "label": f"第 {index} 段"})
first = refreshed[0]
finish_generating_message(
message,
assets=assets,
meta={
**_meta_from_task(first, message),
"kind": "video_segments",
"segment_count": len(refreshed),
"total_duration": payload.get("total_duration") or "",
"needs_merge": True,
"segments": payload.get("segments") or [],
},
)
return True
def start_segmented_video_merge(*, conversation: CreationConversation, message: CreationMessage, user):
"""用户明确点击后才建立合并任务;这之前绝不下载片段或调用 ffmpeg。"""
from .models import AITask
payload = dict(message.payload or {})
if message.kind != CreationMessage.Kind.RESULT or not payload.get("needs_merge"):
raise ValueError("这条结果不需要合并")
state = str(payload.get("merge_state") or "")
if state in {"queued", "processing", "completed"}:
raise ValueError("这组片段已经在合并中或已合并")
task_ids = [str(value) for value in payload.get("task_ids") or [] if value]
source_tasks = list(AITask.objects.filter(id__in=task_ids).select_related("model_config"))
if len(source_tasks) != len(task_ids):
raise ValueError("分段视频不完整,无法合并")
model_config = source_tasks[0].model_config
merge_task = AITask.objects.create(
team=conversation.team,
created_by=user,
project=None,
task_type=AITask.Type.EXPORT,
status=AITask.Status.SUBMITTED,
model_config=model_config,
idempotency_key=f"omni-video-merge:{conversation.id}:{message.id}:{uuid.uuid4()}",
request_payload={
"feature": "omni_video_merge",
"source_message_id": str(message.id),
"source_task_ids": task_ids,
"prompt": payload.get("prompt") or "",
},
)
payload["merge_state"] = "queued"
payload["merge_task_id"] = str(merge_task.id)
message.payload = payload
message.save(update_fields=["payload", "updated_at"])
generating = append_message(
conversation,
role="assistant",
kind=CreationMessage.Kind.GENERATING,
payload={"task_id": str(merge_task.id), "kind": "video_merge", "prompt": payload.get("prompt") or ""},
task=merge_task,
)
return merge_task, generating
def run_segmented_video_merge(task_id: str) -> None:
"""后台执行用户已确认的合并。任务创建不等于执行;仅此函数调用 ffmpeg。"""
from .free_video import _store_free_video_media
from .models import AITask
from .services import asset_stable_url
from .video_replace import _concat_shot_media
from apps.assets.models import Asset
task = AITask.objects.select_related("team", "created_by").filter(id=task_id).first()
if task is None or task.task_type != AITask.Type.EXPORT:
return
source_message_id = str((task.request_payload or {}).get("source_message_id") or "")
source_message = CreationMessage.objects.filter(id=source_message_id).first()
if source_message is None:
return
try:
with transaction.atomic():
locked = AITask.objects.select_for_update().get(id=task.id)
if locked.status not in (AITask.Status.SUBMITTED, AITask.Status.POLLING):
return
locked.status = AITask.Status.POSTPROCESSING
locked.save(update_fields=["status", "updated_at"])
urls: list[str] = []
for source_id in (locked.request_payload or {}).get("source_task_ids") or []:
source_asset = (
Asset.objects.filter(origin_task_id=source_id, is_deleted=False, purged_at__isnull=True)
.prefetch_related("files")
.first()
)
url, _cover = asset_stable_url(source_asset)
if not url:
raise ValueError("有片段尚未准备好,暂时不能合并")
urls.append(url)
merged_bytes = _concat_shot_media(urls)
_store_free_video_media(task=locked, video_bytes=merged_bytes)
with transaction.atomic():
locked = AITask.objects.select_for_update().get(id=task.id)
if locked.status != AITask.Status.POSTPROCESSING:
return
locked.status = AITask.Status.SUCCEEDED
locked.completed_at = timezone.now()
locked.save(update_fields=["status", "completed_at", "updated_at"])
source = CreationMessage.objects.select_for_update().get(id=source_message.id)
source_payload = dict(source.payload or {})
source_payload["merge_state"] = "completed"
source.payload = source_payload
source.save(update_fields=["payload", "updated_at"])
except Exception as exc: # noqa: BLE001 - 合并失败不影响已完成的分段视频
import logging
logging.getLogger(__name__).exception("omni video merge failed for %s", task_id)
with transaction.atomic():
locked = AITask.objects.select_for_update().get(id=task.id)
if locked.status == AITask.Status.POSTPROCESSING:
locked.status = AITask.Status.FAILED
locked.error_code = "MergeError"
locked.error_message = str(exc)[:500]
locked.completed_at = timezone.now()
locked.save(update_fields=["status", "error_code", "error_message", "completed_at", "updated_at"])
source = CreationMessage.objects.select_for_update().filter(id=source_message.id).first()
if source is not None:
source_payload = dict(source.payload or {})
source_payload["merge_state"] = "failed"
source.payload = source_payload
source.save(update_fields=["payload", "updated_at"])
finally:
# 复用既有「生成中 → 结果/错误」转换,生成的合成片会成为独立结果卡。
fresh = AITask.objects.filter(id=task_id).first()
if fresh is not None:
sync_generating_for_task(fresh)
def sync_generating_messages(conversation: CreationConversation) -> int:
"""把一条会话里所有已结束的 GENERATING 回填。返回改了几条。"""
pending = list(