Spaces:
Runtime error
Runtime error
| """Legacy frontend API compatibility layer backed by v2 services.""" | |
| from __future__ import annotations | |
| import logging | |
| import shutil | |
| import tempfile | |
| import time | |
| from pathlib import Path | |
| from typing import Any | |
| from fastapi import ( | |
| APIRouter, | |
| BackgroundTasks, | |
| Depends, | |
| File, | |
| Header, | |
| HTTPException, | |
| Query, | |
| Response, | |
| UploadFile, | |
| status, | |
| ) | |
| from starlette.concurrency import run_in_threadpool | |
| from pydantic import BaseModel, Field | |
| from backend.api.deps import get_current_tenant | |
| from backend.api import security | |
| from backend.config import settings | |
| from backend.core import template_discoverer | |
| from backend.core.content_similarity import scan_similar_content | |
| from backend.core import document_library | |
| from backend.core.legacy_generation import run_legacy_generation | |
| from backend.core.notes_extract import extract_notes_from_upload | |
| from backend.core.notes_routing import route_lines_to_l3_sections | |
| from backend.core import photo_store | |
| from backend.core.photo_store import ALLOWED_CONTENT_TYPES | |
| from backend.core.rag_store import TIER_REFERENCE, get_rag_store | |
| from backend.core.report_assembler import to_docx | |
| from backend.core.report_session import ( | |
| ReportSession, | |
| create_report_session, | |
| get_document, | |
| list_documents, | |
| load_session, | |
| save_session, | |
| ) | |
| logger = logging.getLogger(__name__) | |
| router = APIRouter(tags=["legacy-compat"]) | |
| _GROUP_LABELS = { | |
| "A": "A — About the inspection", | |
| "B": "B — Overall opinion and summary of ratings", | |
| "C": "C — About the property", | |
| "D": "D — Outside the property", | |
| "E": "E — Inside the property", | |
| "F": "F — Services", | |
| "G": "G — Grounds (including shared areas for flats)", | |
| "H": "H — Issues for your legal advisers", | |
| "I": "I — Risks", | |
| "J": "J — Energy matters", | |
| "K": "K — Surveyor's declaration", | |
| "L": "L — What to do now", | |
| "M": "M — Description of the RICS Home Survey - Level 3 service and terms of engagement", | |
| "N": "N — Typical house diagram", | |
| } | |
| class LegacyGenerateBody(BaseModel): | |
| template_id: str | |
| template_ids: list[str] = Field(default_factory=list) | |
| bullets: list[str] = Field(default_factory=list) | |
| bullets_by_section: dict[str, list[str]] = Field(default_factory=dict) | |
| mode: str = "generate" | |
| interference_level: str | None = None | |
| retrieval_level: str = "paragraph" | |
| force_regenerate: bool = True | |
| strict_uploaded_only: bool = False | |
| reference_document_ids: list[str] = Field(default_factory=list) | |
| draft_paragraph: str | None = None | |
| class BatchStatusBody(BaseModel): | |
| document_ids: list[str] = Field(default_factory=list) | |
| class RouteNotesBody(BaseModel): | |
| lines: list[str] = Field(default_factory=list, description="Raw note lines to route into L3 sections.") | |
| class SurveyClassifyBody(BaseModel): | |
| document_ids: list[str] = Field(default_factory=list) | |
| force_refresh_corpus: bool = False | |
| class SurveyApplyBody(BaseModel): | |
| items: list[dict[str, Any]] = Field(default_factory=list) | |
| class PhotoAiSelectionBody(BaseModel): | |
| photo_ids: list[str] = Field(default_factory=list) | |
| class SimilarContentBody(BaseModel): | |
| text: str = "" | |
| section_code: str | None = None | |
| peer_sections: dict[str, str] = Field(default_factory=dict) | |
| limit: int = 8 | |
| min_relevance_percent: float = 28.0 | |
| exclude_document_ids: list[str] = Field(default_factory=list) | |
| class SectionTextPatch(BaseModel): | |
| text: str | |
| def _section_group(code: str) -> str: | |
| code = (code or "").strip().upper() | |
| if len(code) == 1: | |
| return code | |
| return code[0] if code else "A" | |
| def _resolve_tenant_token( | |
| authorization: str = Header(default=""), | |
| access_token: str = Query(default=""), | |
| ) -> str: | |
| token = "" | |
| if authorization.lower().startswith("bearer "): | |
| token = authorization.split(" ", 1)[1].strip() | |
| elif access_token.strip(): | |
| token = access_token.strip() | |
| if not token: | |
| raise HTTPException(status_code=401, detail="Missing bearer token") | |
| payload = security.decode_token(token) | |
| tenant_id = payload.get("sub") | |
| if not tenant_id: | |
| raise HTTPException(status_code=401, detail="Token missing subject") | |
| return tenant_id | |
| def _require_session(tenant_id: str, report_id: str) -> ReportSession: | |
| session = load_session(tenant_id, report_id) | |
| if session is None: | |
| raise HTTPException(status_code=404, detail="Report not found") | |
| return session | |
| def _legacy_photo_url(report_id: str, section_id: str, photo_id: str) -> str: | |
| return f"/reports/{report_id}/sections/{section_id}/photos/{photo_id}" | |
| def _photo_rows(tenant_id: str, session: ReportSession, section_id: str) -> list[dict]: | |
| rows = photo_store.list_section_photos(tenant_id, session.draft_id, section_id) | |
| return [ | |
| { | |
| "photo_id": p.id, | |
| "original_filename": p.filename, | |
| "url": _legacy_photo_url(session.report_id, section_id, p.id), | |
| "selected_for_ai": p.selected_for_ai, | |
| } | |
| for p in rows | |
| ] | |
| async def _run_generation(tenant_id: str, report_id: str, body: LegacyGenerateBody) -> None: | |
| """Run blocking RAG + LLM generation off the async event loop.""" | |
| await run_in_threadpool(run_legacy_generation, tenant_id, report_id, body) | |
| def _ingest_upload_file( | |
| tenant_id: str, | |
| tmp_path: Path, | |
| original_filename: str, | |
| ): | |
| return document_library.ingest_and_register( | |
| tenant_id, | |
| tmp_path, | |
| original_filename=original_filename, | |
| ) | |
| # ── Upload & documents ─────────────────────────────────────────────────────── | |
| async def upload_batch( | |
| files: list[UploadFile] = File(...), | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| from backend.ingest.zip_extract import extract_reference_documents | |
| allowed_ext = {".pdf", ".docx", ".docm", ".doc", ".zip"} | |
| items: list[dict] = [] | |
| accepted = 0 | |
| rejected = 0 | |
| for upload in files: | |
| suffix = Path(upload.filename or "").suffix.lower() | |
| if suffix not in allowed_ext: | |
| rejected += 1 | |
| items.append({ | |
| "document_id": None, | |
| "filename": upload.filename, | |
| "status": "rejected", | |
| "message": f"Unsupported type {suffix}. Use PDF, Word, or ZIP.", | |
| }) | |
| continue | |
| upload_tmp_root = settings.data_dir_path / "tmp" / "uploads" | |
| from backend.utils.runtime_paths import ensure_data_drive_runtime_dirs | |
| ensure_data_drive_runtime_dirs() | |
| upload_tmp_root.mkdir(parents=True, exist_ok=True) | |
| tmp_dir = Path(tempfile.mkdtemp(prefix="v2_batch_", dir=str(upload_tmp_root))) | |
| safe_name = Path(upload.filename or f"upload{suffix}").name | |
| tmp_path = tmp_dir / safe_name | |
| try: | |
| content = await upload.read() | |
| if not content: | |
| raise ValueError("Empty upload body") | |
| tmp_path.write_bytes(content) | |
| if suffix == ".zip": | |
| extract_dir = tmp_dir / "extracted" | |
| extracted = extract_reference_documents(tmp_path, extract_dir) | |
| if not extracted: | |
| rejected += 1 | |
| items.append({ | |
| "document_id": None, | |
| "filename": upload.filename, | |
| "status": "rejected", | |
| "message": "ZIP contained no supported PDF or Word files.", | |
| }) | |
| continue | |
| for _inner, extracted_path in extracted: | |
| doc = _ingest_upload_file( | |
| tenant_id, | |
| extracted_path, | |
| extracted_path.name, | |
| ) | |
| accepted += 1 | |
| items.append({ | |
| "document_id": doc.document_id, | |
| "filename": doc.filename, | |
| "status": "complete", | |
| "message": f"Ingested {doc.ingested_chunks} reference chunks.", | |
| }) | |
| else: | |
| if tmp_path.stat().st_size > settings.max_single_upload_bytes: | |
| rejected += 1 | |
| items.append({ | |
| "document_id": None, | |
| "filename": upload.filename, | |
| "status": "rejected", | |
| "message": f"File exceeds max size ({settings.max_single_upload_bytes} bytes).", | |
| }) | |
| continue | |
| doc = _ingest_upload_file( | |
| tenant_id, | |
| tmp_path, | |
| upload.filename or tmp_path.name, | |
| ) | |
| accepted += 1 | |
| items.append({ | |
| "document_id": doc.document_id, | |
| "filename": upload.filename, | |
| "status": "complete", | |
| "message": f"Ingested {doc.ingested_chunks} reference chunks.", | |
| }) | |
| except Exception as exc: # noqa: BLE001 | |
| logger.exception("Batch upload failed for %s", upload.filename) | |
| rejected += 1 | |
| items.append({ | |
| "document_id": None, | |
| "filename": upload.filename, | |
| "status": "failed", | |
| "message": str(exc), | |
| }) | |
| finally: | |
| shutil.rmtree(tmp_dir, ignore_errors=True) | |
| return { | |
| "accepted": accepted, | |
| "rejected": rejected, | |
| "message": f"{accepted} file(s) indexed as past-report reference.", | |
| "items": items, | |
| } | |
| def batch_status( | |
| body: BatchStatusBody, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| docs = list_documents(tenant_id) | |
| items = [] | |
| complete = failed = pending = processing = 0 | |
| for doc_id in body.document_ids: | |
| row = docs.get(doc_id) | |
| if row is None: | |
| failed += 1 | |
| items.append({ | |
| "document_id": doc_id, | |
| "status": "not_found", | |
| "filename": "", | |
| "error": "Unknown document id", | |
| }) | |
| continue | |
| st = row.status | |
| if st == "complete": | |
| complete += 1 | |
| elif st == "failed": | |
| failed += 1 | |
| else: | |
| pending += 1 | |
| items.append({ | |
| "document_id": doc_id, | |
| "status": st, | |
| "filename": row.filename, | |
| "error": row.error, | |
| }) | |
| return { | |
| "pending": pending, | |
| "processing": processing, | |
| "complete": complete, | |
| "failed": failed, | |
| "items": items, | |
| } | |
| def tenant_chunk_summary(tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| count = get_rag_store().count(tenant_id, TIER_REFERENCE) | |
| return {"tenant_id": tenant_id, "indexed_chunk_count": count} | |
| def survey_level_classify( | |
| body: SurveyClassifyBody, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| docs = list_documents(tenant_id) | |
| items = [] | |
| for doc_id in body.document_ids: | |
| row = docs.get(doc_id) | |
| items.append({ | |
| "document_id": doc_id, | |
| "filename": row.filename if row else doc_id, | |
| "predicted_survey_level": 3, | |
| "confidence": 0.85, | |
| "rationale": "Inferred from template schema (v2 reference mapping).", | |
| }) | |
| return {"items": items, "corpus_labelled_files": len(items)} | |
| def survey_level_apply( | |
| body: SurveyApplyBody, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| return {"updated": body.items} | |
| # ── Reports ────────────────────────────────────────────────────────────────── | |
| def create_report( | |
| document_id: str, | |
| survey_level: int = 3, | |
| confirm_tier_mismatch: bool = False, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| if not get_document(tenant_id, document_id): | |
| raise HTTPException(status_code=404, detail="Document not found") | |
| template_discoverer.ensure_canonical_schema(tenant_id) | |
| if get_rag_store().count(tenant_id, TIER_REFERENCE) == 0: | |
| raise HTTPException( | |
| status_code=400, | |
| detail="Upload at least one past report before creating a report job.", | |
| ) | |
| session = create_report_session( | |
| tenant_id, | |
| survey_level=max(1, min(3, survey_level)), | |
| primary_document_id=document_id, | |
| document_ids=[document_id], | |
| ) | |
| return {"report_id": session.report_id, "status": "ready", "message": "Report session created."} | |
| def template_catalog( | |
| survey_level: int = 3, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| schema = template_discoverer.ensure_canonical_schema(tenant_id) | |
| sections = [] | |
| for sec in schema.ordered_sections(): | |
| grp = _section_group(sec.id) | |
| hint = ", ".join(sec.keywords[:6]) if sec.keywords else sec.title | |
| sections.append({ | |
| "code": sec.id, | |
| "group": grp, | |
| "title": sec.title, | |
| "hint": hint, | |
| }) | |
| labels = {k: v for k, v in _GROUP_LABELS.items() if any(s["group"] == k for s in sections)} | |
| return { | |
| "survey_level": max(1, min(3, survey_level)), | |
| "sections": sections, | |
| "group_labels": labels, | |
| } | |
| def photo_policy(report_id: str, tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| schema = template_discoverer.ensure_canonical_schema(tenant_id) | |
| sections = [] | |
| for sec in schema.ordered_sections(): | |
| sections.append({"code": sec.id, "policy": "allowed", "reason": "Photos supported in v2."}) | |
| return { | |
| "report_id": session.report_id, | |
| "sections": sections, | |
| "photo_limits": { | |
| "max_photos_per_section": settings.max_section_photos_per_section, | |
| "max_photos_for_ai": settings.max_section_photos_for_ai, | |
| }, | |
| } | |
| def list_legacy_photos( | |
| report_id: str, | |
| section_code: str, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| photos = _photo_rows(tenant_id, session, section_code) | |
| selected = sum(1 for p in photos if p.get("selected_for_ai")) | |
| return { | |
| "report_id": report_id, | |
| "section_code": section_code, | |
| "photos": photos, | |
| "max_photos_per_section": settings.max_section_photos_per_section, | |
| "max_photos_for_ai": settings.max_section_photos_for_ai, | |
| "selected_for_ai_count": selected, | |
| } | |
| async def upload_legacy_photos( | |
| report_id: str, | |
| section_code: str, | |
| files: list[UploadFile] = File(...), | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| saved = [] | |
| for upload in files: | |
| ct = (upload.content_type or "image/jpeg").lower().split(";")[0].strip() | |
| if ct not in ALLOWED_CONTENT_TYPES: | |
| raise HTTPException(status_code=400, detail=f"Unsupported type: {upload.content_type}") | |
| data = await upload.read() | |
| row = photo_store.add_section_photo( | |
| tenant_id, | |
| session.draft_id, | |
| section_code, | |
| data, | |
| content_type=ct, | |
| original_filename=upload.filename or "", | |
| ) | |
| saved.append({ | |
| "photo_id": row.id, | |
| "original_filename": row.filename, | |
| }) | |
| return { | |
| "report_id": report_id, | |
| "section_code": section_code, | |
| "saved": saved, | |
| "photos": _photo_rows(tenant_id, session, section_code), | |
| } | |
| def legacy_photo_ai_selection( | |
| report_id: str, | |
| section_code: str, | |
| body: PhotoAiSelectionBody, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| try: | |
| photo_store.set_ai_selection(tenant_id, session.draft_id, section_code, body.photo_ids) | |
| except ValueError as exc: | |
| raise HTTPException(status_code=422, detail=str(exc)) from exc | |
| return { | |
| "report_id": report_id, | |
| "section_code": section_code, | |
| "photos": _photo_rows(tenant_id, session, section_code), | |
| } | |
| def delete_legacy_photo( | |
| report_id: str, | |
| section_code: str, | |
| photo_id: str, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| if not photo_store.delete_section_photo(tenant_id, session.draft_id, section_code, photo_id): | |
| raise HTTPException(status_code=404, detail="Photo not found") | |
| return {"deleted": True, "photo_id": photo_id} | |
| def get_legacy_photo( | |
| report_id: str, | |
| section_code: str, | |
| photo_id: str, | |
| tenant_id: str = Depends(_resolve_tenant_token), | |
| ) -> Response: | |
| session = _require_session(tenant_id, report_id) | |
| rows = photo_store.list_section_photos(tenant_id, session.draft_id, section_code) | |
| row = next((r for r in rows if r.id == photo_id), None) | |
| if row is None: | |
| raise HTTPException(status_code=404, detail="Photo not found") | |
| path = photo_store.photo_file_path( | |
| tenant_id, session.draft_id, section_code, photo_id, content_type=row.content_type | |
| ) | |
| if not path or not path.is_file(): | |
| raise HTTPException(status_code=404, detail="Photo file missing") | |
| return Response(content=path.read_bytes(), media_type=row.content_type) | |
| async def generate_report_legacy( | |
| report_id: str, | |
| body: LegacyGenerateBody, | |
| background_tasks: BackgroundTasks, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| from backend.core.section_mapper import estimate_active_sections_from_generate_body | |
| session = _require_session(tenant_id, report_id) | |
| active = estimate_active_sections_from_generate_body( | |
| body, | |
| tenant_id=tenant_id, | |
| draft_id=session.draft_id, | |
| ) | |
| total = len(active) if active else 1 | |
| session.status = "generating" | |
| session.generation_section_total = total | |
| session.generation_started_at = time.time() | |
| session.error_message = None | |
| save_session(session) | |
| background_tasks.add_task(_run_generation, tenant_id, report_id, body) | |
| return { | |
| "report_id": report_id, | |
| "status": "generating", | |
| "message": "Generation started (v2 reference mapping).", | |
| } | |
| def report_status(report_id: str, tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| sections_saved = len(session.sections_payload) | |
| sections_with_text = sum( | |
| 1 for p in session.sections_payload.values() if (p.get("text") or "").strip() | |
| ) | |
| elapsed = None | |
| if session.generation_started_at and session.status == "generating": | |
| elapsed = round(time.time() - session.generation_started_at, 1) | |
| return { | |
| "report_id": report_id, | |
| "status": session.status if session.status != "idle" else "complete", | |
| "created_at": session.created_at, | |
| "updated_at": session.updated_at, | |
| "error_message": session.error_message, | |
| "survey_level": session.survey_level, | |
| "sections_saved": sections_saved, | |
| "sections_with_text": sections_with_text, | |
| "generation_section_total": session.generation_section_total, | |
| "generation_elapsed_seconds": elapsed, | |
| } | |
| def cancel_generation(report_id: str, tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| if session.status == "generating": | |
| session.status = "failed" | |
| session.error_message = "Cancelled by user (timeout)." | |
| save_session(session) | |
| return {"report_id": report_id, "status": session.status} | |
| def get_sections(report_id: str, tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| return {"report_id": report_id, "sections": session.sections_payload} | |
| def patch_section( | |
| report_id: str, | |
| section_code: str, | |
| body: SectionTextPatch, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session = _require_session(tenant_id, report_id) | |
| payload = session.sections_payload.get(section_code, { | |
| "text": "", | |
| "confidence": 0.5, | |
| "provenance": [], | |
| "mode": "generate", | |
| }) | |
| payload["text"] = body.text | |
| session.sections_payload[section_code] = payload | |
| save_session(session) | |
| return payload | |
| def export_report( | |
| report_id: str, | |
| format: str = "docx", | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> Response: | |
| if format.lower() != "docx": | |
| raise HTTPException(status_code=400, detail="Only format=docx is supported.") | |
| session = _require_session(tenant_id, report_id) | |
| schema = template_discoverer.ensure_canonical_schema(tenant_id) | |
| from backend.core.report_assembler import from_sections_payload | |
| result = from_sections_payload( | |
| tenant_id, | |
| schema, | |
| session.sections_payload, | |
| property_type=session.property_type, | |
| tenure=session.tenure, | |
| ) | |
| from backend.core.photo_store import draft_section_photo_paths | |
| photo_paths = draft_section_photo_paths( | |
| tenant_id, session.draft_id, schema.section_ids() | |
| ) | |
| data = to_docx(result, schema, section_photo_paths=photo_paths, include_footer=True) | |
| filename = f"RICS_Report_{report_id[:8]}.docx" | |
| return Response( | |
| content=data, | |
| media_type="application/vnd.openxmlformats-officedocument.wordprocessingml.document", | |
| headers={"Content-Disposition": f'attachment; filename="{filename}"'}, | |
| ) | |
| # ── Optional stubs (document manager / similarity — off main path) ─────────── | |
| def list_documents_route( | |
| limit: int = 500, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| docs = list_documents(tenant_id) | |
| rows = sorted(docs.values(), key=lambda d: d.filename)[:limit] | |
| return { | |
| "documents": [ | |
| { | |
| "document_id": d.document_id, | |
| "filename": d.filename, | |
| "status": d.status, | |
| "created_at": document_library.document_created_at_iso(d), | |
| "file_size": d.file_size or None, | |
| "document_purpose": "reference", | |
| "error": d.error, | |
| } | |
| for d in rows | |
| ], | |
| "reingest_running": document_library.is_reingest_running(tenant_id), | |
| "reingest_progress": document_library.reingest_progress(tenant_id), | |
| } | |
| def delete_document_route( | |
| document_id: str, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| session_reports = [ | |
| s | |
| for s in _reports_generating_for_document(tenant_id, document_id) | |
| ] | |
| if session_reports: | |
| raise HTTPException( | |
| status_code=409, | |
| detail=( | |
| "A report is still generating from this document. " | |
| "Wait for it to finish before deleting the source file." | |
| ), | |
| ) | |
| try: | |
| removed = document_library.remove_reference_document(tenant_id, document_id) | |
| except KeyError as exc: | |
| raise HTTPException(status_code=404, detail="Document not found") from exc | |
| return { | |
| "document_id": document_id, | |
| "deleted": True, | |
| "chunks_removed": removed, | |
| } | |
| def reingest_one_document( | |
| document_id: str, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| if _reports_generating_for_document(tenant_id, document_id): | |
| raise HTTPException( | |
| status_code=409, | |
| detail="A report is still generating from this document. Wait before re-ingesting.", | |
| ) | |
| try: | |
| result = document_library.schedule_reingest_reference_document( | |
| tenant_id, document_id | |
| ) | |
| except KeyError as exc: | |
| raise HTTPException(status_code=404, detail="Document not found") from exc | |
| except FileNotFoundError: | |
| return { | |
| "queued": 0, | |
| "skipped_missing_file": 1, | |
| "detail": "Source file is no longer on disk; cannot re-ingest.", | |
| } | |
| return result | |
| def reingest_all_route(tenant_id: str = Depends(get_current_tenant)) -> dict: | |
| active_doc_ids = { | |
| doc_id | |
| for doc_id in _all_generating_document_ids(tenant_id) | |
| } | |
| return document_library.schedule_reingest_all_documents( | |
| tenant_id, | |
| skip_document_ids=active_doc_ids, | |
| ) | |
| def _reports_dir(tenant_id: str) -> Path: | |
| from backend.core.report_session import _reports_dir as rd # noqa: PLC0415 | |
| return rd(tenant_id) | |
| def _reports_generating_for_document(tenant_id: str, document_id: str) -> list[str]: | |
| out: list[str] = [] | |
| reports_dir = _reports_dir(tenant_id) | |
| if not reports_dir.is_dir(): | |
| return out | |
| import json | |
| for path in reports_dir.glob("*.json"): | |
| try: | |
| data = json.loads(path.read_text(encoding="utf-8")) | |
| except Exception: # noqa: BLE001 | |
| continue | |
| if data.get("status") != "generating": | |
| continue | |
| doc_ids = set(data.get("document_ids") or []) | |
| if data.get("primary_document_id"): | |
| doc_ids.add(data["primary_document_id"]) | |
| if document_id in doc_ids: | |
| out.append(str(data.get("report_id") or path.stem)) | |
| return out | |
| def _all_generating_document_ids(tenant_id: str) -> set[str]: | |
| ids: set[str] = set() | |
| reports_dir = _reports_dir(tenant_id) | |
| if not reports_dir.is_dir(): | |
| return ids | |
| import json | |
| for path in reports_dir.glob("*.json"): | |
| try: | |
| data = json.loads(path.read_text(encoding="utf-8")) | |
| except Exception: # noqa: BLE001 | |
| continue | |
| if data.get("status") != "generating": | |
| continue | |
| ids.update(data.get("document_ids") or []) | |
| if data.get("primary_document_id"): | |
| ids.add(data["primary_document_id"]) | |
| return ids | |
| def content_similar( | |
| body: SimilarContentBody, | |
| tenant_id: str = Depends(get_current_tenant), | |
| ) -> dict: | |
| return scan_similar_content( | |
| tenant_id, | |
| text=body.text, | |
| section_code=body.section_code, | |
| peer_sections=body.peer_sections, | |
| limit=body.limit, | |
| min_relevance_percent=body.min_relevance_percent, | |
| exclude_document_ids=body.exclude_document_ids, | |
| ) | |
| async def route_notes_route(body: RouteNotesBody) -> dict: | |
| """Route raw note lines into canonical L3 section buckets (server-side cascade).""" | |
| lines = [ln.strip() for ln in (body.lines or []) if (ln or "").strip()] | |
| routed, unmatched = route_lines_to_l3_sections(lines) | |
| return { | |
| "routed": routed, | |
| "unmatched": unmatched, | |
| "line_count": len(lines), | |
| "routed_line_count": sum(len(v) for v in routed.values()), | |
| } | |
| async def extract_notes_route(file: UploadFile = File(...)) -> dict: | |
| suffix = Path(file.filename or "").suffix.lower() | |
| if suffix not in {".docx", ".pdf", ".txt"}: | |
| raise HTTPException( | |
| status_code=422, | |
| detail="Unsupported file type. Upload a .docx, .pdf, or .txt notes file.", | |
| ) | |
| data = await file.read(settings.max_notes_extract_bytes + 1) | |
| if len(data) > settings.max_notes_extract_bytes: | |
| raise HTTPException(status_code=413, detail="Notes file too large.") | |
| try: | |
| return extract_notes_from_upload(data, file.filename or "unknown") | |
| except ValueError as exc: | |
| raise HTTPException(status_code=422, detail=str(exc)) from exc | |
| except Exception as exc: # noqa: BLE001 | |
| raise HTTPException(status_code=500, detail=f"Failed to parse file: {exc}") from exc | |