#!/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()