diff --git a/apps/api/app/api/routes/ingest_jobs.py b/apps/api/app/api/routes/ingest_jobs.py index c6a23bd90..54fab8fb1 100644 --- a/apps/api/app/api/routes/ingest_jobs.py +++ b/apps/api/app/api/routes/ingest_jobs.py @@ -4,6 +4,7 @@ from app.dependencies import get_ingest_job_repository from app.schemas.ingest_job import IngestJobResponse, SubmitIngestJobRequest from packages.adapters.in_memory import InMemoryIngestJobRepository from packages.application import SubmitIngestJobCommand, SubmitIngestJobUseCase +from apps.worker.worker_app.tasks.ingest import ingest_asset router = APIRouter() @@ -22,7 +23,9 @@ def submit_ingest_job( storage_key=request.storage_key, ) ) - # TODO: enqueue async worker task here + # Enqueue async worker task + ingest_asset.delay(job.id) + return IngestJobResponse( id=job.id, workspace_id=job.workspace_id, diff --git a/apps/worker/worker_app/tasks/__init__.py b/apps/worker/worker_app/tasks/__init__.py index e765b7666..75cbbe4f0 100644 --- a/apps/worker/worker_app/tasks/__init__.py +++ b/apps/worker/worker_app/tasks/__init__.py @@ -1 +1,6 @@ """Task modules.""" + +from .health import healthcheck +from .ingest import ingest_asset + +__all__ = ["healthcheck", "ingest_asset"] diff --git a/apps/worker/worker_app/tasks/ingest.py b/apps/worker/worker_app/tasks/ingest.py new file mode 100644 index 000000000..007f94e5f --- /dev/null +++ b/apps/worker/worker_app/tasks/ingest.py @@ -0,0 +1,80 @@ +from worker_app.celery_app import celery_app +from packages.domain import Asset, IngestJob, IngestJobStatus +from packages.adapters.in_memory import InMemoryAssetRepository, InMemoryIngestJobRepository + + +@celery_app.task(name="worker.healthcheck") +def healthcheck() -> dict: + return {"ok": True, "service": "worker"} + + +@celery_app.task(name="worker.ingest_asset") +def ingest_asset(job_id: str) -> dict: + """ + Ingest asset task. + + Steps: + 1. Fetch IngestJob from repository + 2. Extract metadata from storage_key (placeholder: mock metadata) + 3. Create Asset entity + 4. Update IngestJob status to COMPLETED + 5. Return result + """ + # TODO: Replace with real repository injection + job_repo = InMemoryIngestJobRepository() + asset_repo = InMemoryAssetRepository() + + job = job_repo.get(job_id) + if job is None: + return {"status": "failed", "error": "job not found"} + + try: + # Update job status to PROCESSING + job.status = IngestJobStatus.PROCESSING + job_repo.update(job) + + # Mock metadata extraction (in real implementation: use ffprobe, Pillow, etc.) + mime_type = "video/mp4" if job.storage_key.endswith(".mp4") else "image/jpeg" + metadata = { + "duration": 10.5, + "width": 1920, + "height": 1080, + "size_bytes": 1024000, + } + + # Extract filename from storage_key + filename = job.storage_key.split("/")[-1] + + # Create Asset + asset = Asset.create( + workspace_id=job.workspace_id, + project_id=job.project_id, + library_id=job.library_id, + name=filename, + storage_key=job.storage_key, + mime_type=mime_type, + metadata=metadata, + ) + asset_repo.create(asset) + + # Update job status to COMPLETED + job.status = IngestJobStatus.COMPLETED + job.result_asset_id = asset.id + job_repo.update(job) + + return { + "status": "completed", + "job_id": job.id, + "asset_id": asset.id, + } + except Exception as e: + # Update job status to FAILED + job.status = IngestJobStatus.FAILED + job.error_message = str(e) + job_repo.update(job) + + return { + "status": "failed", + "job_id": job.id, + "error": str(e), + } diff --git a/tests/integration/test_ingest_pipeline.py b/tests/integration/test_ingest_pipeline.py new file mode 100644 index 000000000..e6526a8fe --- /dev/null +++ b/tests/integration/test_ingest_pipeline.py @@ -0,0 +1,103 @@ +from packages.application import SubmitIngestJobCommand, SubmitIngestJobUseCase +from packages.adapters.in_memory import InMemoryAssetRepository, InMemoryIngestJobRepository +from packages.domain import Asset, IngestJob, IngestJobStatus + + +def simulate_ingest_asset(job_id: str, job_repo: InMemoryIngestJobRepository, asset_repo: InMemoryAssetRepository) -> dict: + """ + Simulate ingest asset logic without Celery. + This is the core business logic that would run inside the worker task. + """ + job = job_repo.get(job_id) + if job is None: + return {"status": "failed", "error": "job not found"} + + try: + # Update job status to PROCESSING + job.status = IngestJobStatus.PROCESSING + job_repo.update(job) + + # Mock metadata extraction + mime_type = "video/mp4" if job.storage_key.endswith(".mp4") else "image/jpeg" + metadata = { + "duration": 10.5, + "width": 1920, + "height": 1080, + "size_bytes": 1024000, + } + + # Extract filename from storage_key + filename = job.storage_key.split("/")[-1] + + # Create Asset + asset = Asset.create( + workspace_id=job.workspace_id, + project_id=job.project_id, + library_id=job.library_id, + name=filename, + storage_key=job.storage_key, + mime_type=mime_type, + metadata=metadata, + ) + asset_repo.create(asset) + + # Update job status to COMPLETED + job.status = IngestJobStatus.COMPLETED + job.result_asset_id = asset.id + job_repo.update(job) + + return { + "status": "completed", + "job_id": job.id, + "asset_id": asset.id, + } + except Exception as e: + # Update job status to FAILED + job.status = IngestJobStatus.FAILED + job.error_message = str(e) + job_repo.update(job) + + return { + "status": "failed", + "job_id": job.id, + "error": str(e), + } + + +def test_ingest_asset_pipeline(): + """Test the full ingest pipeline: submit job -> worker processes -> asset created.""" + job_repo = InMemoryIngestJobRepository() + asset_repo = InMemoryAssetRepository() + + # Submit ingest job + use_case = SubmitIngestJobUseCase(job_repo) + job = use_case.execute( + SubmitIngestJobCommand( + workspace_id="ws-1", + project_id="proj-1", + library_id="lib-1", + storage_key="uploads/test-video.mp4", + ) + ) + + assert job.status == IngestJobStatus.PENDING + assert job.result_asset_id == "" + + # Simulate worker task execution + result = simulate_ingest_asset(job.id, job_repo, asset_repo) + + assert result["status"] == "completed" + assert "asset_id" in result + + # Verify job was updated + updated_job = job_repo.get(job.id) + assert updated_job is not None + assert updated_job.status == IngestJobStatus.COMPLETED + assert updated_job.result_asset_id != "" + + # Verify asset was created + assets = asset_repo.list_by_library("lib-1") + assert len(assets) == 1 + assert assets[0].id == updated_job.result_asset_id + assert assets[0].name == "test-video.mp4" + assert assets[0].metadata["duration"] == 10.5