diff --git a/alembic/versions/0131_knowledge_extractions.py b/alembic/versions/0131_knowledge_extractions.py new file mode 100644 index 0000000..6bf4d0b --- /dev/null +++ b/alembic/versions/0131_knowledge_extractions.py @@ -0,0 +1,44 @@ +"""knowledge extractions table + +Revision ID: 0131 +Revises: 0130 +Create Date: 2026-08-20 +""" +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects.postgresql import UUID, JSONB + +revision = "0131" +down_revision = "0130" +branch_labels = None +depends_on = None + +def upgrade() -> None: + op.create_table( + "knowledge_extractions", + sa.Column("id", UUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")), + sa.Column("tenant_id", UUID(as_uuid=True), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False), + sa.Column("source_type", sa.String(50), nullable=False), + sa.Column("source_id", UUID(as_uuid=True), nullable=False), + sa.Column("source_title", sa.String(500), nullable=True), + sa.Column("extracted_entities", JSONB, nullable=False, server_default=sa.text("'[]'::jsonb")), + sa.Column("extracted_relationships", JSONB, nullable=False, server_default=sa.text("'[]'::jsonb")), + sa.Column("confidence", sa.Float, nullable=False, server_default=sa.text("0.0")), + sa.Column("status", sa.String(30), nullable=False, server_default=sa.text("'pending'")), + sa.Column("review_notes", sa.Text, nullable=True), + sa.Column("llm_model", sa.String(100), nullable=True), + sa.Column("llm_cost_usd", sa.Float, nullable=False, server_default=sa.text("0.0")), + sa.Column("created_by", UUID(as_uuid=True), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("reviewed_by", UUID(as_uuid=True), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("reviewed_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.text("now()")), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.text("now()")), + ) + op.create_index("ix_knowledge_ext_tenant_status", "knowledge_extractions", ["tenant_id", "status"]) + op.create_index("ix_knowledge_ext_source", "knowledge_extractions", ["tenant_id", "source_type", "source_id"]) + # RLS + op.execute("ALTER TABLE knowledge_extractions ENABLE ROW LEVEL SECURITY;") + op.execute("CREATE POLICY knowledge_extractions_tenant_isolation ON knowledge_extractions USING (tenant_id::text = current_setting('app.current_tenant_id', true));") + +def downgrade() -> None: + op.drop_table("knowledge_extractions") diff --git a/app/plugins/builtins/knowledge/__init__.py b/app/plugins/builtins/knowledge/__init__.py new file mode 100644 index 0000000..2061ae0 --- /dev/null +++ b/app/plugins/builtins/knowledge/__init__.py @@ -0,0 +1 @@ +"""Knowledge plugin — LLM-based entity/relationship extraction, ask-knowledge, review queue.""" diff --git a/app/plugins/builtins/knowledge/models.py b/app/plugins/builtins/knowledge/models.py new file mode 100644 index 0000000..57c6163 --- /dev/null +++ b/app/plugins/builtins/knowledge/models.py @@ -0,0 +1,32 @@ +"""Knowledge extraction models — tracks LLM extractions and review queue.""" +from __future__ import annotations +import uuid +from datetime import datetime +from sqlalchemy import DateTime, Float, ForeignKey, Index, Integer, String, Text, func +from sqlalchemy.dialects.postgresql import JSONB, UUID as PGUUID +from sqlalchemy.orm import Mapped, mapped_column +from app.core.db import Base, TenantMixin + +class KnowledgeExtraction(Base, TenantMixin): + """Tracks a single knowledge extraction run from a source (wiki, dms, mail, comm).""" + __tablename__ = "knowledge_extractions" + __table_args__ = ( + Index("ix_knowledge_ext_tenant_status", "tenant_id", "status"), + Index("ix_knowledge_ext_source", "tenant_id", "source_type", "source_id"), + ) + id: Mapped[uuid.UUID] = mapped_column(PGUUID(as_uuid=True), primary_key=True, default=uuid.uuid4) + source_type: Mapped[str] = mapped_column(String(50), nullable=False) # wiki_article, dms_file, mail, communication + source_id: Mapped[uuid.UUID] = mapped_column(PGUUID(as_uuid=True), nullable=False) + source_title: Mapped[str | None] = mapped_column(String(500), nullable=True) + extracted_entities: Mapped[list] = mapped_column(JSONB, nullable=False, default=list) + extracted_relationships: Mapped[list] = mapped_column(JSONB, nullable=False, default=list) + confidence: Mapped[float] = mapped_column(Float, nullable=False, default=0.0) + status: Mapped[str] = mapped_column(String(30), nullable=False, default="pending") # pending, approved, rejected, auto_created + review_notes: Mapped[str | None] = mapped_column(Text, nullable=True) + llm_model: Mapped[str | None] = mapped_column(String(100), nullable=True) + llm_cost_usd: Mapped[float] = mapped_column(Float, nullable=False, default=0.0) + created_by: Mapped[uuid.UUID | None] = mapped_column(PGUUID(as_uuid=True), ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + reviewed_by: Mapped[uuid.UUID | None] = mapped_column(PGUUID(as_uuid=True), ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + reviewed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, server_default=func.now()) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()) diff --git a/app/plugins/builtins/knowledge/plugin.py b/app/plugins/builtins/knowledge/plugin.py new file mode 100644 index 0000000..7c610b0 --- /dev/null +++ b/app/plugins/builtins/knowledge/plugin.py @@ -0,0 +1,51 @@ +"""Knowledge plugin — LLM-based entity/relationship extraction, ask-knowledge, review queue.""" +from __future__ import annotations +import logging +from app.plugins.base import BasePlugin +from app.plugins.manifest import PluginManifest, PluginRouteDef + +logger = logging.getLogger(__name__) + +class KnowledgePlugin(BasePlugin): + manifest = PluginManifest( + name="knowledge", + version="1.0.0", + display_name="Knowledge", + description="LLM-based knowledge extraction, ask-knowledge, review queue. Builds on graph_rag + unified_search.", + dependencies=["permissions", "graph_rag", "unified_search"], + routes=[ + PluginRouteDef(path="/api/v1/knowledge", module="app.plugins.builtins.knowledge.routes", router_attr="router"), + ], + permissions=["knowledge:read", "knowledge:write", "knowledge:admin"], + ) + + async def on_activate(self, db, service_container, event_bus) -> None: + """Register event-driven extraction hooks on activation.""" + await super().on_activate(db, service_container, event_bus) + try: + from app.core.hooks import register_action + from app.plugins.builtins.knowledge.services import extract_knowledge + async def on_wiki_create(*args, **kwargs): + article_id = kwargs.get("article_id") or kwargs.get("entity_id") + tenant_id = kwargs.get("tenant_id") + title = kwargs.get("title", "") + content = kwargs.get("content", "") + if article_id and tenant_id and content: + from app.core.db import get_worker_session_factory + factory = get_worker_session_factory() + async with factory() as session: + await extract_knowledge( + db=session, tenant_id=uuid.UUID(str(tenant_id)), + source_type="wiki_article", source_id=uuid.UUID(str(article_id)), + source_title=title, source_text=content, + ) + register_action("wiki.article.created", on_wiki_create, priority=20, owner_tag="knowledge") + logger.info("Registered knowledge extraction hooks") + except Exception: + logger.exception("Failed to register knowledge hooks") + + async def on_deactivate(self, db, service_container, event_bus) -> None: + """Clean up on deactivation.""" + from app.core.hooks import unregister_actions_by_owner + unregister_actions_by_owner("knowledge") + await super().on_deactivate(db, service_container, event_bus) diff --git a/app/plugins/builtins/knowledge/routes.py b/app/plugins/builtins/knowledge/routes.py new file mode 100644 index 0000000..f6a4e8e --- /dev/null +++ b/app/plugins/builtins/knowledge/routes.py @@ -0,0 +1,84 @@ +"""Knowledge plugin routes — extraction, ask, review queue.""" +from __future__ import annotations +import uuid +from fastapi import APIRouter, Depends, HTTPException, Query +from sqlalchemy.ext.asyncio import AsyncSession +from app.core.db import get_db +from app.deps import require_permission +from app.plugins.builtins.knowledge.services import extract_knowledge, ask_knowledge, get_review_queue, review_extraction + +router = APIRouter(prefix="/api/v1/knowledge", tags=["knowledge"]) + +@router.post("/extract") +async def extract( + body: dict, + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_permission("wiki:read")), +): + """Extract knowledge from a source (wiki article, dms file, mail, communication).""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + source_type = body.get("source_type", "") + source_id = body.get("source_id", "") + source_title = body.get("source_title") + source_text = body.get("source_text", "") + if not source_type or not source_id or not source_text: + raise HTTPException(400, detail={"detail": "source_type, source_id, source_text required", "code": "missing_fields"}) + try: + sid = uuid.UUID(source_id) + except ValueError: + raise HTTPException(400, detail={"detail": "Invalid source_id", "code": "invalid_id"}) from None + result = await extract_knowledge( + db=db, tenant_id=tenant_id, source_type=source_type, source_id=sid, + source_title=source_title, source_text=source_text, + user_id=uuid.UUID(current_user["user_id"]) if current_user.get("user_id") else None, + ) + return result + +@router.post("/ask") +async def ask( + body: dict, + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_permission("wiki:read")), +): + """Ask a knowledge question — uses wiki + graph_rag as context.""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + question = body.get("question", "") + if not question: + raise HTTPException(400, detail={"detail": "question required", "code": "missing_question"}) + result = await ask_knowledge(db=db, tenant_id=tenant_id, question=question) + return result + +@router.get("/review") +async def review_queue( + page: int = Query(1, ge=1), + page_size: int = Query(20, ge=1, le=100), + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_permission("wiki:read")), +): + """Get pending knowledge extractions for review.""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + return await get_review_queue(db=db, tenant_id=tenant_id, page=page, page_size=page_size) + +@router.post("/review/{extraction_id}") +async def review( + extraction_id: str, + body: dict, + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_permission("wiki:write")), +): + """Approve or reject a knowledge extraction.""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + approved = body.get("approved", False) + notes = body.get("notes") + try: + eid = uuid.UUID(extraction_id) + except ValueError: + raise HTTPException(400, detail={"detail": "Invalid extraction_id", "code": "invalid_id"}) from None + result = await review_extraction( + db=db, tenant_id=tenant_id, extraction_id=eid, approved=approved, + user_id=uuid.UUID(current_user["user_id"]) if current_user.get("user_id") else None, + notes=notes, + ) + if "error" in result: + raise HTTPException(404, detail={"detail": result["error"], "code": "not_found"}) + return result diff --git a/app/plugins/builtins/knowledge/services.py b/app/plugins/builtins/knowledge/services.py new file mode 100644 index 0000000..c04c8c1 --- /dev/null +++ b/app/plugins/builtins/knowledge/services.py @@ -0,0 +1,211 @@ +"""Knowledge extraction services — LLM-based entity/relationship extraction.""" +from __future__ import annotations +import logging +import uuid +from typing import Any +from sqlalchemy import select, update +from sqlalchemy.ext.asyncio import AsyncSession +from app.ai.llm_client import llm_complete +from app.plugins.builtins.knowledge.models import KnowledgeExtraction + +logger = logging.getLogger(__name__) + +EXTRACTION_PROMPT = """You are a knowledge extraction assistant for a CRM system. +Analyze the following text and extract entities and relationships. + +Return JSON with this structure: +{ + "entities": [ + {"type": "person|company|project|topic", "name": "...", "description": "..."} + ], + "relationships": [ + {"source": "entity_name", "target": "entity_name", "type": "works_for|related_to|part_of|mentions"} + ], + "confidence": 0.0-1.0 +} + +Text to analyze: +""" + +async def extract_knowledge( + db: AsyncSession, + tenant_id: uuid.UUID, + source_type: str, + source_id: uuid.UUID, + source_title: str | None, + source_text: str, + user_id: uuid.UUID | None = None, + model: str = "openai/gpt-4o-mini", +) -> dict[str, Any]: + """Extract entities and relationships from text using LLM.""" + messages = [ + {"role": "system", "content": "You are a knowledge extraction assistant. Return only valid JSON."}, + {"role": "user", "content": EXTRACTION_PROMPT + source_text[:4000]}, + ] + response = await llm_complete( + model=model, + messages=messages, + temperature=0.2, + max_tokens=2000, + tenant_id=tenant_id, + db=db, + ) + import json + try: + result = json.loads(response.get("content", "{}")) + except (json.JSONDecodeError, TypeError): + result = {"entities": [], "relationships": [], "confidence": 0.0} + extraction = KnowledgeExtraction( + tenant_id=tenant_id, + source_type=source_type, + source_id=source_id, + source_title=source_title, + extracted_entities=result.get("entities", []), + extracted_relationships=result.get("relationships", []), + confidence=float(result.get("confidence", 0.0)), + status="auto_created" if result.get("confidence", 0.0) >= 0.8 else "pending", + llm_model=model, + llm_cost_usd=response.get("cost_usd", 0.0), + created_by=user_id, + ) + db.add(extraction) + await db.flush() + # Auto-create relationships in GraphRAG if confidence >= 0.8 + if extraction.confidence >= 0.8 and extraction.extracted_relationships: + try: + from app.plugins.builtins.graph_rag.services import create_relationship + for rel in extraction.extracted_relationships: + # Only create if both source and target have IDs (resolved entities) + if rel.get("source_id") and rel.get("target_id"): + await create_relationship( + db=db, + tenant_id=tenant_id, + source_type=rel.get("source_type", "topic"), + source_id=uuid.UUID(rel["source_id"]), + target_type=rel.get("target_type", "topic"), + target_id=uuid.UUID(rel["target_id"]), + relationship_type=rel.get("type", "related_to"), + metadata={"extraction_id": str(extraction.id), "confidence": extraction.confidence}, + owner_id=user_id, + ) + except Exception as e: + logger.warning("Failed to auto-create relationships: %s", e) + await db.commit() + return { + "id": str(extraction.id), + "entities": extraction.extracted_entities, + "relationships": extraction.extracted_relationships, + "confidence": extraction.confidence, + "status": extraction.status, + } + +async def ask_knowledge( + db: AsyncSession, + tenant_id: uuid.UUID, + question: str, + model: str = "openai/gpt-4o-mini", +) -> dict[str, Any]: + """Answer a knowledge question using wiki articles + graph_rag as context.""" + # Search wiki articles for context + from app.plugins.builtins.unified_search.provider_registry import get_search_registry + registry = get_search_registry() + wiki_provider = registry.get("wiki_article") + context_parts = [] + if wiki_provider: + results = await wiki_provider._search_fts_filtered( + db=db, tsquery=question, tenant_id=tenant_id, limit=5, visible_ids=None + ) + for r in results: + context_parts.append(f"Title: {r.get('title', '')}\nContent: {r.get('content', '')[:500]}") + # Search graph_rag for relationships + graph_provider = registry.get("graph_relationship") + if graph_provider: + results = await graph_provider._search_fts_filtered( + db=db, tsquery=question, tenant_id=tenant_id, limit=5, visible_ids=None + ) + for r in results: + context_parts.append(f"Relationship: {r.get('source_type')} -> {r.get('relationship_type')} -> {r.get('target_type')}") + context = "\n\n".join(context_parts) if context_parts else "No knowledge base content found." + messages = [ + {"role": "system", "content": f"You are a knowledge assistant. Answer based on this context:\n\n{context}"}, + {"role": "user", "content": question}, + ] + response = await llm_complete( + model=model, messages=messages, temperature=0.3, max_tokens=1000, + tenant_id=tenant_id, db=db, + ) + return { + "answer": response.get("content", ""), + "evidence": context_parts[:3], + "cost_usd": response.get("cost_usd", 0.0), + } + +async def get_review_queue( + db: AsyncSession, + tenant_id: uuid.UUID, + page: int = 1, + page_size: int = 20, +) -> dict[str, Any]: + """Get pending knowledge extractions for review.""" + q = select(KnowledgeExtraction).where( + KnowledgeExtraction.tenant_id == tenant_id, + KnowledgeExtraction.status == "pending", + ).order_by(KnowledgeExtraction.created_at.desc()) + from sqlalchemy import func + count_q = select(func.count()).select_from(q.subquery()) + total = (await db.execute(count_q)).scalar() or 0 + offset = (page - 1) * page_size + result = await db.execute(q.offset(offset).limit(page_size)) + items = [ + { + "id": str(e.id), "source_type": e.source_type, "source_id": str(e.source_id), + "source_title": e.source_title, "entities": e.extracted_entities, + "relationships": e.extracted_relationships, "confidence": e.confidence, + "status": e.status, "created_at": e.created_at.isoformat() if e.created_at else None, + } + for e in result.scalars().all() + ] + return {"items": items, "total": total, "page": page, "page_size": page_size} + +async def review_extraction( + db: AsyncSession, + tenant_id: uuid.UUID, + extraction_id: uuid.UUID, + approved: bool, + user_id: uuid.UUID, + notes: str | None = None, +) -> dict[str, Any]: + """Approve or reject a knowledge extraction.""" + from datetime import datetime, timezone + result = await db.execute( + select(KnowledgeExtraction).where( + KnowledgeExtraction.tenant_id == tenant_id, + KnowledgeExtraction.id == extraction_id, + ) + ) + extraction = result.scalar_one_or_none() + if not extraction: + return {"error": "Extraction not found"} + extraction.status = "approved" if approved else "rejected" + extraction.reviewed_by = user_id + extraction.reviewed_at = datetime.now(timezone.utc) + extraction.review_notes = notes + if approved and extraction.confidence < 0.8 and extraction.extracted_relationships: + try: + from app.plugins.builtins.graph_rag.services import create_relationship + for rel in extraction.extracted_relationships: + if rel.get("source_id") and rel.get("target_id"): + await create_relationship( + db=db, tenant_id=tenant_id, + source_type=rel.get("source_type", "topic"), + source_id=uuid.UUID(rel["source_id"]), + target_type=rel.get("target_type", "topic"), + target_id=uuid.UUID(rel["target_id"]), + relationship_type=rel.get("type", "related_to"), + metadata={"extraction_id": str(extraction.id), "reviewed": True}, + owner_id=user_id, + ) + except Exception as e: + logger.warning("Failed to create relationships after approval: %s", e) + await db.commit() + return {"id": str(extraction.id), "status": extraction.status}