Files
xiaoxia-saas/apps/worker/tasks.py
T
2026-06-17 18:14:35 +08:00

76 lines
2.6 KiB
Python

from datetime import datetime, timezone
import json
from .celery_app import celery_app
from packages.adapters.sqlalchemy_impl.session import SessionLocal
from packages.adapters.sqlalchemy_impl.ingest_job_repository import SQLAlchemyIngestJobRepository
from packages.adapters.sqlalchemy_impl.asset_repository import SQLAlchemyAssetRepository
from packages.domain import Asset, IngestJobStatus
@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:
session = SessionLocal()
try:
ingest_repo = SQLAlchemyIngestJobRepository(session)
asset_repo = SQLAlchemyAssetRepository(session)
job = ingest_repo.get(job_id)
if job is None:
return {"ok": False, "error": f"job {job_id} not found"}
job.status = IngestJobStatus.PROCESSING
job.updated_at = datetime.now(timezone.utc)
ingest_repo.update(job)
storage_key = job.storage_key
filename = storage_key.split("/")[-1]
lower_name = filename.lower()
if lower_name.endswith((".mp4", ".mov", ".avi", ".mkv")):
mime_type = "video/mp4"
elif lower_name.endswith((".mp3", ".wav", ".aac")):
mime_type = "audio/mpeg"
elif lower_name.endswith((".jpg", ".jpeg")):
mime_type = "image/jpeg"
elif lower_name.endswith((".png",)):
mime_type = "image/png"
else:
mime_type = "application/octet-stream"
asset = Asset.create(
workspace_id=job.workspace_id,
project_id=job.project_id,
library_id=job.library_id,
name=filename,
storage_key=storage_key,
mime_type=mime_type,
metadata={"source": "ingest_task"},
)
asset_repo.create(asset)
job.status = IngestJobStatus.COMPLETED
job.result_asset_id = asset.id
job.updated_at = datetime.now(timezone.utc)
ingest_repo.update(job)
return {"ok": True, "job_id": job.id, "asset_id": asset.id}
except Exception as e:
try:
ingest_repo = SQLAlchemyIngestJobRepository(session)
job = ingest_repo.get(job_id)
if job is not None:
job.status = IngestJobStatus.FAILED
job.error_message = str(e)
job.updated_at = datetime.now(timezone.utc)
ingest_repo.update(job)
except Exception:
pass
return {"ok": False, "job_id": job_id, "error": str(e)}
finally:
session.close()