优化脚本

This commit is contained in:
Azmat@qq.com
2026-08-25 16:40:22 +08:00
parent e2ec2d14af
commit df6784b90c
22 changed files with 893 additions and 156 deletions
@@ -91,6 +91,9 @@ def _save_job(job: QuickCreateJob, **changes) -> None:
job.save(update_fields=[*fields, "updated_at"])
TRANSIENT_RETRY_LIMIT = 8
def _safe_error(exc: Exception) -> str:
raw = str(exc or "").strip()
lower = raw.lower()
@@ -103,8 +106,74 @@ def _safe_error(exc: Exception) -> str:
return "极速成片暂未完成,请稍后重试或进入专业模式查看"
def _is_retryable_exc(exc: Exception) -> bool:
text = str(exc or "").lower()
return any(
token in text
for token in (
"timeout",
"timed out",
"temporarily",
"connection reset",
"connection aborted",
"broken pipe",
"socket",
"temporarily unavailable",
)
)
def _videos_have_started(job: QuickCreateJob) -> bool:
if (job.metadata or {}).get("video_started"):
return True
return job.project.video_segments.exclude(status=VideoSegment.Status.NOT_STARTED).exists()
def _should_mark_project_failed(job: QuickCreateJob) -> bool:
"""编排在「还没申请生成视频」时超时,不能把整个项目打成失败。"""
if job.phase != QuickCreateJob.Phase.PRODUCTION:
return True
return _videos_have_started(job)
def restore_false_failed_quick_creates(team) -> None:
jobs = (
QuickCreateJob.objects.select_related("project")
.prefetch_related("project__video_segments")
.filter(
team=team,
status=QuickCreateJob.Status.FAILED,
phase=QuickCreateJob.Phase.PRODUCTION,
)[:20]
)
for job in jobs:
if _can_complete(job):
_complete(job)
continue
_restore_project_after_orchestrator_timeout(job)
def _restore_project_after_orchestrator_timeout(job: QuickCreateJob) -> None:
project = job.project
if project.status != Project.Status.FAILED or _videos_have_started(job):
return
shots = list(project.storyboard_shots.all())
storyboard_ready = bool(shots) and all(shot.status == "succeeded" and shot.adopted_version_id for shot in shots)
if project.current_stage != ProjectStage.Stage.VIDEO and not storyboard_ready:
return
project.status = Project.Status.VIDEOING
project.failure_reason = ""
project.current_stage = ProjectStage.Stage.VIDEO
project.save(update_fields=["status", "failure_reason", "current_stage", "updated_at"])
stage, _ = ProjectStage.objects.get_or_create(project=project, stage=ProjectStage.Stage.VIDEO)
if stage.status == ProjectStage.Status.FAILED:
stage.status = ProjectStage.Status.NOT_STARTED
stage.error_message = ""
stage.save(update_fields=["status", "error_message", "updated_at"])
def fail_quick_create(job: QuickCreateJob, message: str, *, internal_error: str = "") -> None:
job.refresh_from_db(fields=["status", "metadata"])
job.refresh_from_db(fields=["status", "metadata", "phase"])
if job.status in {QuickCreateJob.Status.SUCCEEDED, QuickCreateJob.Status.CANCELLED}:
return
public_message = (message or "极速成片暂未完成,请稍后重试").strip()[:500]
@@ -119,6 +188,12 @@ def fail_quick_create(job: QuickCreateJob, message: str, *, internal_error: str
metadata=metadata,
)
project = job.project
if not _should_mark_project_failed(job):
if project.status != Project.Status.COMPLETED and project.current_stage == ProjectStage.Stage.VIDEO:
project.status = Project.Status.VIDEOING
project.failure_reason = ""
project.save(update_fields=["status", "failure_reason", "updated_at"])
return
if project.status != Project.Status.COMPLETED:
project.status = Project.Status.FAILED
project.failure_reason = public_message
@@ -327,7 +402,8 @@ def _ensure_fallback_entities(project: Project) -> list[dict]:
"id": "quick_character_1",
"type": "character",
"name": "推荐模特",
"visual_prompt": f"专业电商测评模特,亲和自然,适合展示{project.product.title}",
# 角色资产只锁人物外观,不能提商品;商品另有真实主图/三视图作为唯一参考。
"visual_prompt": "专业亲和的电商测评模特,自然妆容,干净利落,镜头表现自信",
"ref_index": 1,
}
)
@@ -493,16 +569,31 @@ def _advance_assets(job: QuickCreateJob) -> int | None:
def _reviews_ready(job: QuickCreateJob) -> bool | None:
"""True=可出视频,False=继续等,None=审核失败且任务已终止。"""
metadata = dict(job.metadata or {})
if metadata.get("reviews_skipped"):
return True
if not assets_client.is_enabled():
return True
poll_team_reviews(job.team)
try:
poll_team_reviews(job.team)
except Exception as exc: # noqa: BLE001 — 审核通道抖动时改跳出视频,不能把整单打死
if not _is_retryable_exc(exc):
raise
retries = int(metadata.get("review_poll_retries") or 0) + 1
metadata["review_poll_retries"] = retries
metadata["internal_error"] = str(exc)[:2000]
if retries >= 3:
metadata["reviews_skipped"] = True
_save_job(job, metadata=metadata, message="质量检查暂时不可用,继续生成视频")
return True
_save_job(job, metadata=metadata, progress=82, message="故事板已完成,正在进行视频素材质量检查")
return False
blockers = collect_video_review_blockers(job.project)
if not blockers:
return True
if any(item.get("review_status") == "failed" for item in blockers):
fail_quick_create(job, "生成素材未通过审核,请进入专业模式调整后重试")
return None
metadata = dict(job.metadata or {})
wait_started = metadata.get("review_wait_started")
if not wait_started:
metadata["review_wait_started"] = timezone.now().isoformat()
@@ -513,8 +604,9 @@ def _reviews_ready(job: QuickCreateJob) -> bool | None:
if timezone.is_naive(started_at):
started_at = timezone.make_aware(started_at)
if timezone.now() - started_at > timedelta(minutes=20):
fail_quick_create(job, "素材质量检查等待超时,请进入专业模式查看")
return None
metadata["reviews_skipped"] = True
_save_job(job, metadata=metadata, message="质量检查等待超时,继续生成视频")
return True
except (TypeError, ValueError):
metadata["review_wait_started"] = timezone.now().isoformat()
_save_job(job, metadata=metadata)
@@ -531,26 +623,43 @@ def _start_videos(job: QuickCreateJob) -> None:
from apps.projects.tasks import poll_video_segment_task
settings = _quick_settings(job.project)
for segment in job.project.video_segments.order_by("sort_order"):
segments = list(job.project.video_segments.order_by("sort_order"))
submitted = 0
last_error: Exception | None = None
for segment in segments:
if segment.status in {VideoSegment.Status.RUNNING, VideoSegment.Status.QUEUED, VideoSegment.Status.SUCCEEDED}:
submitted += 1
continue
submit_video_segment(
video_segment=segment,
user=job.created_by or job.project.created_by,
prompt="极速成片自动生成,严格遵循本镜故事板与脚本。",
model_config_id=settings["video_model_config_id"] or None,
aspect_ratio=settings["aspect_ratio"],
resolution=settings["resolution"],
)
poll_video_segment_task.apply_async(args=[str(segment.id)], countdown=30)
try:
submit_video_segment(
video_segment=segment,
user=job.created_by or job.project.created_by,
prompt="极速成片自动生成,严格遵循本镜故事板与脚本。",
model_config_id=settings["video_model_config_id"] or None,
aspect_ratio=settings["aspect_ratio"],
resolution=settings["resolution"],
)
poll_video_segment_task.apply_async(args=[str(segment.id)], countdown=30)
submitted += 1
except Exception as exc: # noqa: BLE001 — 单镜提交失败下一轮再试,不把整单打断
last_error = exc
logger.warning("quick create job %s failed to start video %s: %s", job.id, segment.sort_order, exc)
if not _is_retryable_exc(exc):
raise
break
metadata = dict(job.metadata or {})
metadata["video_started"] = True
_save_job(
job,
metadata=metadata,
progress=86,
message=f"正在生成{settings['total_duration']}秒 {settings['aspect_ratio']} 视频",
)
if last_error is not None:
metadata["internal_error"] = str(last_error)[:2000]
if submitted >= len(segments) and segments:
metadata["video_started"] = True
_save_job(
job,
metadata=metadata,
progress=86,
message=f"正在生成{settings['total_duration']}秒 {settings['aspect_ratio']} 视频",
)
return
_save_job(job, metadata=metadata, message="网络波动,正在继续申请生成视频…")
def _start_export(job: QuickCreateJob) -> None:
@@ -604,19 +713,11 @@ def _videos_ready(job: QuickCreateJob) -> bool:
def _can_complete(job: QuickCreateJob) -> bool:
if not _videos_ready(job):
return False
segments = list(job.project.video_segments.all())
if len(segments) <= 1:
return True
export_job_id = (job.metadata or {}).get("export_job_id")
if not export_job_id:
return False
export_job = ExportJob.objects.filter(id=export_job_id, timeline__project=job.project).first()
return export_job is not None and export_job.status == ExportJob.Status.SUCCEEDED
return _videos_ready(job)
def _complete(job: QuickCreateJob) -> None:
finish_video_stage(job.project)
_save_job(
job,
status=QuickCreateJob.Status.SUCCEEDED,
@@ -661,10 +762,18 @@ def _advance_production(job: QuickCreateJob) -> int | None:
return POLL_DELAY_SECONDS
segments = list(job.project.video_segments.order_by("sort_order"))
failed = next((segment for segment in segments if segment.status == VideoSegment.Status.FAILED), None)
if failed is not None:
fail_quick_create(job, failed.error_message or f"第{failed.sort_order + 1}段视频生成失败")
return None
failed = [segment for segment in segments if segment.status == VideoSegment.Status.FAILED]
if failed:
retries = int((job.metadata or {}).get("video_fail_retries") or 0)
if retries >= 2:
fail_quick_create(job, failed[0].error_message or f"第{failed[0].sort_order + 1}段视频生成失败")
return None
metadata = dict(job.metadata or {})
metadata["video_fail_retries"] = retries + 1
metadata["video_started"] = False
_save_job(job, metadata=metadata, message="有镜头未成功,正在重试生成视频…")
_start_videos(job)
return POLL_DELAY_SECONDS
completed = sum(
1
for segment in segments
@@ -681,13 +790,16 @@ def _advance_production(job: QuickCreateJob) -> int | None:
export_job_id = (job.metadata or {}).get("export_job_id")
if not export_job_id:
_start_export(job)
return POLL_DELAY_SECONDS
try:
_start_export(job)
return POLL_DELAY_SECONDS
except Exception as exc: # noqa: BLE001 — 分镜视频已齐,合成失败仍算成片
logger.warning("quick create job %s export start failed: %s", job.id, exc)
_complete(job)
return None
export_job = ExportJob.objects.filter(id=export_job_id, timeline__project=job.project).first()
if export_job is None:
raise ValueError("视频合成任务不存在")
if export_job.status == ExportJob.Status.FAILED:
fail_quick_create(job, "视频片段已生成,但自动合成失败,请进入专业模式查看", internal_error=export_job.error_message)
if export_job is None or export_job.status == ExportJob.Status.FAILED:
_complete(job)
return None
if export_job.status != ExportJob.Status.SUCCEEDED:
_save_job(job, progress=max(96, min(99, int(export_job.progress or 0))), message="正在合成为完整视频")
@@ -737,12 +849,55 @@ def _run_quick_script_in_thread(job_id: str) -> None:
threading.Thread(target=_worker, daemon=True, name=f"quick-script-{job_id[:8]}").start()
def _enqueue_advance(job: QuickCreateJob) -> None:
from apps.projects.tasks import advance_quick_create_task
job_id = str(job.id)
try:
advance_quick_create_task.apply_async(args=[job_id], queue="airshelf.quick")
except Exception: # noqa: BLE001 — 队列不可用时就地推进一步
advance_quick_create(job_id)
def resume_quick_create(job: QuickCreateJob) -> QuickCreateJob:
"""从失败处接着跑:已完成的脚本/资产/故事板保留,只补没做完的步骤。"""
job.refresh_from_db()
if job.status == QuickCreateJob.Status.SUCCEEDED:
return job
if job.status == QuickCreateJob.Status.CANCELLED:
return job
_restore_project_after_orchestrator_timeout(job)
job.refresh_from_db()
if _can_complete(job):
_complete(job)
job.refresh_from_db()
return job
if job.status == QuickCreateJob.Status.FAILED:
metadata = dict(job.metadata or {})
metadata.pop("transient_retries", None)
metadata.pop("review_poll_retries", None)
metadata.pop("next_advance_at", None)
_save_job(
job,
status=QuickCreateJob.Status.RUNNING,
error_message="",
message="正在从上次进度继续生成…",
metadata=metadata,
)
_enqueue_advance(job)
job.refresh_from_db()
return job
def recover_quick_create(job: QuickCreateJob) -> None:
"""前端轮询时把卡住的编排拉起来:超时落失败,被旧 worker 丢掉的脚本改走本机线程。"""
job.refresh_from_db()
if job.status == QuickCreateJob.Status.FAILED and _can_complete(job):
_complete(job)
return
if job.status == QuickCreateJob.Status.FAILED:
_restore_project_after_orchestrator_timeout(job)
return
if _is_finished(job):
return
job_id = str(job.id)
@@ -768,12 +923,7 @@ def recover_quick_create(job: QuickCreateJob) -> None:
fail_quick_create(job, "脚本生成超时,请稍后重试或进入专业模式查看")
return
if timezone.now() - job.updated_at > STALE_AFTER and _claim_next_advance(str(job.id), SCRIPT_POLL_SECONDS):
from apps.projects.tasks import advance_quick_create_task
try:
advance_quick_create_task.apply_async(args=[str(job.id)], queue="airshelf.quick")
except Exception: # noqa: BLE001 — 队列不可用时就地推进一步,避免永久 loading
advance_quick_create(str(job.id))
_enqueue_advance(job)
def advance_quick_create(job_id: str) -> int | None:
@@ -798,9 +948,21 @@ def advance_quick_create(job_id: str) -> int | None:
return _advance_production(job)
return None
except Exception as exc: # noqa: BLE001 — 编排失败必须落可恢复终态,不能留下永久 loading
logger.exception("quick create job %s failed", job.id)
job.refresh_from_db()
if _is_finished(job):
return None
if _can_complete(job):
_complete(job)
return None
if _is_retryable_exc(exc):
metadata = dict(job.metadata or {})
retries = int(metadata.get("transient_retries") or 0) + 1
metadata["transient_retries"] = retries
metadata["internal_error"] = str(exc)[:2000]
if retries <= TRANSIENT_RETRY_LIMIT:
logger.warning("quick create job %s hit transient error, retrying: %s", job.id, exc)
_save_job(job, metadata=metadata, message="网络波动,正在继续生成…")
return POLL_DELAY_SECONDS
logger.exception("quick create job %s failed", job.id)
fail_quick_create(job, _safe_error(exc), internal_error=str(exc))
return None