from __future__ import annotations from uuid import UUID from fastapi import APIRouter, Depends, HTTPException from sqlalchemy.orm import Session, joinedload, subqueryload from app.core.dependencies import get_current_user from app.database.session import get_db from app.models.knowledge_item import KnowledgeItem, KnowledgeStatus, KnowledgeType from app.models.user import User from app.models.workspace import Workspace router = APIRouter( prefix="/knowledge", tags=["Knowledge"], ) def _knowledge_response(item: KnowledgeItem) -> dict: dv = item.document_version evidence_list = [] for ev in (item.evidence or []): evidence_list.append({ "quote": ev.quote, "page_number": ev.page_number, "section": ev.section, "confidence": ev.confidence, "source_type": ev.source_type, }) # Include the most recent proposal info for governance context proposal_info = None if item.proposals: # Get the most recent proposal (by created_at or just first) latest_proposal = sorted( item.proposals, key=lambda p: p.created_at or "", reverse=True, )[0] proposal_info = { "id": str(latest_proposal.id), "status": latest_proposal.status.value, "proposal_type": latest_proposal.proposal_type.value, } return { "id": str(item.id), "workspace_id": str(item.workspace_id), "document_version_id": str(item.document_version_id), "filename": dv.filename if dv else None, "type": item.type.value, "title": item.title, "value": item.value, "summary": item.summary, "attributes": item.attributes, "confidence": item.confidence, "status": item.status.value, "proposal": proposal_info, "evidence": evidence_list, "created_at": item.created_at, "updated_at": item.updated_at, } @router.get("") def list_knowledge( workspace_id: UUID, status: str | None = None, type: str | None = None, current_user: User = Depends(get_current_user), db: Session = Depends(get_db), ): workspace = ( db.query(Workspace) .filter( Workspace.id == workspace_id, Workspace.created_by == current_user.id, ) .first() ) if workspace is None: raise HTTPException( status_code=403, detail="You do not have access to this workspace.", ) query = db.query(KnowledgeItem).filter( KnowledgeItem.workspace_id == workspace_id, ) if status: status_upper = status.upper() # Special case: "ARCHIVED" means knowledge items whose proposal # has been archived or the knowledge item itself is marked ARCHIVED. if status_upper == "ARCHIVED": from app.models.proposal import Proposal, ProposalStatus from sqlalchemy import select archived_ki_ids = ( select(Proposal.knowledge_item_id) .filter( Proposal.knowledge_item_id.isnot(None), Proposal.status == ProposalStatus.ARCHIVED, ) ) query = query.filter( (KnowledgeItem.status == KnowledgeStatus.ARCHIVED) | KnowledgeItem.id.in_(archived_ki_ids) ) elif status_upper == "PENDING": # PENDING means genuinely untouched — PENDING knowledge with # a PENDING proposal (not archived, not decided). from app.models.proposal import Proposal, ProposalStatus from sqlalchemy import select archived_or_rejected_ki_ids = ( select(Proposal.knowledge_item_id) .filter( Proposal.knowledge_item_id.isnot(None), Proposal.status.in_([ProposalStatus.ARCHIVED, ProposalStatus.REJECTED]), ) ) query = query.filter( KnowledgeItem.status == KnowledgeStatus.PENDING, ~KnowledgeItem.id.in_(archived_or_rejected_ki_ids), ) else: try: ks = KnowledgeStatus(status_upper) query = query.filter(KnowledgeItem.status == ks) except ValueError: pass if type: try: kt = KnowledgeType(type.upper()) query = query.filter(KnowledgeItem.type == kt) except ValueError: pass items = ( query .options( joinedload(KnowledgeItem.document_version), # 1 JOIN — gets filename subqueryload(KnowledgeItem.evidence), # 1 extra query for all evidence subqueryload(KnowledgeItem.proposals), # 1 extra query for all proposals ) .order_by(KnowledgeItem.created_at.desc()) .limit(200) .all() ) return [_knowledge_response(i) for i in items] @router.get("/search") def search_knowledge( workspace_id: UUID, q: str, current_user: User = Depends(get_current_user), db: Session = Depends(get_db), ): """ Hybrid search: lexical (ILIKE) + vector (pgvector cosine similarity), merged with Reciprocal Rank Fusion (RRF). """ workspace = ( db.query(Workspace) .filter( Workspace.id == workspace_id, Workspace.created_by == current_user.id, ) .first() ) if workspace is None: raise HTTPException( status_code=403, detail="You do not have access to this workspace.", ) # --- Lexical search (substring match on title/value/summary) --- lexical_items = ( db.query(KnowledgeItem) .options( joinedload(KnowledgeItem.document_version), subqueryload(KnowledgeItem.evidence), subqueryload(KnowledgeItem.proposals), ) .filter( KnowledgeItem.workspace_id == workspace_id, ( KnowledgeItem.title.ilike(f"%{q}%") | KnowledgeItem.value.ilike(f"%{q}%") | KnowledgeItem.summary.ilike(f"%{q}%") ), ) .order_by(KnowledgeItem.created_at.desc()) .limit(30) .all() ) # --- Vector search (embed query, find nearest chunks, map to knowledge items) --- vector_items = [] try: from app.services.embedding_service import EmbeddingService from app.models.document_chunk import DocumentChunk from app.models.document_version import DocumentVersion embedding_service = EmbeddingService() query_embedding = embedding_service.embed(q) if query_embedding: # Find nearest chunks in this workspace's documents nearest_chunks = ( db.query(DocumentChunk) .join(DocumentVersion, DocumentChunk.document_version_id == DocumentVersion.id) .filter( DocumentVersion.document.has(workspace_id=workspace_id), DocumentChunk.embedding.isnot(None), ) .order_by(DocumentChunk.embedding.cosine_distance(query_embedding)) .limit(20) .all() ) # Map chunks to knowledge items via document_version_id if nearest_chunks: version_ids = list({c.document_version_id for c in nearest_chunks}) vector_items = ( db.query(KnowledgeItem) .options( joinedload(KnowledgeItem.document_version), subqueryload(KnowledgeItem.evidence), subqueryload(KnowledgeItem.proposals), ) .filter( KnowledgeItem.workspace_id == workspace_id, KnowledgeItem.document_version_id.in_(version_ids), ) .limit(30) .all() ) except Exception: # If vector search fails (model not loaded, etc.), fall back to lexical only pass # --- RRF merge (Reciprocal Rank Fusion, k=60) --- k = 60 scores: dict[str, float] = {} item_map: dict[str, KnowledgeItem] = {} # Score lexical results for rank, item in enumerate(lexical_items): item_id = str(item.id) scores[item_id] = scores.get(item_id, 0) + 1.0 / (k + rank + 1) item_map[item_id] = item # Score vector results for rank, item in enumerate(vector_items): item_id = str(item.id) scores[item_id] = scores.get(item_id, 0) + 1.0 / (k + rank + 1) item_map[item_id] = item # Sort by RRF score descending ranked_ids = sorted(scores.keys(), key=lambda x: scores[x], reverse=True) results = [item_map[item_id] for item_id in ranked_ids[:30]] return [_knowledge_response(i) for i in results] @router.get("/{item_id}") def get_knowledge_item( item_id: UUID, current_user: User = Depends(get_current_user), db: Session = Depends(get_db), ): from app.models.proposal import Proposal from app.models.review import Review from app.models.commit import Commit item = db.query(KnowledgeItem).filter(KnowledgeItem.id == item_id).first() if item is None: raise HTTPException(status_code=404, detail="Knowledge item not found.") workspace = ( db.query(Workspace) .filter( Workspace.id == item.workspace_id, Workspace.created_by == current_user.id, ) .first() ) if workspace is None: raise HTTPException(status_code=403, detail="You do not have access to this item.") # Get proposals related to this knowledge item proposals = ( db.query(Proposal) .filter(Proposal.knowledge_item_id == item_id) .order_by(Proposal.created_at.desc()) .all() ) history = [] for p in proposals: entry = { "type": "proposal", "proposal_id": str(p.id), "proposal_type": p.proposal_type.value, "status": p.status.value, "summary": p.summary, "timestamp": p.created_at.isoformat() if p.created_at else None, } history.append(entry) # If proposal was reviewed, add the review event if p.reviewed_at: history.append({ "type": "decision", "proposal_id": str(p.id), "status": p.status.value, "timestamp": p.reviewed_at.isoformat(), }) # If approved, find the commit if p.status.value == "APPROVED": commit = ( db.query(Commit) .filter(Commit.proposal_id == p.id) .first() ) if commit: history.append({ "type": "commit", "commit_id": str(commit.id), "message": commit.message, "timestamp": commit.committed_at.isoformat() if commit.committed_at else None, }) # Sort history by timestamp history.sort(key=lambda h: h.get("timestamp") or "", reverse=True) resp = _knowledge_response(item) resp["history"] = history return resp