Files
xiaoxia-saas/packages/domain/generation_task.py
T
saas-backend-agent 48f6f49f7b
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m22s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m32s
AI Code Review / AI Code Review (pull_request) Successful in 6m51s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 10m53s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 0s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 30s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 29s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 3m46s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 4m59s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 5m33s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 10m57s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 11m4s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
feat(generation): #2024 defer video finalization until cover confirmed
Two-phase pipeline: worker renders and precomputes dedup metadata, but does
not create GeneratedVideo until user confirms cover in Step 5.

Domain
- Add GenerationTaskStatus.AWAITING_COVER (non-terminal)
- Add mark_awaiting_cover(): progress=100, clears error, leaves completed_at unset
- Transitions: running → {awaiting_cover, completed, failed, cancelled};
  awaiting_cover → {completed, failed, cancelled}
- _missing_ aliases: waiting_cover/video_ready/rendered/pending_cover
- Keep running→completed for backward compat / legacy paths

Worker (apps/worker/worker_app/tasks/generation.py)
- Replace _record_video_and_dedup with _precompute_render_metadata:
  calls compute_render_fingerprint_and_dedup, returns dict without DB writes
- Remove runtime batch-rerender decision (should_rerender_for_batch_dedup path);
  dedup now happens at finalize time against all finished records
- Persist precomputed metadata into task.extra_meta['rendered_output']
- Final status: mark_awaiting_cover instead of mark_completed

Dedup helpers (apps/worker/video_processing/dedup_helpers.py)
- New compute_render_fingerprint_and_dedup(video_path,...): local fingerprint
  + historical dedup + batch dedup, returns fully serializable dict
  (fingerprint_dict, fingerprint_chunks list, is_duplicate, duplicate_of, etc.)
- create_video_record_and_dedup() retained for compat/tests; now supports
  pre_dedup_result to reuse worker-precomputed data without re-reading video
- VideoFingerprint.from_dict() added to reconstruct from serialized form

Application layer (packages/application/generated_video_finalize.py)
- RenderedOutput dataclass + from_dict() for deserializing worker output
- finalize_generated_video(): creates GeneratedVideo, bulk-inserts
  VideoFingerprintChunk rows, commits; API-layer free

API service (apps/api/app/services/generation_finalize_service.py)
- GenerationFinalizeService.finalize_task(task_id, user_id, cover_url):
  * permission / existence check
  * idempotent: if GeneratedVideo already exists for this task, just ensure
    task is marked completed and return (safe for double-click)
  * status gate: only awaiting_cover (or legacy completed) accepted
  * cover resolution: explicit cover_url arg > task.cover_url
  * delegates to finalize_generated_video, marks task completed,
    clears extra_meta['rendered_output']

