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

87 lines
3.7 KiB
Python

from datetime import datetime, timezone
import json
import random
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.adapters.sqlalchemy_impl.classification_job_repository import SQLAlchemyClassificationJobRepository
from packages.domain import Asset, AssetClassification, ClassificationJobStatus, IngestJobStatus
@celery_app.task(name="worker.healthcheck")
def healthcheck() -> dict:
return {"ok": True, "service": "worker"}
@celery_app.task(name="worker.classify_asset")
def classify_asset(job_id: str) -> dict:
session = SessionLocal()
try:
classification_repo = SQLAlchemyClassificationJobRepository(session)
asset_repo = SQLAlchemyAssetRepository(session)
job = classification_repo.get(job_id)
if job is None:
return {"ok": False, "error": f"classification job {job_id} not found"}
job.status = ClassificationJobStatus.PROCESSING
job.updated_at = datetime.now(timezone.utc)
classification_repo.update(job)
asset = asset_repo.get(job.asset_id)
if asset is None:
raise ValueError(f"asset {job.asset_id} not found")
name = asset.name.lower()
if any(token in name for token in ["food", "meal", "cook"]):
classification = AssetClassification.FOOD.value
elif any(token in name for token in ["person", "human", "portrait"]):
classification = AssetClassification.PERSON.value
elif any(token in name for token in ["music", "song", "audio"]):
classification = AssetClassification.MUSIC.value
elif any(token in name for token in ["product", "sku", "item"]):
classification = AssetClassification.PRODUCT.value
elif any(token in name for token in ["animal", "pet", "cat", "dog"]):
classification = AssetClassification.ANIMAL.value
elif any(token in name for token in ["sport", "run", "ball"]):
classification = AssetClassification.SPORT.value
elif any(token in name for token in ["tech", "phone", "device", "pc"]):
classification = AssetClassification.TECH.value
elif any(token in name for token in ["view", "travel", "mountain", "sea"]):
classification = AssetClassification.SCENIC.value
else:
classification = AssetClassification.OTHER.value
confidence = round(random.uniform(0.72, 0.96), 2)
asset.metadata = {
**asset.metadata,
"classification": classification,
"classification_confidence": confidence,
}
asset_repo.update(asset)
job.status = ClassificationJobStatus.COMPLETED
job.classification = classification
job.confidence = confidence
job.updated_at = datetime.now(timezone.utc)
classification_repo.update(job)
return {"ok": True, "job_id": job.id, "asset_id": asset.id, "classification": classification}
except Exception as e:
try:
classification_repo = SQLAlchemyClassificationJobRepository(session)
job = classification_repo.get(job_id)
if job is not None:
job.status = ClassificationJobStatus.FAILED
job.error_message = str(e)
job.updated_at = datetime.now(timezone.utc)
classification_repo.update(job)
except Exception:
pass
return {"ok": False, "job_id": job_id, "error": str(e)}
finally:
session.close()