"""Queue document ingestion with a concurrency limit so large batches do not overwhelm the server. For batch uploads of thousands of files, tasks queue behind the semaphore and execute as slots free up — the server stays responsive for status polls and new requests throughout. """ import asyncio import logging from pathlib import Path from app.config import settings logger = logging.getLogger(__name__) _background_tasks: set[asyncio.Task] = set() # type: ignore[type-arg] _sem: asyncio.Semaphore | None = None def _ingest_semaphore() -> asyncio.Semaphore: global _sem if _sem is None: _sem = asyncio.Semaphore(max(1, settings.max_concurrent_ingests)) return _sem async def _ingest_worker(doc_id: str, file_path: Path) -> None: async with _ingest_semaphore(): from app.ingest import pipeline await pipeline.ingest_document(doc_id=doc_id, file_path=file_path) def schedule_ingest(doc_id: str, file_path: Path) -> None: """Fire-and-forget background ingestion for a committed Document row.""" task = asyncio.create_task(_ingest_worker(doc_id, file_path)) _background_tasks.add(task) def _done(t: asyncio.Task) -> None: # type: ignore[type-arg] _background_tasks.discard(t) try: t.result() except Exception: logger.exception("Background ingest task failed doc_id=%s", doc_id) task.add_done_callback(_done) logger.info("Queued ingestion doc=%s path=%s (queue_size=%d)", doc_id, file_path, len(_background_tasks)) def pending_count() -> int: """How many ingest tasks are currently queued or running.""" return len(_background_tasks)