RICS / backend /core /document_library.py
StormShadow308's picture
Add demo documentation and Docker setup for v2 report generation system
aad7814
Raw
History Blame Contribute Delete
10.8 kB
"""Reference document persistence, delete, and re-ingest."""
from __future__ import annotations
import logging
import shutil
import threading
import time
from datetime import UTC, datetime
from pathlib import Path
from backend.config import settings
from backend.core import ingest
from backend.core.rag_store import TIER_REFERENCE, get_rag_store
from backend.core.report_session import (
UploadedDocument,
delete_document,
get_document,
list_documents,
save_document,
)
from backend.core.style_profile import invalidate_style_profile
from backend.utils import tenant_store
logger = logging.getLogger(__name__)
_reingest_lock = threading.Lock()
_reingest_running: set[str] = set()
def is_reingest_running(tenant_id: str) -> bool:
with _reingest_lock:
return tenant_id in _reingest_running
def recover_stale_processing_documents(tenant_id: str) -> int:
"""Reset orphaned ``processing`` rows when no re-ingest worker is active."""
if is_reingest_running(tenant_id):
return 0
recovered = 0
for _doc_id, doc in list_documents(tenant_id).items():
if doc.status == "processing":
doc.status = "complete"
doc.error = None
save_document(tenant_id, doc)
recovered += 1
if recovered:
logger.info(
"Recovered %d stale processing document(s) for tenant=%s",
recovered,
tenant_id,
)
return recovered
def recover_all_tenants_stale_processing() -> int:
"""On startup, clear processing flags left by a killed background worker."""
tenants_root = settings.data_dir_path / "tenants"
if not tenants_root.is_dir():
return 0
total = 0
for tenant_dir in sorted(tenants_root.iterdir()):
if not tenant_dir.is_dir():
continue
if not (tenant_dir / "compat_documents.json").is_file():
continue
total += recover_stale_processing_documents(tenant_dir.name)
return total
def reingest_progress(tenant_id: str) -> dict:
docs = list_documents(tenant_id)
counts = {"complete": 0, "processing": 0, "failed": 0, "pending": 0}
for doc in docs.values():
key = doc.status if doc.status in counts else "pending"
counts[key] = counts.get(key, 0) + 1
return {
"total": len(docs),
"running": is_reingest_running(tenant_id),
**counts,
}
def document_created_at_iso(doc: UploadedDocument) -> str:
"""Serialize ``created_at`` for API responses (ISO-8601 UTC)."""
ts = doc.created_at
if isinstance(ts, (int, float)) and ts > 0:
return datetime.fromtimestamp(ts, tz=UTC).isoformat()
if isinstance(ts, str) and ts.strip():
return ts.strip()
return datetime.now(tz=UTC).isoformat()
def persist_reference_file(
tenant_id: str,
document_id: str,
source_path: Path,
*,
original_filename: str,
) -> Path:
"""Copy an uploaded file into tenant storage for later re-ingest."""
suffix = source_path.suffix.lower() or Path(original_filename).suffix.lower()
dest = tenant_store.reference_upload_path(tenant_id, document_id, suffix)
shutil.copy2(source_path, dest)
return dest
def ingest_and_register(
tenant_id: str,
source_path: Path,
*,
original_filename: str,
document_id: str | None = None,
) -> UploadedDocument:
"""Ingest a reference file and record it in the document library."""
from backend.core.report_session import new_document_id
doc_id = document_id or new_document_id()
stored = persist_reference_file(
tenant_id, doc_id, source_path, original_filename=original_filename
)
chunks = ingest.ingest_reference(tenant_id, stored)
doc = UploadedDocument(
document_id=doc_id,
filename=original_filename or stored.name,
status="complete",
ingested_chunks=chunks,
storage_path=str(stored),
file_size=stored.stat().st_size if stored.is_file() else 0,
created_at=time.time(),
)
save_document(tenant_id, doc)
invalidate_style_profile(tenant_id)
return doc
def remove_reference_document(tenant_id: str, document_id: str) -> int:
"""Remove chunks, stored file, and library record. Returns chunks removed."""
doc = get_document(tenant_id, document_id)
if doc is None:
raise KeyError("Document not found")
removed = get_rag_store().remove_document(
tenant_id,
TIER_REFERENCE,
source_filename=doc.filename,
doc_id=f"reference:{doc.filename}",
)
if doc.storage_path:
path = Path(doc.storage_path)
if path.is_file():
path.unlink(missing_ok=True)
delete_document(tenant_id, document_id)
invalidate_style_profile(tenant_id)
return removed
def reingest_reference_document(tenant_id: str, document_id: str) -> UploadedDocument:
"""Re-process one stored reference file through the current pipeline."""
doc = get_document(tenant_id, document_id)
if doc is None:
raise KeyError("Document not found")
path = Path(doc.storage_path) if doc.storage_path else None
if path is None or not path.is_file():
raise FileNotFoundError("Source file is no longer on disk; cannot re-ingest.")
doc.status = "processing"
doc.error = None
save_document(tenant_id, doc)
get_rag_store().remove_document(
tenant_id,
TIER_REFERENCE,
source_filename=doc.filename,
doc_id=f"reference:{doc.filename}",
)
chunks = ingest.ingest_reference(tenant_id, path)
doc.status = "complete"
doc.error = None
doc.ingested_chunks = chunks
doc.file_size = path.stat().st_size
save_document(tenant_id, doc)
invalidate_style_profile(tenant_id)
return doc
def reingest_all_documents(
tenant_id: str,
*,
skip_document_ids: set[str] | None = None,
) -> dict:
skip = skip_document_ids or set()
docs = list_documents(tenant_id)
queued: list[str] = []
skipped_missing = 0
skipped_active = 0
for doc_id, doc in docs.items():
if doc_id in skip:
skipped_active += 1
continue
try:
logger.info("Re-ingesting %s for tenant=%s", doc.filename, tenant_id)
updated = reingest_reference_document(tenant_id, doc_id)
queued.append(doc_id)
logger.info(
"Re-ingested %s (%d chunks)",
updated.filename,
updated.ingested_chunks,
)
except FileNotFoundError:
skipped_missing += 1
except Exception as exc: # noqa: BLE001
doc.status = "failed"
doc.error = str(exc)
save_document(tenant_id, doc)
return {
"queued": len(queued),
"document_ids": queued,
"skipped_active": skipped_active,
"skipped_missing_file": skipped_missing,
"detail": (
f"Re-ingested {len(queued)} document(s); "
f"skipped {skipped_active} blocked and {skipped_missing} missing-file."
),
}
def schedule_reingest_all_documents(
tenant_id: str,
*,
skip_document_ids: set[str] | None = None,
) -> dict:
"""Queue a full-library re-ingest on a background thread (non-blocking HTTP).
Only the document currently being embedded is marked ``processing`` so a
server restart cannot strand the whole library in that state.
"""
skip = skip_document_ids or set()
docs = list_documents(tenant_id)
to_queue = [doc_id for doc_id in docs if doc_id not in skip]
with _reingest_lock:
if tenant_id in _reingest_running:
progress = reingest_progress(tenant_id)
return {
"queued": 0,
"document_ids": [],
"skipped_active": len(skip),
"skipped_missing_file": 0,
"reingest_running": True,
"progress": progress,
"detail": (
f"Re-ingest already running "
f"({progress['processing']} processing, "
f"{progress['complete']} ready)."
),
}
_reingest_running.add(tenant_id)
# Clear orphaned processing flags from a prior killed worker.
recover_stale_processing_documents(tenant_id)
def _worker() -> None:
try:
logger.info(
"Background re-ingest started for tenant=%s (%d documents)",
tenant_id,
len(to_queue),
)
reingest_all_documents(tenant_id, skip_document_ids=skip)
except Exception: # noqa: BLE001
logger.exception("Background re-ingest failed for tenant=%s", tenant_id)
recover_stale_processing_documents(tenant_id)
finally:
with _reingest_lock:
_reingest_running.discard(tenant_id)
logger.info("Background re-ingest finished for tenant=%s", tenant_id)
threading.Thread(
target=_worker,
name=f"reingest-{tenant_id}",
daemon=True,
).start()
progress = reingest_progress(tenant_id)
return {
"queued": len(to_queue),
"document_ids": to_queue,
"skipped_active": len(skip),
"skipped_missing_file": 0,
"reingest_running": True,
"progress": progress,
"detail": (
f"Re-ingest started in the background for {len(to_queue)} document(s). "
"Status updates every few seconds as each file completes."
),
}
def schedule_reingest_reference_document(tenant_id: str, document_id: str) -> dict:
"""Queue a single-document re-ingest on a background thread."""
doc = get_document(tenant_id, document_id)
if doc is None:
raise KeyError("Document not found")
path = Path(doc.storage_path) if doc.storage_path else None
if path is None or not path.is_file():
raise FileNotFoundError("Source file is no longer on disk; cannot re-ingest.")
doc.status = "processing"
doc.error = None
save_document(tenant_id, doc)
def _worker() -> None:
try:
reingest_reference_document(tenant_id, document_id)
except Exception as exc: # noqa: BLE001
failed = get_document(tenant_id, document_id)
if failed is not None:
failed.status = "failed"
failed.error = str(exc)
save_document(tenant_id, failed)
threading.Thread(
target=_worker,
name=f"reingest-{tenant_id}-{document_id[:8]}",
daemon=True,
).start()
return {
"queued": 1,
"document_ids": [document_id],
"detail": f"Re-ingest started in the background for {doc.filename}.",
}