Files
xiaoxia-saas/apps/worker/tasks.py
T
2026-06-18 20:56:31 +08:00

72 lines
2.6 KiB
Python

from datetime import datetime, timezone
from app.config import get_settings
from app.core.storage import get_minio_service
from .celery_app import celery_app
from packages.adapters.sqlalchemy_impl.session import SessionLocal, build_session_factory
from packages.adapters.sqlalchemy_impl.generation_task_repository import SQLAlchemyGenerationTaskRepository
from packages.adapters.sqlalchemy_impl.generated_video_repository import SQLAlchemyGeneratedVideoRepository
from packages.domain import GeneratedVideo, GenerationTaskStatus
settings = get_settings()
if SessionLocal is None:
build_session_factory(settings.database_url)
@celery_app.task(name="worker.generate_video")
def generate_video(task_id: str) -> dict:
session = SessionLocal()
try:
task_repo = SQLAlchemyGenerationTaskRepository(session)
video_repo = SQLAlchemyGeneratedVideoRepository(session)
storage_service = get_minio_service()
task = task_repo.get(task_id)
if task is None:
return {"ok": False, "error": f"generation task {task_id} not found"}
task.status = GenerationTaskStatus.RUNNING
task.progress = 10.0
task.started_at = task.started_at or datetime.now(timezone.utc)
task_repo.update(task)
file_name = f"{task.id}.mp4"
storage_key = f"workspaces/{task.workspace_id}/projects/{task.project_id}/generated/{task.id}/{file_name}"
file_url = storage_service.get_url(storage_key)
video = GeneratedVideo.create(
workspace_id=task.workspace_id,
project_id=task.project_id,
generation_task_id=task.id,
name=file_name,
file_url=file_url,
file_size=1024,
duration=10.0,
width=1920,
height=1080,
fps=25.0,
)
video_repo.create(video)
task.status = GenerationTaskStatus.COMPLETED
task.progress = 100.0
task.result_count = 1
task.completed_at = datetime.now(timezone.utc)
task_repo.update(task)
return {"ok": True, "task_id": task.id, "video_id": video.id, "file_url": file_url}
except Exception as error:
try:
task_repo = SQLAlchemyGenerationTaskRepository(session)
task = task_repo.get(task_id)
if task is not None:
task.status = GenerationTaskStatus.FAILED
task.error_message = str(error)
task.completed_at = datetime.now(timezone.utc)
task_repo.update(task)
except Exception:
pass
return {"ok": False, "task_id": task_id, "error": str(error)}
finally:
session.close()