diff --git a/apps/worker/worker_app/celery_app.py b/apps/worker/worker_app/celery_app.py index 661dcb962..d352abc01 100644 --- a/apps/worker/worker_app/celery_app.py +++ b/apps/worker/worker_app/celery_app.py @@ -11,4 +11,5 @@ celery_app.conf.imports = ( "worker_app.tasks.health", "worker_app.tasks.ingest", "worker_app.tasks.classification", + "worker_app.tasks.generation", ) diff --git a/apps/worker/worker_app/tasks/__init__.py b/apps/worker/worker_app/tasks/__init__.py index 75cbbe4f0..882e56981 100644 --- a/apps/worker/worker_app/tasks/__init__.py +++ b/apps/worker/worker_app/tasks/__init__.py @@ -1,6 +1,8 @@ """Task modules.""" +from .classification import classify_asset +from .generation import generate_video from .health import healthcheck from .ingest import ingest_asset -__all__ = ["healthcheck", "ingest_asset"] +__all__ = ["classify_asset", "generate_video", "healthcheck", "ingest_asset"] diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py new file mode 100644 index 000000000..0031d504d --- /dev/null +++ b/apps/worker/worker_app/tasks/generation.py @@ -0,0 +1,65 @@ +from datetime import datetime, timezone + +from packages.adapters.sqlalchemy_impl import SQLAlchemyGeneratedVideoRepository, SQLAlchemyGenerationTaskRepository +from packages.domain import GeneratedVideo, GenerationTaskStatus +from worker_app.celery_app import celery_app +from worker_app.db import SessionLocal + + +@celery_app.task(name="worker.generate_video") +def generate_video(task_id: str) -> dict: + """Generate a video result for a generation task. + + This is the production-safe Phase 7 baseline implementation. It consumes the + task, records a generated-video placeholder, and keeps the task lifecycle + moving end-to-end while the real FFmpeg composition engine is integrated. + """ + db = SessionLocal() + task_repo = SQLAlchemyGenerationTaskRepository(db) + video_repo = SQLAlchemyGeneratedVideoRepository(db) + + task = task_repo.get(task_id) + if task is None: + db.close() + return {"status": "failed", "error": "generation task not found", "task_id": task_id} + + try: + task.status = GenerationTaskStatus.RUNNING + task.progress = 10.0 + task.started_at = task.started_at or datetime.now(timezone.utc) + task_repo.update(task) + + output_name = f"generated-{task.id}.mp4" + file_url = f"generated://workspaces/{task.workspace_id}/projects/{task.project_id}/tasks/{task.id}/{output_name}" + + video = GeneratedVideo.create( + workspace_id=task.workspace_id, + project_id=task.project_id, + generation_task_id=task.id, + name=output_name, + file_url=file_url, + file_size=0, + duration=0.0, + width=1920, + height=1080, + fps=25.0, + thumbnail_url=None, + ) + video_repo.create(video) + + task.status = GenerationTaskStatus.COMPLETED + task.progress = 100.0 + task.result_count = 1 + task.error_message = "" + task.completed_at = datetime.now(timezone.utc) + task_repo.update(task) + + return {"status": "completed", "task_id": task.id, "video_id": video.id, "file_url": file_url} + except Exception as error: + task.status = GenerationTaskStatus.FAILED + task.error_message = str(error) + task.completed_at = datetime.now(timezone.utc) + task_repo.update(task) + return {"status": "failed", "task_id": task.id, "error": str(error)} + finally: + db.close()