import uuid

from app.core.celery_app import celery_app
from app.db.session import SessionLocal
from app.services.rag import index_knowledge_source
from app.ws.events import publish_job_event


@celery_app.task(name="rag.index_knowledge_source")
def index_knowledge_source_task(source_id: str) -> None:
    db = SessionLocal()
    try:
        publish_job_event(source_id, {"status": "indexing"})
        chunk_count = index_knowledge_source(db, uuid.UUID(source_id))
        publish_job_event(source_id, {"status": "indexed", "chunk_count": chunk_count})
    except Exception as exc:  # noqa: BLE001
        db.rollback()
        publish_job_event(source_id, {"status": "failed", "error": str(exc)})
        raise
    finally:
        db.close()
