"""Ingestion endpoints.""" from __future__ import annotations import uuid from pathlib import Path from fastapi import APIRouter, Depends, Header, HTTPException, Response, status from app.api.schemas import ( AssetManifestEnvelope, IngestFolderRequest, IngestFolderResponse, KnowledgeIngestResponse, ) from app.api.security import require_service_api_key from app.config import settings from app.db.audit import record_audit from app.db.session import session_scope from app.ingestion.knowledge_ingest import KnowledgeIngestError, accept_knowledge_ingest from app.logging_config import get_logger logger = get_logger(__name__) router = APIRouter(tags=["ingestion"]) @router.post( "/knowledge-ingest", response_model=KnowledgeIngestResponse, dependencies=[Depends(require_service_api_key)], ) def knowledge_ingest( req: AssetManifestEnvelope, response: Response, idempotency_key: str | None = Header(default=None, alias="Idempotency-Key"), ) -> KnowledgeIngestResponse: """Accept a TeamHUB AssetManifest and queue LegacyHUB projection ingest.""" try: result = accept_knowledge_ingest(req, idempotency_key=idempotency_key) except KnowledgeIngestError as exc: raise HTTPException( status_code=exc.status_code, detail={"reason_code": exc.reason_code, "message": exc.message}, ) from exc response.status_code = ( status.HTTP_202_ACCEPTED if result.status == "accepted" else status.HTTP_200_OK ) return result @router.post( "/ingest/folder", response_model=IngestFolderResponse, deprecated=True, dependencies=[Depends(require_service_api_key)], ) def ingest_folder(req: IngestFolderRequest, response: Response) -> IngestFolderResponse: """Discover all PDFs under ``path`` and queue them for processing. The request returns immediately after the discovery pass. Per-document OCR / extraction / indexing happens asynchronously in Celery workers. Deprecated for TeamHUB integrations: use ``POST /api/v1/knowledge-ingest`` with an AssetManifest and object-storage reference instead. """ response.headers["Deprecation"] = "true" response.headers["Link"] = '; rel="successor-version"' if not settings.enable_folder_ingest: raise HTTPException( status_code=status.HTTP_410_GONE, detail="folder ingest is disabled; use POST /api/v1/knowledge-ingest", ) folder = Path(req.path) if not folder.exists() or not folder.is_dir(): raise HTTPException(status_code=400, detail=f"Folder not found: {req.path}") # Lazy import - keeps module load light. from app.ingestion.scanner import discover_documents # noqa: PLC0415 from app.workers.tasks import process_document # noqa: PLC0415 run_id = uuid.uuid4() discovered, queued, dups, invalid = 0, 0, 0, 0 for record in discover_documents(folder, recursive=req.recursive, force=req.force): discovered += 1 if record.duplicate and not req.force: dups += 1 continue if not record.document_id: invalid += 1 continue process_document.delay(str(record.document_id), str(run_id)) queued += 1 logger.info( "ingest.folder.queued", path=str(folder), discovered=discovered, queued=queued, skipped_duplicates=dups, invalid=invalid, run_id=str(run_id), ) with session_scope() as db: record_audit( db, action="ingest.folder.queued", actor="service:ingest-folder", entity_type="ingestion_run", entity_id=str(run_id), details={"discovered": discovered, "queued": queued, "skipped_duplicates": dups}, ) return IngestFolderResponse( run_id=run_id, discovered=discovered, queued=queued, skipped_duplicates=dups, invalid_files=invalid, )