Spaces:
Sleeping
Sleeping
| #!/usr/bin/env python3 | |
| """Re-ingest tenant documents into Qdrant after switching from FAISS. | |
| Usage (from repo root, with Qdrant running and VECTORSTORE_BACKEND=qdrant): | |
| python scripts/reingest_qdrant.py --tenant-id YOUR_TENANT | |
| Optional: mark all completed documents for a tenant and schedule ingest jobs via the API | |
| or call ingest pipeline directly for each file path in the database. | |
| """ | |
| from __future__ import annotations | |
| import argparse | |
| import asyncio | |
| import logging | |
| import sys | |
| from pathlib import Path | |
| ROOT = Path(__file__).resolve().parents[1] | |
| if str(ROOT) not in sys.path: | |
| sys.path.insert(0, str(ROOT)) | |
| logging.basicConfig(level=logging.INFO) | |
| logger = logging.getLogger("reingest_qdrant") | |
| async def _reingest_tenant(tenant_id: str, *, limit: int | None) -> int: | |
| from sqlalchemy import select | |
| from app.config import settings | |
| from app.db.database import get_session_factory, init_db | |
| from app.db.models import Document, IngestStatus | |
| from app.ingest.schedule import schedule_ingest | |
| from app.vectorstore.factory import reset_vectorstore | |
| backend = (settings.vectorstore_backend or "faiss").strip().lower() | |
| if backend != "qdrant": | |
| raise SystemExit( | |
| f"VECTORSTORE_BACKEND={backend!r}; set VECTORSTORE_BACKEND=qdrant in .env" | |
| ) | |
| await init_db() | |
| reset_vectorstore() | |
| factory = get_session_factory() | |
| n = 0 | |
| async with factory() as db: | |
| q = select(Document).where(Document.tenant_id == tenant_id) | |
| if limit is not None: | |
| q = q.limit(limit) | |
| result = await db.execute(q) | |
| docs = list(result.scalars().all()) | |
| for doc in docs: | |
| if doc.status != IngestStatus.complete: | |
| logger.warning("Skipping non-complete doc %s status=%s", doc.id, doc.status) | |
| continue | |
| path = Path(doc.file_path) | |
| if not path.is_file(): | |
| logger.warning("Missing file for doc %s: %s", doc.id, doc.file_path) | |
| continue | |
| schedule_ingest(doc_id=doc.id, file_path=path) | |
| n += 1 | |
| logger.info("Queued %d document(s) for re-ingest tenant=%s", n, tenant_id) | |
| return n | |
| def main() -> None: | |
| parser = argparse.ArgumentParser(description="Re-queue ingest jobs into Qdrant") | |
| parser.add_argument("--tenant-id", required=True, help="Tenant id to re-ingest") | |
| parser.add_argument("--limit", type=int, default=None, help="Max documents to queue") | |
| args = parser.parse_args() | |
| asyncio.run(_reingest_tenant(args.tenant_id, limit=args.limit)) | |
| if __name__ == "__main__": | |
| main() | |