API endpoint (apps/api/app/api/routes/generation_tasks.py)
- POST /api/v1/generation/tasks/{task_id}/finalize
- Body: { cover_url?: string }; returns video_id/cover_url/file_url/is_duplicate
- Adjust confirm_generation fast path: mark_confirmed keeps task in
  awaiting_cover (don't auto-complete); historical 'completed' previews
  migrated back to awaiting_cover
- Preview-reuse accepts awaiting_cover tasks

Preview / task center
- GET /preview/{task_id} constructs lightweight _PreviewVideo from
  extra_meta.rendered_output when task is awaiting_cover
- Task center step map: awaiting_cover → 等待确认封面; status filter includes it

Tests
- tests/unit/test_finalize_generation.py: 13 new tests covering status
  transitions, mark_awaiting_cover, RenderedOutput parsing, finalize use case
  (success / missing metadata / cover fallback)
- test_generation_task.py updated for 6th status value
- Full suite: 15978 passed, 28 skipped (matches #2023 baseline, no regressions)
- black/isort/ruff clean

Out of scope (intentionally untouched)
- AI avatar pipeline uses separate AiAvatarRenderJob; finalize_job and
  /{job_id}/finalize are unchanged
- PR #2023 subtitle font scaling files (ass_subtitle_builder,
  video_filter_builder, subtitle_generator) not modified
- No DB migration: GenerationTaskModel.status is String(20) without CHECK
- Temp file cleanup: leave to existing periodic job
2026-09-24 11:29:36 +08:00

420 lines
15 KiB
Python
Executable File
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.
"""GenerationTask 领域模型 — 视频生成任务.
状态机:
pending → running → awaiting_cover → completed
↘ failed → pending (重试)
↘ cancelled
``awaiting_cover`` 表示渲染已完成、视频文件已上传、封面候选已就绪,
但用户尚未在 Step5 确认封面并点击「完成」,此时不创建 GeneratedVideo 成品记录。
用户调用 finalize 接口后才进入 ``completed`` 并正式入库。
"""
from __future__ import annotations
import json
import sys
from dataclasses import dataclass, field
from datetime import UTC, datetime
if sys.version_info >= (3, 11):
from enum import StrEnum
else:
from enum import Enum
class StrEnum(str, Enum):
pass
from uuid import uuid4
class GenerationTaskStatus(StrEnum):
"""生成任务状态枚举。"""
PENDING = "pending"
"""待处理(任务已创建,等待执行)"""
RUNNING = "running"
"""运行中(正在生成视频)"""
AWAITING_COVER = "awaiting_cover"
"""视频已渲染上传、封面候选已就绪,等待用户在 Step5 确认封面(finalize 前的中间态)"""
COMPLETED = "completed"
"""已完成(用户已确认封面,视频已正式入库)"""
FAILED = "failed"
"""失败(生成失败)"""
CANCELLED = "cancelled"
"""已取消(用户取消或系统取消)"""
@classmethod
def _missing_(cls, value: object) -> "GenerationTaskStatus":
"""兼容历史脏数据,避免枚举转换失败导致500。
- success/done/finished/complete → COMPLETED
- fail/error/err → FAILED
- process/processing/run/running → RUNNING
- cancel/canceled → CANCELLED
- 其他未知值 → PENDING(兜底,不阻塞业务)
"""
if isinstance(value, str):
normalized = value.strip().lower()
if normalized in ("done", "success", "finished", "complete", "completed"):
return cls.COMPLETED
if normalized in ("fail", "failed", "error", "err"):
return cls.FAILED
if normalized in ("process", "processing", "run", "running", "in_progress"):
return cls.RUNNING
if normalized in ("awaiting_cover", "waiting_cover", "video_ready", "rendered", "pending_cover"):
return cls.AWAITING_COVER
if normalized in ("cancel", "cancelled", "canceled"):
return cls.CANCELLED
return cls.PENDING
# 终态集合
TERMINAL_STATUSES = frozenset(
{GenerationTaskStatus.COMPLETED, GenerationTaskStatus.FAILED, GenerationTaskStatus.CANCELLED}
)
# 合法状态转换
_VALID_TRANSITIONS: dict[GenerationTaskStatus, set[GenerationTaskStatus]] = {
GenerationTaskStatus.PENDING: {
GenerationTaskStatus.RUNNING,
GenerationTaskStatus.FAILED,
GenerationTaskStatus.CANCELLED,
},
GenerationTaskStatus.RUNNING: {
GenerationTaskStatus.AWAITING_COVER,
GenerationTaskStatus.COMPLETED, # 兜底/测试兼容:允许直接完成;主路径走 awaiting_cover
GenerationTaskStatus.FAILED,
GenerationTaskStatus.CANCELLED,
},
GenerationTaskStatus.AWAITING_COVER: {
GenerationTaskStatus.COMPLETED,
GenerationTaskStatus.FAILED,
GenerationTaskStatus.CANCELLED,
},
GenerationTaskStatus.FAILED: {GenerationTaskStatus.PENDING}, # 重试回到 pending
}
@dataclass(slots=True)
class GenerationTask:
id: str
project_id: str
asset_library_id: str
strategy_id: str = ""
voice_library_id: str = ""
template_id: str = ""
asset_ids: list[str] = field(default_factory=list)
title_ids: list[str] = field(default_factory=list)
voice_ids: list[str] = field(default_factory=list)
status: GenerationTaskStatus = GenerationTaskStatus.PENDING
progress: float = 0.0
result_count: int = 0
error_message: str = ""
error_info: dict = field(default_factory=dict)
retry_count: int = 0
auto_retry_enabled: bool = False
auto_retry_max: int = 0
started_at: datetime | None = None
completed_at: datetime | None = None
source_edit_plan_id: str = ""
created_by_user_id: str = ""
asset_select_mode: str = ""
batch_id: str = ""
video_title: str = ""
resolution: str = ""
bgm_config: dict = field(default_factory=dict)
is_preview: bool = False
source_task_id: str = ""
celery_task_id: str = ""
output_width: int = 1280
output_height: int = 720
cover_url: str = ""
title_config: dict = field(default_factory=dict)
extra_meta: dict = field(default_factory=dict)
logs: str = "[]"
created_at: datetime = field(default_factory=lambda: datetime.now(UTC))
updated_at: datetime = field(default_factory=lambda: datetime.now(UTC))
@classmethod
def create(
cls,
project_id: str,
asset_library_id: str,
*,
strategy_id: str = "",
voice_library_id: str = "",
template_id: str = "",
asset_ids: list[str] | None = None,
title_ids: list[str] | None = None,
voice_ids: list[str] | None = None,
created_by_user_id: str = "",
source_edit_plan_id: str = "",
asset_select_mode: str = "",
batch_id: str = "",
video_title: str = "",
resolution: str = "",
bgm_config: dict | None = None,
auto_retry_enabled: bool = False,
auto_retry_max: int = 0,
is_preview: bool = False,
source_task_id: str = "",
output_width: int = 1280,
output_height: int = 720,
cover_url: str = "",
title_config: dict | None = None,
extra_meta: dict | None = None,
) -> "GenerationTask":
if not project_id.strip() and not template_id.strip():
raise ValueError("project_id 或 template_id 至少需要提供一个")
if not asset_library_id.strip() and not (asset_ids or title_ids or voice_ids):
raise ValueError("asset_library_id 或 asset_ids/title_ids/voice_ids 至少需要提供一个")
return cls(
id=uuid4().hex,
project_id=project_id.strip(),
asset_library_id=asset_library_id.strip(),
strategy_id=strategy_id.strip(),
voice_library_id=voice_library_id.strip(),
template_id=template_id.strip(),
asset_ids=list(asset_ids) if asset_ids else [],
title_ids=list(title_ids) if title_ids else [],
voice_ids=list(voice_ids) if voice_ids else [],
created_by_user_id=created_by_user_id.strip(),
source_edit_plan_id=source_edit_plan_id.strip(),
asset_select_mode=asset_select_mode,
batch_id=batch_id,
video_title=video_title.strip(),
resolution=resolution.strip(),
bgm_config=dict(bgm_config) if bgm_config else {},
auto_retry_enabled=auto_retry_enabled,
auto_retry_max=auto_retry_max,
is_preview=is_preview,
source_task_id=source_task_id,
output_width=output_width,
output_height=output_height,
cover_url=cover_url,
title_config=dict(title_config) if title_config else {},
extra_meta=dict(extra_meta) if extra_meta else {},
)
# ── 状态查询 ────────────────────────────────────────────────────────────
@property
def is_terminal(self) -> bool:
"""是否处于终态(completed / failed / cancelled)。"""
return self.status in TERMINAL_STATUSES
@property
def is_completed(self) -> bool:
"""是否已完成。"""
return self.status == GenerationTaskStatus.COMPLETED
@property
def is_failed(self) -> bool:
"""是否失败。"""
return self.status == GenerationTaskStatus.FAILED
@property
def is_running(self) -> bool:
"""是否运行中。"""
return self.status == GenerationTaskStatus.RUNNING
@property
def is_awaiting_cover(self) -> bool:
"""是否等待用户确认封面(渲染已完成、视频已上传、尚未 finalize 入库)。"""
return self.status == GenerationTaskStatus.AWAITING_COVER
# ── 状态转换 ────────────────────────────────────────────────────────────
def transition_to(self, new_status: GenerationTaskStatus | str) -> None:
"""执行状态转换。
Args:
new_status: 目标状态
Raises:
ValueError: 非法状态转换
"""
if isinstance(new_status, str):
try:
new_status = GenerationTaskStatus(new_status)
except ValueError as _e:
raise ValueError(f"无效状态: {new_status}") from _e
allowed = _VALID_TRANSITIONS.get(self.status, set())
if new_status not in allowed:
raise ValueError(
f"非法状态转换: {self.status.value} → {new_status.value},"
f"允许: {{{', '.join(sorted(s.value for s in allowed))}}}"
)
self.status = new_status
def mark_processing(self) -> None:
"""标记为处理中(pending → running)。
设置 started_at,清除 error_message。
Raises:
ValueError: 当前状态不允许转换到 running
"""
self.transition_to(GenerationTaskStatus.RUNNING)
self.started_at = datetime.now(UTC)
self.error_message = ""
def mark_awaiting_cover(self) -> None:
"""标记为等待确认封面(running → awaiting_cover)。
渲染与上传已完成、封面候选已就绪,等待用户在 Step5 选封面并点「完成」。
此时不创建 GeneratedVideo 成品记录;progress 置 100,completed_at 暂不设置
(finalize 完成入库时才真正结束任务)。
Raises:
ValueError: 当前状态不允许转换到 awaiting_cover
"""
self.transition_to(GenerationTaskStatus.AWAITING_COVER)
self.progress = 100.0
self.error_message = ""
def mark_completed(self, result_count: int = 1) -> None:
"""标记为已完成(awaiting_cover → completed,由 finalize 调用)。
设置 completed_at、progress=100.0、result_count,清除 error_message。
Args:
result_count: 入库的视频数量,默认为 1
Raises:
ValueError: 当前状态不允许转换到 completed
"""
self.transition_to(GenerationTaskStatus.COMPLETED)
self.completed_at = datetime.now(UTC)
self.progress = 100.0
self.result_count = result_count
self.error_message = ""
def mark_failed(self, error_message: str, error_info: dict | None = None) -> None:
"""标记为失败(pending / running → failed)。
设置 error_message、error_info、completed_at。
Args:
error_message: 错误信息
error_info: 结构化错误信息(error_type, stack_trace, stage, failed_at等)
Raises:
ValueError: 当前状态不允许转换到 failed
"""
self.transition_to(GenerationTaskStatus.FAILED)
self.error_message = error_message
self.completed_at = datetime.now(UTC)
if error_info is not None:
self.error_info = error_info
else:
self.error_info = {
"error_type": "UnknownError",
"message": error_message,
"failed_at": datetime.now(UTC).isoformat(),
}
def mark_cancelled(self) -> None:
"""标记为已取消(pending / running → cancelled)。
设置 completed_at。
Raises:
ValueError: 当前状态不允许转换到 cancelled
"""
self.transition_to(GenerationTaskStatus.CANCELLED)
self.completed_at = datetime.now(UTC)
def mark_confirmed(
self,
*,
cover_url: str = "",
extra_meta: dict | None = None,
output_width: int = 0,
output_height: int = 0,
title_config: dict | None = None,
) -> None:
"""将预览任务确认为正式产出。
预览渲染品质已与正式生成一致(1080p, CRF 23, medium),
确认时直接复用已有产物,无需重新渲染。
"""
self.is_preview = False
if cover_url:
self.cover_url = cover_url
if output_width > 0:
self.output_width = output_width
if output_height > 0:
self.output_height = output_height
if title_config:
self.title_config = dict(title_config)
if extra_meta:
self.extra_meta.update(extra_meta)
self.updated_at = datetime.now(UTC)
# ── 日志辅助 ────────────────────────────────────────────────────────────
_MAX_LOGS = 200
def append_log(self, stage: str, message: str, level: str = "INFO", **kwargs) -> None:
"""追加一条结构化日志到 logs 字段。
Args:
stage: 阶段名称(如 "接收任务"、"下载素材"、"渲染")
message: 日志消息
level: 日志级别(INFO / WARN / ERROR)
**kwargs: 额外字段(如 asset_id、duration 等)
"""
try:
entries = json.loads(self.logs) if self.logs else []
except (json.JSONDecodeError, TypeError):
entries = []
entry = {
"ts": datetime.now(UTC).isoformat(),
"level": level,
"stage": stage,
"message": message,
**kwargs,
}
entries.append(entry)
# 限制最多保留 _MAX_LOGS 条,防止字段过大
if len(entries) > self._MAX_LOGS:
entries = entries[-self._MAX_LOGS :]
self.logs = json.dumps(entries, ensure_ascii=False)
def get_logs(self) -> list[dict]:
"""解析 logs 字段为 list[dict]。"""
try:
return json.loads(self.logs) if self.logs else []
except (json.JSONDecodeError, TypeError):
return []
def mark_pending_from_failed(self) -> None:
"""从失败状态重置为待处理(用于重试)。
清除 error_message、error_info、started_at、completed_at、progress,
递增 retry_count。
Raises:
ValueError: 当前状态不是 failed
"""
if self.status != GenerationTaskStatus.FAILED:
raise ValueError(f"只有 failed 状态的任务可以重置为 pending,当前状态: {self.status.value}")
self.transition_to(GenerationTaskStatus.PENDING)
self.error_message = ""
self.error_info = {}
self.started_at = None
self.completed_at = None
self.progress = 0.0
self.result_count = 0
self.retry_count += 1