diff --git a/app/ai/agent_workstream.py b/app/ai/agent_workstream.py deleted file mode 100644 index af9a943..0000000 --- a/app/ai/agent_workstream.py +++ /dev/null @@ -1,201 +0,0 @@ -"""Agent → Workstream integration. - -Posts agent messages to the central communication system (kommunikation plugin). -Supports text, action_card, entity_card, task_card, approval, and miniapp block types. -All agent messages are marked as AI-generated via transparency metadata. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - -from app.ai.transparency import mark_as_ai_generated - -logger = logging.getLogger(__name__) - - -async def post_agent_message( - db: AsyncSession, - tenant_id: uuid.UUID, - agent_id: uuid.UUID, - agent_run_id: uuid.UUID, - content: str, - block_type: str = "text", - block_data: dict[str, Any] | None = None, - conversation_id: uuid.UUID | None = None, -) -> uuid.UUID: - """Post a message from an agent to the communication system. - - Creates a CommMessage in the agent's conversation channel. - If no conversation_id is provided, uses the agent's default channel. - - Args: - db: Database session. - tenant_id: Tenant ID. - agent_id: Agent definition ID. - agent_run_id: Agent run ID for traceability. - content: Message text content. - block_type: Block type (text, action_card, entity_card, task_card, approval, miniapp). - block_data: Additional block data (e.g. action buttons, entity reference). - conversation_id: Optional conversation to post to. If None, uses agent channel. - - Returns: - Message ID. - """ - from app.plugins.builtins.kommunikation.models import CommMessage, CommMessageBlock - from app.plugins.builtins.kommunikation.services import send_message - - # Mark as AI-generated - ai_metadata = mark_as_ai_generated(content, { - "agent_id": str(agent_id), - "agent_run_id": str(agent_run_id), - }) - - # Build block if not plain text - blocks: list[dict[str, Any]] = [] - if block_type != "text" and block_data: - blocks.append({ - "type": block_type, - "data": block_data, - }) - - # Post via kommunikation service - message_id = await send_message( - db=db, - tenant_id=tenant_id, - sender_id=agent_id, - sender_type="agent", - conversation_id=conversation_id or await _get_or_create_agent_channel(db, tenant_id, agent_id), - content=content, - blocks=blocks, - metadata=ai_metadata, - ) - - logger.info( - "Agent %s posted message %s (block_type=%s, run=%s)", - agent_id, message_id, block_type, agent_run_id, - ) - return message_id - - -async def _get_or_create_agent_channel( - db: AsyncSession, tenant_id: uuid.UUID, agent_id: uuid.UUID -) -> uuid.UUID: - """Get or create a dedicated conversation channel for an agent.""" - from app.plugins.builtins.kommunikation.models import CommConversation - - # Try to find existing agent channel - result = await db.execute( - select(CommConversation).where( - CommConversation.tenant_id == tenant_id, - CommConversation.entity_type == "agent", - CommConversation.entity_id == agent_id, - CommConversation.deleted_at.is_(None), - ) - ) - conv = result.scalar_one_or_none() - if conv: - return conv.id - - # Create new channel - conv = CommConversation( - tenant_id=tenant_id, entity_type="agent", - entity_id=agent_id, - title=f"Agent Channel", - conversation_type="channel", - is_system=False, - ) - db.add(conv) - await db.flush() - return conv.id - - -async def post_agent_step( - db: AsyncSession, - tenant_id: uuid.UUID, - agent_id: uuid.UUID, - agent_run_id: uuid.UUID, - step_number: int, - thought: str, - action: str | None = None, - observation: str | None = None, -) -> uuid.UUID | None: - """Post a ReAct step as an action_card to the workstream. - - Only posts if the agent's trace_mode is 'extended'. - """ - block_data = { - "step_number": step_number, - "thought": thought[:500], # Truncate for display - "action": action, - "observation": (observation or "")[:500], - } - return await post_agent_message( - db=db, - tenant_id=tenant_id, - agent_id=agent_id, - agent_run_id=agent_run_id, - content=f"Step {step_number}: {action or 'Thinking...'}", - block_type="action_card", - block_data=block_data, - ) - - -async def post_agent_result( - db: AsyncSession, - tenant_id: uuid.UUID, - agent_id: uuid.UUID, - agent_run_id: uuid.UUID, - final_content: str, - total_cost_usd: float, - steps_taken: int, - status: str, -) -> uuid.UUID: - """Post the final result of an agent run to the workstream.""" - block_data = { - "status": status, - "steps_taken": steps_taken, - "total_cost_usd": round(total_cost_usd, 6), - "run_id": str(agent_run_id), - } - return await post_agent_message( - db=db, - tenant_id=tenant_id, - agent_id=agent_id, - agent_run_id=agent_run_id, - content=final_content, - block_type="action_card", - block_data=block_data, - ) - - -async def post_approval_request( - db: AsyncSession, - tenant_id: uuid.UUID, - agent_id: uuid.UUID, - agent_run_id: uuid.UUID, - approval_id: uuid.UUID, - action: str, - description: str, -) -> uuid.UUID: - """Post an approval request card to the workstream.""" - block_data = { - "approval_id": str(approval_id), - "action": action, - "description": description, - "status": "pending", - } - return await post_agent_message( - db=db, - tenant_id=tenant_id, - agent_id=agent_id, - agent_run_id=agent_run_id, - content=f"Approval required: {action}", - block_type="approval", - block_data=block_data, - ) diff --git a/app/ai/dashboard.py b/app/ai/dashboard.py deleted file mode 100644 index 3fbc4fe..0000000 --- a/app/ai/dashboard.py +++ /dev/null @@ -1,300 +0,0 @@ -"""Platform dashboard & analytics (I-DASH, I-COST, I-USE). - -Provides aggregated metrics for: -- Agent status, workflow stats, search metrics, knowledge coverage -- LLM cost tracking per agent/workflow/user, budget alerts -- Feature usage, search queries, agent runs, workflow executions -- Proactive suggestions accepted/rejected -""" - -from __future__ import annotations - -import logging -import uuid -from datetime import UTC, datetime, timedelta -from typing import Any - -from sqlalchemy import func, select -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── I-DASH: Platform Dashboard ───────────────────────────────────────────── - - -async def get_platform_dashboard( - db: AsyncSession, - tenant_id: uuid.UUID, -) -> dict[str, Any]: - """Get aggregated platform metrics for the dashboard (I-DASH). - - Returns: Agent status, workflow stats, search metrics, - knowledge coverage, workstream metrics, system health. - """ - dashboard: dict[str, Any] = { - "agents": {}, - "workflows": {}, - "search": {}, - "knowledge": {}, - "workstream": {}, - "system_health": {}, - "generated_at": datetime.now(UTC).isoformat(), - } - - # Agent metrics - try: - from app.models.workflow import AgentDefinition, AgentRun - active_agents = await db.scalar( - select(func.count(AgentDefinition.id)).where( - AgentDefinition.tenant_id == tenant_id, - AgentDefinition.is_active == True, # noqa: E712 - ) - ) - total_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - ) - ) - recent_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= datetime.now(UTC) - timedelta(days=7), - ) - ) - dashboard["agents"] = { - "active_agents": active_agents or 0, - "total_runs": total_runs or 0, - "recent_runs_7d": recent_runs or 0, - } - except Exception as e: - logger.warning("Dashboard agent metrics failed: %s", e) - dashboard["agents"] = {"error": str(e)} - - # Workflow metrics - try: - from app.models.workflow import WorkflowDefinition, WorkflowInstance - active_workflows = await db.scalar( - select(func.count(WorkflowDefinition.id)).where( - WorkflowDefinition.tenant_id == tenant_id, - WorkflowDefinition.is_active == True, # noqa: E712 - ) - ) - running_instances = await db.scalar( - select(func.count(WorkflowInstance.id)).where( - WorkflowInstance.tenant_id == tenant_id, - WorkflowInstance.status.in_(["pending", "running", "waiting"]), - ) - ) - completed_instances = await db.scalar( - select(func.count(WorkflowInstance.id)).where( - WorkflowInstance.tenant_id == tenant_id, - WorkflowInstance.status == "completed", - ) - ) - dashboard["workflows"] = { - "active_workflows": active_workflows or 0, - "running_instances": running_instances or 0, - "completed_instances": completed_instances or 0, - } - except Exception as e: - logger.warning("Dashboard workflow metrics failed: %s", e) - dashboard["workflows"] = {"error": str(e)} - - # Knowledge metrics - try: - from app.plugins.builtins.wiki.models import WikiArticle - wiki_articles = await db.scalar( - select(func.count(WikiArticle.id)).where( - WikiArticle.tenant_id == tenant_id, - WikiArticle.deleted_at.is_(None), - ) - ) - dashboard["knowledge"] = { - "wiki_articles": wiki_articles or 0, - } - except Exception as e: - logger.warning("Dashboard knowledge metrics failed: %s", e) - dashboard["knowledge"] = {"error": str(e)} - - # System health (from Redis check) - try: - from app.core.redis import get_redis - redis = await get_redis() - if redis: - await redis.ping() - dashboard["system_health"] = {"redis": "up", "status": "healthy"} - else: - dashboard["system_health"] = {"redis": "down", "status": "degraded"} - except Exception as e: - dashboard["system_health"] = {"redis": "error", "status": "degraded", "error": str(e)} - - return dashboard - - -# ─── I-COST: Cost Tracking Dashboard ───────────────────────────────────────── - - -async def get_cost_dashboard( - db: AsyncSession, - tenant_id: uuid.UUID, - days: int = 30, -) -> dict[str, Any]: - """Get LLM cost tracking metrics (I-COST). - - Returns: Live costs, budget utilization, alert history, - cost per tenant/agent/workflow, hard-stop events. - """ - since = datetime.now(UTC) - timedelta(days=days) - - cost_data: dict[str, Any] = { - "period_days": days, - "total_cost_usd": 0.0, - "by_agent": {}, - "by_workflow": {}, - "by_user": {}, - "daily_trend": [], - "budget": {}, - "alerts": [], - "generated_at": datetime.now(UTC).isoformat(), - } - - # Aggregate costs from AgentRun - try: - from app.models.workflow import AgentRun - - # Total cost - total_cost = await db.scalar( - select(func.sum(AgentRun.total_cost_usd)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ) - ) - cost_data["total_cost_usd"] = float(total_cost or 0.0) - - # Cost by agent - agent_costs = await db.execute( - select( - AgentRun.agent_id, - func.sum(AgentRun.total_cost_usd).label("cost"), - func.count(AgentRun.id).label("runs"), - ) - .where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ) - .group_by(AgentRun.agent_id) - ) - for row in agent_costs: - cost_data["by_agent"][str(row.agent_id)] = { - "cost_usd": float(row.cost or 0.0), - "runs": row.runs, - } - - except Exception as e: - logger.warning("Cost dashboard failed: %s", e) - cost_data["error"] = str(e) - - # Budget info from config - try: - from app.config import get_settings - settings = get_settings() - monthly_budget = getattr(settings, "llm_monthly_budget_usd", None) - if monthly_budget: - cost_data["budget"] = { - "monthly_limit_usd": monthly_budget, - "utilization_pct": (cost_data["total_cost_usd"] / monthly_budget) * 100, - } - except Exception: - pass - - return cost_data - - -# ─── I-USE: Usage & Collaboration Analytics ───────────────────────────────── - - -async def get_usage_analytics( - db: AsyncSession, - tenant_id: uuid.UUID, - days: int = 30, -) -> dict[str, Any]: - """Get feature usage and collaboration analytics (I-USE). - - Returns: Feature usage, search queries, agent runs, - workflow executions, proactive suggestions accepted/rejected. - """ - since = datetime.now(UTC) - timedelta(days=days) - - analytics: dict[str, Any] = { - "period_days": days, - "agent_runs": {}, - "workflow_executions": {}, - "search_queries": {}, - "proactive_suggestions": {}, - "generated_at": datetime.now(UTC).isoformat(), - } - - # Agent run stats - try: - from app.models.workflow import AgentRun - total_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ) - ) - completed_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - AgentRun.status == "completed", - ) - ) - failed_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - AgentRun.status.in_(["stopped_error", "stopped_timeout"]), - ) - ) - analytics["agent_runs"] = { - "total": total_runs or 0, - "completed": completed_runs or 0, - "failed": failed_runs or 0, - "success_rate": (completed_runs / total_runs * 100) if total_runs else 0.0, - } - except Exception as e: - analytics["agent_runs"] = {"error": str(e)} - - # Workflow execution stats - try: - from app.models.workflow import WorkflowInstance - total_instances = await db.scalar( - select(func.count(WorkflowInstance.id)).where( - WorkflowInstance.tenant_id == tenant_id, - WorkflowInstance.created_at >= since, - ) - ) - completed_instances = await db.scalar( - select(func.count(WorkflowInstance.id)).where( - WorkflowInstance.tenant_id == tenant_id, - WorkflowInstance.created_at >= since, - WorkflowInstance.status == "completed", - ) - ) - analytics["workflow_executions"] = { - "total": total_instances or 0, - "completed": completed_instances or 0, - } - except Exception as e: - analytics["workflow_executions"] = {"error": str(e)} - - return analytics - - -__all__ = [ - "get_platform_dashboard", - "get_cost_dashboard", - "get_usage_analytics", -] diff --git a/app/ai/dsgvo_export.py b/app/ai/dsgvo_export.py deleted file mode 100644 index 3da57bb..0000000 --- a/app/ai/dsgvo_export.py +++ /dev/null @@ -1,346 +0,0 @@ -"""DSGVO-Betroffenenrechte & Compliance Export (I-DSGVO, I-DSAR, I-COMP-EXPORT). - -Provides: -- Full platform data subject access export (JSON/ZIP) -- Data subject rights workflow (access/correction/erasure/restriction) -- AI/Compliance evidence export (audit, oversight, approval records) - -Sensitive/Exposure rules are always respected. No blind auto-delete -over legal retention obligations. -""" - -from __future__ import annotations - -import logging -import uuid -from datetime import UTC, datetime, timedelta -from typing import Any, Literal - -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── I-DSGVO: Platform Data Subject Access Export ─────────────────────────── - - -async def export_user_data( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, -) -> dict[str, Any]: - """Export all personal data for a user across core and active plugins (I-DSGVO). - - Collects data from: CRM (contacts, companies), Mail, Calendar, DMS, - Communication/Workstreams, Agents, Workflows, Knowledge, Audit. - - Returns structured JSON ready for ZIP packaging. - Sensitive fields are masked per data_policy rules. - """ - export: dict[str, Any] = { - "export_metadata": { - "exported_at": datetime.now(UTC).isoformat(), - "tenant_id": str(tenant_id), - "user_id": str(user_id), - "export_type": "dsgvo_data_subject_access", - "version": "1.0", - }, - "core": {}, - "mail": {}, - "calendar": {}, - "dms": {}, - "communication": {}, - "agents": {}, - "workflows": {}, - "knowledge": {}, - "audit": {}, - } - - # Core: User profile - try: - from app.models.user import User - user = await db.get(User, user_id) - if user: - export["core"]["user"] = { - "id": str(user.id), - "email": user.email, - "full_name": getattr(user, "full_name", None), - "is_active": user.is_active, - "is_system_admin": getattr(user, "is_system_admin", False), - "created_at": user.created_at.isoformat() if user.created_at else None, - } - except Exception as e: - export["core"]["error"] = str(e) - - # Core: Contacts owned by user - try: - from app.models.contact import Contact - result = await db.execute( - select(Contact).where( - Contact.tenant_id == tenant_id, - Contact.owner_id == user_id, - Contact.deleted_at.is_(None), - ) - ) - contacts = result.scalars().all() - export["core"]["contacts"] = [ - { - "id": str(c.id), - "first_name": c.first_name, - "last_name": c.last_name, - "email": c.email, - "phone": c.phone, - "created_at": c.created_at.isoformat() if c.created_at else None, - } - for c in contacts - ] - except Exception as e: - export["core"]["contacts_error"] = str(e) - - # Agents: Agent runs by user - try: - from app.models.workflow import AgentRun - result = await db.execute( - select(AgentRun).where( - AgentRun.tenant_id == tenant_id, - AgentRun.user_id == user_id, - ).limit(100) - ) - runs = result.scalars().all() - export["agents"]["agent_runs"] = [ - { - "id": str(r.id), - "status": r.status, - "total_cost_usd": float(r.total_cost_usd or 0), - "created_at": r.created_at.isoformat() if r.created_at else None, - } - for r in runs - ] - except Exception as e: - export["agents"]["error"] = str(e) - - # Audit: User's audit entries - try: - from app.models.audit import AuditLog - result = await db.execute( - select(AuditLog).where( - AuditLog.tenant_id == tenant_id, - AuditLog.user_id == user_id, - ).limit(200) - ) - entries = result.scalars().all() - export["audit"]["entries"] = [ - { - "id": str(e.id), - "action": e.action, - "entity_type": e.entity_type, - "created_at": e.created_at.isoformat() if e.created_at else None, - } - for e in entries - ] - except Exception as e: - export["audit"]["error"] = str(e) - - return export - - -# ─── I-DSAR: Data Subject Rights Workflow ──────────────────────────────────── - - -DSARType = Literal["access", "correction", "erasure", "restriction"] - - -async def create_dsar_request( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - subject_user_id: uuid.UUID, - request_type: DSARType, - description: str = "", -) -> dict[str, Any]: - """Create a data subject rights request (I-DSAR). - - Creates a trackable Task for the DSGVO request. Finds affected sources, - calls domain handlers, tracks derived data via lifecycle, documents - exceptions/retention. No generic blind hard-delete. - """ - from app.plugins.builtins.tasks.services import create_task - - task_data: dict[str, Any] = { - "title": f"DSAR: {request_type} for user {subject_user_id}", - "description": description or f"Data subject {request_type} request", - "task_type": "dsar", - "assignee_type": "user", - "assignee_id": str(user_id), - "entity_type": "user", - "entity_id": str(subject_user_id), - "status": "open", - "priority": "high", - } - - task = await create_task(db, tenant_id, user_id, task_data) - - # Find affected data sources - affected_sources = await _find_affected_sources(db, tenant_id, subject_user_id) - - return { - "task": task, - "request_type": request_type, - "subject_user_id": str(subject_user_id), - "affected_sources": affected_sources, - } - - -async def _find_affected_sources( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, -) -> list[dict[str, str]]: - """Find all data sources containing personal data for a user.""" - sources: list[dict[str, str]] = [] - - # Check each source - source_checks = [ - ("core.contacts", "Contact", "owner_id"), - ("mail.accounts", "MailAccount", "user_id"), - ("dms.files", "DmsFile", "owner_id"), - ("communication.messages", "CommMessage", "sender_id"), - ("agents.runs", "AgentRun", "user_id"), - ] - - for source_name, model_name, id_field in source_checks: - try: - # Dynamic import would be needed here; for now just list the source - sources.append({ - "source": source_name, - "model": model_name, - "id_field": id_field, - "status": "identified", - }) - except Exception: - pass - - return sources - - -# ─── I-COMP-EXPORT: AI/Compliance Evidence Export ──────────────────────────── - - -async def export_compliance_evidence( - db: AsyncSession, - tenant_id: uuid.UUID, - days: int = 90, -) -> dict[str, Any]: - """Export AI/Compliance evidence package (I-COMP-EXPORT). - - Returns: AI use case metadata, provider/model references, - agent/workflow versions, audit/oversight/approval evidence, - and technical policies as exportable evidence package. - """ - since = datetime.now(UTC) - timedelta(days=days) - - evidence: dict[str, Any] = { - "export_metadata": { - "exported_at": datetime.now(UTC).isoformat(), - "tenant_id": str(tenant_id), - "export_type": "compliance_evidence", - "period_days": days, - "version": "1.0", - }, - "ai_use_cases": [], - "agent_definitions": [], - "workflow_definitions": [], - "audit_entries": [], - "approval_records": [], - "oversight_records": [], - "technical_policies": {}, - } - - # Agent definitions with AI metadata - try: - from app.models.workflow import AgentDefinition - result = await db.execute( - select(AgentDefinition).where( - AgentDefinition.tenant_id == tenant_id, - AgentDefinition.is_active == True, # noqa: E712 - ) - ) - agents = result.scalars().all() - evidence["agent_definitions"] = [ - { - "id": str(a.id), - "name": a.name, - "llm_model": getattr(a, "llm_model", None), - "provider": getattr(a, "provider", None), - "is_active": a.is_active, - "created_at": a.created_at.isoformat() if a.created_at else None, - } - for a in agents - ] - except Exception as e: - evidence["agent_definitions_error"] = str(e) - - # Approval records - try: - from app.core.approval import ApprovalRequest - result = await db.execute( - select(ApprovalRequest).where( - ApprovalRequest.tenant_id == tenant_id, - ApprovalRequest.created_at >= since, - ).limit(100) - ) - approvals = result.scalars().all() - evidence["approval_records"] = [ - { - "id": str(a.id), - "action": a.action, - "status": a.status, - "created_at": a.created_at.isoformat() if a.created_at else None, - } - for a in approvals - ] - except Exception as e: - evidence["approval_records_error"] = str(e) - - # Technical policies - evidence["technical_policies"] = { - "data_policy": { - "sensitive_fields": list(_get_sensitive_fields()), - "provider_compliance": "enforced", - }, - "permission_model": { - "type": "ABAC", - "tenant_isolation": "RLS", - }, - "auth": { - "type": "session_based", - "cookies": "HttpOnly", - }, - "retention": { - "soft_delete": True, - "hard_delete_requires_gdpr_flag": True, - }, - } - - return evidence - - -def _get_sensitive_fields() -> dict[str, set[str]]: - """Get the sensitive fields mapping from data_policy. - - Returns a dict mapping entity types to their sensitive field sets. - """ - try: - from app.ai.data_policy import SENSITIVE_FIELDS - return SENSITIVE_FIELDS - except Exception: - return {"contact": {"email", "phone", "address", "date_of_birth"}} - - -__all__ = [ - "export_user_data", - "create_dsar_request", - "export_compliance_evidence", - "DSARType", -] diff --git a/app/ai/integration_tools.py b/app/ai/integration_tools.py deleted file mode 100644 index bb0a94c..0000000 --- a/app/ai/integration_tools.py +++ /dev/null @@ -1,217 +0,0 @@ -"""Integration tools — Agent → Workflow and Agent → Knowledge (I-AW, I-AK). - -Provides AI agent tools for: -- Starting and checking workflow status (I-AW) -- Querying knowledge base with evidence (I-AK) - -These tools are registered in the AI tool registry and can be used by -agents via the ReAct loop. Each tool respects tenant_id and permissions. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── I-AW: Agent → Workflow Tools ──────────────────────────────────────────── - - -async def start_workflow_tool( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - workflow_id: str, - context: dict[str, Any] | None = None, -) -> dict[str, Any]: - """Agent tool: Start a workflow instance. - - Args: - workflow_id: The workflow definition ID. - context: Optional initial context variables. - - Returns: - Dict with instance_id, status, and workflow info. - """ - from app.services.workflow_service import create_instance - - try: - result = await create_instance( - db, - tenant_id, - user_id, - workflow_id=workflow_id, - context=context or {}, - ) - if result is None: - return {"error": "Workflow not found", "status": "not_found"} - return { - "instance_id": result.get("id"), - "status": result.get("status"), - "workflow_id": workflow_id, - "message": f"Workflow started successfully", - } - except Exception as e: - logger.warning("start_workflow_tool failed: %s", e) - return {"error": str(e), "status": "failed"} - - -async def check_workflow_status_tool( - db: AsyncSession, - tenant_id: uuid.UUID, - instance_id: str, -) -> dict[str, Any]: - """Agent tool: Check the status of a workflow instance. - - Args: - instance_id: The workflow instance ID. - - Returns: - Dict with status, current_step, and step history. - """ - from app.services.workflow_service import get_instance - from app.models.workflow import WorkflowStepHistory - from sqlalchemy import select - - try: - instance = await get_instance(db, tenant_id, instance_id) - if instance is None: - return {"error": "Instance not found", "status": "not_found"} - - # Get step history - history_result = await db.execute( - select(WorkflowStepHistory) - .where( - WorkflowStepHistory.tenant_id == tenant_id, - WorkflowStepHistory.instance_id == uuid.UUID(instance_id), - ) - .order_by(WorkflowStepHistory.created_at.desc()) - .limit(5) - ) - recent_steps = [ - { - "step_index": h.step_index, - "step_type": h.step_type, - "action": h.action, - } - for h in history_result.scalars().all() - ] - - return { - "instance_id": instance_id, - "status": instance.get("status"), - "current_step_index": instance.get("current_step_index"), - "recent_steps": recent_steps, - } - except Exception as e: - logger.warning("check_workflow_status_tool failed: %s", e) - return {"error": str(e), "status": "failed"} - - -# ─── I-AK: Agent → Knowledge Tools ─────────────────────────────────────────── - - -async def ask_knowledge_tool( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - query: str, - source_types: list[str] | None = None, - max_results: int = 5, -) -> dict[str, Any]: - """Agent tool: Query the knowledge base with evidence-backed results. - - Args: - query: Natural language query. - source_types: Optional filter (wiki, dms, mail, communication). - max_results: Maximum results to return. - - Returns: - Dict with answer, evidence references, and workstream blocks. - """ - from app.ai.knowledge_lifecycle import ask_knowledge - - return await ask_knowledge( - db=db, - tenant_id=tenant_id, - user_id=user_id, - query=query, - source_types=source_types, - max_results=max_results, - ) - - -async def search_knowledge_tool( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - query: str, - entity_type: str | None = None, - limit: int = 10, -) -> dict[str, Any]: - """Agent tool: Search across all knowledge sources. - - Args: - query: Search query. - entity_type: Optional entity type filter. - limit: Maximum results. - - Returns: - Dict with search results and evidence references. - """ - try: - from app.plugins.builtins.unified_search.contracts import UnifiedSearchContract - contract = UnifiedSearchContract - search_fn = contract.get_function("unified_search") - if search_fn is None: - return {"error": "Search not available", "results": []} - - results = await search_fn( - db=db, - tenant_id=tenant_id, - query=query, - entity_type=entity_type, - limit=limit, - ) - - # Build evidence references - from app.ai.knowledge_sources import build_evidence_references - refs = build_evidence_references(results or [], max_results=limit) - - return { - "results": [r.to_dict() for r in refs], - "total": len(refs), - "query": query, - } - except Exception as e: - logger.warning("search_knowledge_tool failed: %s", e) - return {"error": str(e), "results": []} - - -# ─── Tool Registration ─────────────────────────────────────────────────────── - - -def register_integration_tools(registry: Any) -> None: - """Register integration tools in the AI tool registry. - - Called during plugin initialization to make workflow and knowledge - tools available to AI agents. - """ - # These would be registered as ToolDefinition objects in the registry. - # The actual registration depends on the ToolRegistry API. - # For now, we expose the functions for manual registration. - pass - - -__all__ = [ - "start_workflow_tool", - "check_workflow_status_tool", - "ask_knowledge_tool", - "search_knowledge_tool", - "register_integration_tools", -] diff --git a/app/ai/knowledge_extraction.py b/app/ai/knowledge_extraction.py deleted file mode 100644 index 294ab6d..0000000 --- a/app/ai/knowledge_extraction.py +++ /dev/null @@ -1,361 +0,0 @@ -"""Knowledge extraction — LLM-based relationship and entity extraction (H-EXT, H-ENT, H-AUTO, H-CONF). - -Analyzes texts from knowledge sources (wiki, DMS, mail, communication) -and extracts: -- Named entities (persons, companies, projects) — H-ENT -- Relationships between entities — H-EXT -- Auto-creates relationships in GraphRAG — H-AUTO -- Confidence scores with low-confidence → review queue — H-CONF -""" - -from __future__ import annotations - -import logging -import uuid -from dataclasses import dataclass, field -from typing import Any - -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -@dataclass -class ExtractedEntity: - """A named entity extracted from text (H-ENT).""" - - name: str - entity_type: str # person, company, project, location, date, other - mentions: list[int] = field(default_factory=list) # character positions - confidence: float = 1.0 - metadata: dict[str, Any] = field(default_factory=dict) - - -@dataclass -class ExtractedRelationship: - """A relationship between two entities extracted from text (H-EXT).""" - - source_entity: str - source_type: str - target_entity: str - target_type: str - relationship_type: str # works_for, related_to, has_email, etc. - confidence: float = 0.0 - evidence: str = "" # Text snippet that supports this relationship - metadata: dict[str, Any] = field(default_factory=dict) - - -@dataclass -class ExtractionResult: - """Result of a knowledge extraction run.""" - - entities: list[ExtractedEntity] = field(default_factory=list) - relationships: list[ExtractedRelationship] = field(default_factory=list) - source_type: str = "" - source_id: str = "" - tenant_id: str = "" - overall_confidence: float = 0.0 - - -# ─── Extraction prompt ─────────────────────────────────────────────────────── - -_EXTRACTION_SYSTEM_PROMPT = """You are a knowledge extraction assistant. Analyze the given text and extract: - -1. Named entities (persons, companies, projects, locations, dates) -2. Relationships between entities (e.g., "works_for", "related_to", "has_email", "part_of") - -Return JSON in this format: -{ - "entities": [ - {"name": "John Doe", "type": "person", "confidence": 0.95} - ], - "relationships": [ - { - "source": "John Doe", - "source_type": "person", - "target": "Acme Corp", - "target_type": "company", - "type": "works_for", - "confidence": 0.9, - "evidence": "John Doe works at Acme Corp" - } - ] -} - -Only extract explicitly stated facts. Do not infer or hallucinate. -If no entities or relationships are found, return empty arrays.""" - - -async def extract_knowledge( - text: str, - tenant_id: uuid.UUID, - source_type: str = "", - source_id: str = "", - llm_model: str | None = None, -) -> ExtractionResult: - """Extract entities and relationships from text using LLM (H-EXT, H-ENT). - - Args: - text: The text to analyze. - tenant_id: Tenant ID for multi-tenancy. - source_type: Source type (wiki, dms, mail, communication). - source_id: Source entity ID. - llm_model: Optional LLM model override. - - Returns: - ExtractionResult with entities and relationships. - """ - if not text or len(text.strip()) < 10: - return ExtractionResult( - source_type=source_type, - source_id=source_id, - tenant_id=str(tenant_id), - ) - - try: - from app.ai.llm_client import llm_complete - - messages = [ - {"role": "system", "content": _EXTRACTION_SYSTEM_PROMPT}, - {"role": "user", "content": f"Analyze this text:\n\n{text[:8000]}"}, - ] - - response = await llm_complete( - model=llm_model or "ollama/deepseek-v4-flash", - messages=messages, - temperature=0.1, - max_tokens=2000, - ) - - # Parse LLM response - import json - - raw = response.get("content", "") - # Try to extract JSON from response - try: - data = json.loads(raw) - except json.JSONDecodeError: - # Try to find JSON in the response - start = raw.find("{") - end = raw.rfind("}") + 1 - if start >= 0 and end > start: - data = json.loads(raw[start:end]) - else: - logger.warning("Failed to parse extraction response as JSON") - return ExtractionResult( - source_type=source_type, - source_id=source_id, - tenant_id=str(tenant_id), - ) - - # Build ExtractionResult - entities: list[ExtractedEntity] = [] - for ent in data.get("entities", []): - entities.append(ExtractedEntity( - name=ent.get("name", ""), - entity_type=ent.get("type", "other"), - confidence=ent.get("confidence", 0.5), - metadata=ent.get("metadata", {}), - )) - - relationships: list[ExtractedRelationship] = [] - for rel in data.get("relationships", []): - relationships.append(ExtractedRelationship( - source_entity=rel.get("source", ""), - source_type=rel.get("source_type", "other"), - target_entity=rel.get("target", ""), - target_type=rel.get("target_type", "other"), - relationship_type=rel.get("type", "related_to"), - confidence=rel.get("confidence", 0.5), - evidence=rel.get("evidence", ""), - metadata=rel.get("metadata", {}), - )) - - # Calculate overall confidence - all_confidences = [e.confidence for e in entities] + [r.confidence for r in relationships] - overall = sum(all_confidences) / len(all_confidences) if all_confidences else 0.0 - - return ExtractionResult( - entities=entities, - relationships=relationships, - source_type=source_type, - source_id=source_id, - tenant_id=str(tenant_id), - overall_confidence=overall, - ) - - except Exception as e: - logger.warning("Knowledge extraction failed: %s", e) - return ExtractionResult( - source_type=source_type, - source_id=source_id, - tenant_id=str(tenant_id), - ) - - -# ─── Confidence scoring (H-CONF) ───────────────────────────────────────────── - -LOW_CONFIDENCE_THRESHOLD = 0.6 - - -def is_low_confidence(confidence: float) -> bool: - """Check if a confidence score is below the review threshold (H-CONF).""" - return confidence < LOW_CONFIDENCE_THRESHOLD - - -def filter_high_confidence( - relationships: list[ExtractedRelationship], - threshold: float = LOW_CONFIDENCE_THRESHOLD, -) -> tuple[list[ExtractedRelationship], list[ExtractedRelationship]]: - """Split relationships into high-confidence and low-confidence (review queue). - - Returns: - Tuple of (high_confidence, low_confidence) lists. - """ - high = [r for r in relationships if r.confidence >= threshold] - low = [r for r in relationships if r.confidence < threshold] - return high, low - - -# ─── Auto-relationship creation in GraphRAG (H-AUTO) ───────────────────────── - - -async def auto_create_relationships( - db: AsyncSession, - tenant_id: uuid.UUID, - extraction: ExtractionResult, - min_confidence: float = LOW_CONFIDENCE_THRESHOLD, -) -> dict[str, Any]: - """Auto-create extracted relationships in GraphRAG (H-AUTO). - - Only creates relationships with confidence >= min_confidence. - Low-confidence relationships are returned for the review queue. - - Returns: - Dict with ``created``, ``skipped_low_confidence``, ``errors`` counts. - """ - from app.plugins.builtins.graph_rag.models import EntityRelationship - - created = 0 - skipped = 0 - errors = 0 - - high_conf, low_conf = filter_high_confidence(extraction.relationships, min_confidence) - - for rel in high_conf: - try: - # Try to resolve entity names to actual entity IDs - # For now, store as typed relationships with name-based references - source_id = await _resolve_entity_id(db, tenant_id, rel.source_entity, rel.source_type) - target_id = await _resolve_entity_id(db, tenant_id, rel.target_entity, rel.target_type) - - if source_id is None or target_id is None: - skipped += 1 - continue - - # Check if relationship already exists - existing = await db.execute( - select(EntityRelationship).where( - EntityRelationship.tenant_id == tenant_id, - EntityRelationship.source_type == rel.source_type, - EntityRelationship.source_id == source_id, - EntityRelationship.target_type == rel.target_type, - EntityRelationship.target_id == target_id, - EntityRelationship.relationship_type == rel.relationship_type, - ) - ) - if existing.scalar_one_or_none() is not None: - skipped += 1 - continue - - # Create new relationship - er = EntityRelationship( - tenant_id=tenant_id, - source_type=rel.source_type, - source_id=source_id, - target_type=rel.target_type, - target_id=target_id, - relationship_type=rel.relationship_type, - confidence=rel.confidence, - metadata={ - "evidence": rel.evidence, - "source_type": extraction.source_type, - "source_id": extraction.source_id, - "auto_extracted": True, - }, - ) - db.add(er) - created += 1 - except Exception as e: - logger.warning("Failed to auto-create relationship: %s", e) - errors += 1 - - await db.flush() - - return { - "created": created, - "skipped_low_confidence": len(low_conf), - "skipped_existing": skipped, - "errors": errors, - "review_queue": [ - { - "source": r.source_entity, - "target": r.target_entity, - "type": r.relationship_type, - "confidence": r.confidence, - "evidence": r.evidence, - } - for r in low_conf - ], - } - - -async def _resolve_entity_id( - db: AsyncSession, - tenant_id: uuid.UUID, - name: str, - entity_type: str, -) -> uuid.UUID | None: - """Try to resolve an entity name to an actual entity ID. - - Searches contacts, companies, etc. by name. - Returns None if no match found. - """ - try: - if entity_type == "person": - from app.models.contact import Contact - result = await db.execute( - select(Contact.id).where( - Contact.tenant_id == tenant_id, - Contact.deleted_at.is_(None), - Contact.name.ilike(f"%{name}%"), - ).limit(1) - ) - return result.scalar_one_or_none() - elif entity_type == "company": - from app.models.contact import Contact - result = await db.execute( - select(Contact.id).where( - Contact.tenant_id == tenant_id, - Contact.deleted_at.is_(None), - Contact.is_company.is_(True), - Contact.name.ilike(f"%{name}%"), - ).limit(1) - ) - return result.scalar_one_or_none() - except Exception: - pass - return None - - -__all__ = [ - "ExtractedEntity", - "ExtractedRelationship", - "ExtractionResult", - "extract_knowledge", - "is_low_confidence", - "filter_high_confidence", - "auto_create_relationships", - "LOW_CONFIDENCE_THRESHOLD", -] diff --git a/app/ai/knowledge_lifecycle.py b/app/ai/knowledge_lifecycle.py deleted file mode 100644 index 5846c47..0000000 --- a/app/ai/knowledge_lifecycle.py +++ /dev/null @@ -1,447 +0,0 @@ -"""Knowledge lifecycle — event-driven extraction, derived-data lifecycle, -retention policy, ask-knowledge, and review queue (H-EVT, H-DATA-LIFE, H-RET, H-ASK, H-REV). - -Event-driven extraction: new mail/dokument/message → ARQ-Job → extraction. -Derived-data lifecycle: correction/delete/erasure of source propagates to -RAG chunks, embeddings, graph references, and agent memory. -Retention: configurable per-source retention policy, ARQ cleans up. -Ask Knowledge: RAG queries with evidence cards via workstream. -Review queue: low-confidence extracted relationships pending review. -""" - -from __future__ import annotations - -import logging -import uuid -from datetime import UTC, datetime, timedelta -from typing import Any - -from sqlalchemy import select, delete -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── H-EVT: Event-Driven Extraction ────────────────────────────────────────── - -# Events that trigger knowledge extraction -EXTRACTION_TRIGGERS = { - "mail.received": {"source_type": "mail", "text_field": "body_text"}, - "dms.file_uploaded": {"source_type": "dms", "text_field": "extracted_text"}, - "wiki.article_published": {"source_type": "wiki", "text_field": "content"}, - "communication.message_created": {"source_type": "communication", "text_field": "content"}, -} - - -def should_extract(event_name: str) -> bool: - """Check if an event should trigger knowledge extraction (H-EVT).""" - return event_name in EXTRACTION_TRIGGERS - - -async def handle_extraction_event( - db: AsyncSession, - tenant_id: uuid.UUID, - event_name: str, - payload: dict[str, Any], -) -> dict[str, Any] | None: - """Handle an event that may trigger knowledge extraction (H-EVT). - - Called by the event bus. If the event matches a configured extraction - trigger, fetches the source content and runs knowledge extraction. - """ - if not should_extract(event_name): - return None - - trigger_config = EXTRACTION_TRIGGERS[event_name] - source_type = trigger_config["source_type"] - entity_id_str = payload.get("entity_id") or payload.get("file_id") or payload.get("message_id") - if not entity_id_str: - return None - - try: - entity_id = uuid.UUID(str(entity_id_str)) - except (ValueError, TypeError): - return None - - # Fetch source content - from app.ai.knowledge_sources import fetch_source_content - content = await fetch_source_content(db, tenant_id, source_type, entity_id) - if content is None or not content.get("text"): - return None - - # Run extraction - from app.ai.knowledge_extraction import extract_knowledge, auto_create_relationships - extraction = await extract_knowledge( - text=content["text"], - tenant_id=tenant_id, - source_type=source_type, - source_id=str(entity_id), - ) - - # Auto-create high-confidence relationships - result = await auto_create_relationships(db, tenant_id, extraction) - - return { - "event": event_name, - "source_type": source_type, - "source_id": str(entity_id), - "entities_found": len(extraction.entities), - "relationships_found": len(extraction.relationships), - **result, - } - - -# ─── H-DATA-LIFE: Derived-Data Lifecycle ───────────────────────────────────── - -async def propagate_source_deletion( - db: AsyncSession, - tenant_id: uuid.UUID, - source_type: str, - source_id: uuid.UUID, -) -> dict[str, int]: - """Propagate correction/delete/erasure of a source to derived data (H-DATA-LIFE). - - When a source (wiki article, DMS file, mail, message) is deleted or corrected, - this removes: - - RAG chunks referencing the source - - Embeddings referencing the source - - Graph relationships with metadata.source_id matching - - Agent memory entries referencing the source - - Returns counts of what was removed. - """ - removed = {"graph_relationships": 0, "agent_memory": 0} - - # Remove graph relationships that were auto-extracted from this source - try: - from app.plugins.builtins.graph_rag.models import EntityRelationship - - result = await db.execute( - select(EntityRelationship).where( - EntityRelationship.tenant_id == tenant_id, - EntityRelationship.metadata["source_id"].astext == str(source_id), - EntityRelationship.metadata["source_type"].astext == source_type, - EntityRelationship.metadata["auto_extracted"].astext == "true", - ) - ) - rels = result.scalars().all() - for rel in rels: - await db.delete(rel) - removed["graph_relationships"] = len(rels) - except Exception as e: - logger.warning("Failed to remove graph relationships for %s/%s: %s", source_type, source_id, e) - - # Remove agent memory entries referencing this source - try: - from app.ai.agent_memory import AgentMemory - - result = await db.execute( - select(AgentMemory).where( - AgentMemory.tenant_id == tenant_id, - AgentMemory.metadata["source_type"].astext == source_type, - AgentMemory.metadata["source_id"].astext == str(source_id), - ) - ) - memories = result.scalars().all() - for mem in memories: - await db.delete(mem) - removed["agent_memory"] = len(memories) - except Exception as e: - logger.warning("Failed to remove agent memory for %s/%s: %s", source_type, source_id, e) - - await db.flush() - return removed - - -async def propagate_source_correction( - db: AsyncSession, - tenant_id: uuid.UUID, - source_type: str, - source_id: uuid.UUID, -) -> dict[str, Any]: - """Propagate source correction — re-extract knowledge from updated content (H-DATA-LIFE). - - Removes old derived data and re-runs extraction on the updated source. - """ - # First remove old derived data - removed = await propagate_source_deletion(db, tenant_id, source_type, source_id) - - # Then re-extract from updated content - from app.ai.knowledge_sources import fetch_source_content - content = await fetch_source_content(db, tenant_id, source_type, source_id) - if content is None: - return {"removed": removed, "re_extracted": False} - - from app.ai.knowledge_extraction import extract_knowledge, auto_create_relationships - extraction = await extract_knowledge( - text=content["text"], - tenant_id=tenant_id, - source_type=source_type, - source_id=str(source_id), - ) - created = await auto_create_relationships(db, tenant_id, extraction) - - return { - "removed": removed, - "re_extracted": True, - "entities_found": len(extraction.entities), - "relationships_found": len(extraction.relationships), - **created, - } - - -# ─── H-RET: Knowledge/Memory Retention ─────────────────────────────────────── - -# Default retention per source type (days). 0 = no retention limit. -DEFAULT_RETENTION_DAYS = { - "wiki": 0, # No limit — wiki articles are persistent knowledge - "dms": 365, # 1 year for document-derived knowledge - "mail": 180, # 6 months for mail-derived knowledge - "communication": 90, # 3 months for communication-derived knowledge -} - - -def get_retention_days(source_type: str) -> int: - """Get retention period for a knowledge source type (H-RET).""" - return DEFAULT_RETENTION_DAYS.get(source_type, 180) - - -async def cleanup_expired_knowledge( - db: AsyncSession, - tenant_id: uuid.UUID, -) -> dict[str, int]: - """Clean up expired knowledge based on retention policy (H-RET). - - Called by ARQ cron job. Removes graph relationships and agent memory - entries that have exceeded their retention period. - """ - cleaned = {"graph_relationships": 0, "agent_memory": 0} - now = datetime.now(UTC) - - # Clean up expired graph relationships - try: - from app.plugins.builtins.graph_rag.models import EntityRelationship - - for source_type, retention_days in DEFAULT_RETENTION_DAYS.items(): - if retention_days == 0: - continue - cutoff = now - timedelta(days=retention_days) - result = await db.execute( - select(EntityRelationship).where( - EntityRelationship.tenant_id == tenant_id, - EntityRelationship.metadata["source_type"].astext == source_type, - EntityRelationship.metadata["auto_extracted"].astext == "true", - EntityRelationship.created_at < cutoff, - ) - ) - rels = result.scalars().all() - for rel in rels: - await db.delete(rel) - cleaned["graph_relationships"] += len(rels) - except Exception as e: - logger.warning("Failed to cleanup expired graph relationships: %s", e) - - # Clean up expired agent memory - try: - from app.ai.agent_memory import AgentMemory - - for source_type, retention_days in DEFAULT_RETENTION_DAYS.items(): - if retention_days == 0: - continue - cutoff = now - timedelta(days=retention_days) - result = await db.execute( - select(AgentMemory).where( - AgentMemory.tenant_id == tenant_id, - AgentMemory.metadata["source_type"].astext == source_type, - AgentMemory.created_at < cutoff, - ) - ) - memories = result.scalars().all() - for mem in memories: - await db.delete(mem) - cleaned["agent_memory"] += len(memories) - except Exception as e: - logger.warning("Failed to cleanup expired agent memory: %s", e) - - await db.flush() - return cleaned - - -# ─── H-ASK: Ask Knowledge in Workstream ────────────────────────────────────── - -async def ask_knowledge( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - query: str, - *, - source_types: list[str] | None = None, - max_results: int = 5, -) -> dict[str, Any]: - """Ask Knowledge — RAG query with evidence cards (H-ASK). - - Runs a unified search query, builds evidence references, - and returns results formatted for workstream display. - """ - if not query: - return {"answer": "", "evidence": [], "query": ""} - - try: - from app.plugins.builtins.unified_search.contracts import SearchContract - contract = SearchContract - search_fn = contract.get_function("unified_search") - if search_fn is None: - return {"answer": "Search not available", "evidence": [], "query": query} - - results = await search_fn( - db=db, - tenant_id=tenant_id, - query=query, - entity_type=None, - limit=max_results * 2, - ) - - # Filter by source types if specified - if source_types and results: - results = [r for r in results if r.get("source_type") in source_types] - - # Build evidence references - from app.ai.knowledge_sources import build_evidence_references - refs = build_evidence_references(results, max_results=max_results) - - # Build answer from top results - if not refs: - return {"answer": "No relevant knowledge found.", "evidence": [], "query": query} - - # Summarize top results - snippets = [f"- {r.title}: {r.snippet[:150]}" for r in refs[:3]] - answer = f"Based on {len(refs)} source(s):\n\n" + "\n".join(snippets) - - return { - "answer": answer, - "evidence": [r.to_dict() for r in refs], - "workstream_blocks": [r.to_workstream_block() for r in refs], - "query": query, - } - except Exception as e: - logger.warning("Ask knowledge failed: %s", e) - return {"answer": f"Knowledge query failed: {e}", "evidence": [], "query": query} - - -# ─── H-REV: Review Queue for extracted relationships ───────────────────────── - -async def get_review_queue( - db: AsyncSession, - tenant_id: uuid.UUID, - *, - page: int = 1, - page_size: int = 20, -) -> dict[str, Any]: - """Get low-confidence extracted relationships pending review (H-REV).""" - from app.plugins.builtins.graph_rag.models import EntityRelationship - from sqlalchemy import func - - query = select(EntityRelationship).where( - EntityRelationship.tenant_id == tenant_id, - EntityRelationship.confidence < 0.6, - EntityRelationship.metadata["auto_extracted"].astext == "true", - EntityRelationship.metadata["reviewed"].astext != "true", - ) - - count_q = select(func.count()).select_from(query.subquery()) - total = (await db.execute(count_q)).scalar() or 0 - - query = query.order_by(EntityRelationship.confidence.asc()).offset((page - 1) * page_size).limit(page_size) - result = await db.execute(query) - items = result.scalars().all() - - return { - "items": [ - { - "id": str(r.id), - "source_type": r.source_type, - "source_id": str(r.source_id), - "target_type": r.target_type, - "target_id": str(r.target_id), - "relationship_type": r.relationship_type, - "confidence": r.confidence, - "evidence": (r.metadata or {}).get("evidence", ""), - "source": (r.metadata or {}).get("source_type", ""), - } - for r in items - ], - "total": total, - "page": page, - "page_size": page_size, - } - - -async def approve_relationship( - db: AsyncSession, - tenant_id: uuid.UUID, - relationship_id: uuid.UUID, - user_id: uuid.UUID, -) -> bool: - """Approve a low-confidence relationship (H-REV). - - Marks the relationship as reviewed and boosts its confidence. - """ - from app.plugins.builtins.graph_rag.models import EntityRelationship - - result = await db.execute( - select(EntityRelationship).where( - EntityRelationship.id == relationship_id, - EntityRelationship.tenant_id == tenant_id, - ) - ) - rel = result.scalar_one_or_none() - if rel is None: - return False - - meta = dict(rel.metadata or {}) - meta["reviewed"] = True - meta["reviewed_by"] = str(user_id) - meta["reviewed_at"] = datetime.now(UTC).isoformat() - rel.metadata = meta - rel.confidence = max(rel.confidence, 0.8) # Boost confidence after review - await db.flush() - return True - - -async def reject_relationship( - db: AsyncSession, - tenant_id: uuid.UUID, - relationship_id: uuid.UUID, -) -> bool: - """Reject a low-confidence relationship — delete it (H-REV).""" - from app.plugins.builtins.graph_rag.models import EntityRelationship - - result = await db.execute( - select(EntityRelationship).where( - EntityRelationship.id == relationship_id, - EntityRelationship.tenant_id == tenant_id, - ) - ) - rel = result.scalar_one_or_none() - if rel is None: - return False - - await db.delete(rel) - await db.flush() - return True - - -__all__ = [ - "should_extract", - "handle_extraction_event", - "propagate_source_deletion", - "propagate_source_correction", - "get_retention_days", - "cleanup_expired_knowledge", - "ask_knowledge", - "get_review_queue", - "approve_relationship", - "reject_relationship", - "EXTRACTION_TRIGGERS", - "DEFAULT_RETENTION_DAYS", -] diff --git a/app/ai/knowledge_sources.py b/app/ai/knowledge_sources.py deleted file mode 100644 index 4f91547..0000000 --- a/app/ai/knowledge_sources.py +++ /dev/null @@ -1,259 +0,0 @@ -"""Knowledge source adapter — connects DMS, Wiki, Mail, Communication to the -existing SearchProvider/RAG pipeline (H-SRC). - -Originalquelle bleibt authoritative; Permissions/Sensitive Fields gelten -durchgängig. No second universal knowledge store — uses existing -unified_search infrastructure. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── Source type registry ──────────────────────────────────────────────────── - -_SOURCE_TYPES: dict[str, dict[str, Any]] = { - "wiki": { - "display_name": "Wiki Articles", - "model_path": "app.plugins.builtins.wiki.models.WikiArticle", - "text_field": "content", - "title_field": "title", - "id_field": "id", - "status_filter": {"status": "published"}, - }, - "dms": { - "display_name": "DMS Documents", - "model_path": "app.plugins.builtins.dms.models.DmsFile", - "text_field": "extracted_text", - "title_field": "filename", - "id_field": "id", - "status_filter": {}, - }, - "mail": { - "display_name": "Mail Messages", - "model_path": "app.plugins.builtins.mail.models.MailMessage", - "text_field": "body_text", - "title_field": "subject", - "id_field": "id", - "status_filter": {}, - }, - "communication": { - "display_name": "Communication Messages", - "model_path": "app.plugins.builtins.kommunikation.models.CommMessage", - "text_field": "content", - "title_field": "content", - "id_field": "id", - "status_filter": {}, - }, -} - - -def get_available_sources() -> list[dict[str, str]]: - """List all available knowledge source types.""" - return [ - {"type": k, "display_name": v["display_name"]} - for k, v in _SOURCE_TYPES.items() - ] - - -def get_source_config(source_type: str) -> dict[str, Any] | None: - """Get configuration for a knowledge source type.""" - return _SOURCE_TYPES.get(source_type) - - -async def fetch_source_content( - db: AsyncSession, - tenant_id: uuid.UUID, - source_type: str, - entity_id: uuid.UUID, -) -> dict[str, Any] | None: - """Fetch content from a knowledge source for RAG indexing. - - Returns a dict with: - - ``title``: Title for the content - - ``text``: Text content for embedding - - ``source_type``: The source type - - ``source_id``: The entity ID - - ``source_url``: Deep link to the original content - - ``metadata``: Additional metadata - """ - config = get_source_config(source_type) - if config is None: - return None - - try: - # Dynamic import of the model - import importlib - module_path, class_name = config["model_path"].rsplit(".", 1) - module = importlib.import_module(module_path) - model = getattr(module, class_name) - - # Fetch the entity - query = select(model).where( - model.id == entity_id, - model.tenant_id == tenant_id, - ) - if hasattr(model, "deleted_at"): - query = query.where(model.deleted_at.is_(None)) - - # Apply status filter - for field, value in config.get("status_filter", {}).items(): - if hasattr(model, field): - query = query.where(getattr(model, field) == value) - - result = await db.execute(query) - entity = result.scalar_one_or_none() - if entity is None: - return None - - # Extract text and title - text = getattr(entity, config["text_field"], "") or "" - title = getattr(entity, config["title_field"], "") or "" - - # Build source URL (deep link) - source_url = _build_source_url(source_type, entity_id) - - # Build metadata - metadata = { - "source_type": source_type, - "source_id": str(entity_id), - "tenant_id": str(tenant_id), - } - if hasattr(entity, "owner_id"): - metadata["owner_id"] = str(entity.owner_id) if entity.owner_id else None - if hasattr(entity, "tags"): - metadata["tags"] = entity.tags or [] - if hasattr(entity, "category_id"): - metadata["category_id"] = str(entity.category_id) if entity.category_id else None - - return { - "title": title, - "text": text, - "source_type": source_type, - "source_id": str(entity_id), - "source_url": source_url, - "metadata": metadata, - } - except Exception as e: - logger.warning("Failed to fetch %s content %s: %s", source_type, entity_id, e) - return None - - -def _build_source_url(source_type: str, entity_id: uuid.UUID) -> str: - """Build a deep-link URL to the original content.""" - url_map = { - "wiki": f"/wiki/articles/{entity_id}", - "dms": f"/dms/files/{entity_id}", - "mail": f"/mail/messages/{entity_id}", - "communication": f"/communication/messages/{entity_id}", - } - return url_map.get(source_type, f"/{source_type}/{entity_id}") - - -# ─── Evidence/Source References (H-CITE) ───────────────────────────────────── - - -class EvidenceReference: - """Structured source reference for RAG/Knowledge results (H-CITE). - - Provides deep-links and cards to original documents, mails, messages, - or business objects. Agents can display these in the workstream. - """ - - def __init__( - self, - source_type: str, - source_id: str, - title: str, - url: str, - snippet: str = "", - confidence: float = 0.0, - metadata: dict[str, Any] | None = None, - ): - self.source_type = source_type - self.source_id = source_id - self.title = title - self.url = url - self.snippet = snippet - self.confidence = confidence - self.metadata = metadata or {} - - def to_dict(self) -> dict[str, Any]: - """Serialize to dict for API responses and workstream blocks.""" - return { - "source_type": self.source_type, - "source_id": self.source_id, - "title": self.title, - "url": self.url, - "snippet": self.snippet, - "confidence": self.confidence, - "metadata": self.metadata, - } - - def to_workstream_block(self) -> dict[str, Any]: - """Convert to a typed workstream block for display in Communication.""" - return { - "type": "evidence_card", - "source_type": self.source_type, - "source_id": self.source_id, - "title": self.title, - "url": self.url, - "snippet": self.snippet[:200], - "confidence": self.confidence, - } - - -def build_evidence_references( - search_results: list[dict[str, Any]], - max_results: int = 5, -) -> list[EvidenceReference]: - """Build evidence references from search/RAG results. - - Args: - search_results: Raw search results with source_type, source_id, title, etc. - max_results: Maximum number of references to return. - - Returns: - List of EvidenceReference objects sorted by confidence. - """ - refs: list[EvidenceReference] = [] - for result in search_results[:max_results]: - source_type = result.get("source_type", "unknown") - source_id = result.get("source_id", "") - title = result.get("title", "") - url = result.get("source_url") or _build_source_url( - source_type, uuid.UUID(source_id) if source_id else uuid.uuid4() - ) - snippet = result.get("snippet", "") or result.get("text", "")[:200] - confidence = result.get("score", 0.0) - - refs.append(EvidenceReference( - source_type=source_type, - source_id=source_id, - title=title, - url=url, - snippet=snippet, - confidence=confidence, - metadata=result.get("metadata", {}), - )) - - # Sort by confidence descending - refs.sort(key=lambda r: r.confidence, reverse=True) - return refs - - -__all__ = [ - "get_available_sources", - "get_source_config", - "fetch_source_content", - "EvidenceReference", - "build_evidence_references", -] diff --git a/app/ai/mcp_exposure.py b/app/ai/mcp_exposure.py deleted file mode 100644 index e35a62e..0000000 --- a/app/ai/mcp_exposure.py +++ /dev/null @@ -1,226 +0,0 @@ -"""MCP-Exposure for platform features (I-MCP). - -Exposes Search, Agents, Workflows, and Knowledge as thin MCP-compatible -tools on top of existing tools/services. MCP possesses no own rights; -the existing auth/run-as context and normal permission checks always apply. - -This is NOT a separate MCP server — it's a thin exposure layer that -maps existing platform functions to MCP tool schemas so external -MCP clients can invoke them. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── MCP Tool Definitions ──────────────────────────────────────────────────── - -MCP_TOOLS: list[dict[str, Any]] = [ - { - "name": "search", - "description": "Search across all entities (contacts, companies, DMS, wiki, mail, communication).", - "input_schema": { - "type": "object", - "properties": { - "query": {"type": "string", "description": "Search query"}, - "entity_type": {"type": "string", "description": "Optional entity type filter"}, - "limit": {"type": "integer", "description": "Max results (default 10)", "default": 10}, - }, - "required": ["query"], - }, - "required_permission": "contacts:read", - "handler": "search", - }, - { - "name": "ask_knowledge", - "description": "Query the knowledge base with RAG and evidence-backed results.", - "input_schema": { - "type": "object", - "properties": { - "query": {"type": "string", "description": "Natural language query"}, - "source_types": { - "type": "array", - "items": {"type": "string"}, - "description": "Optional source filter (wiki, dms, mail, communication)", - }, - "max_results": {"type": "integer", "description": "Max results (default 5)", "default": 5}, - }, - "required": ["query"], - }, - "required_permission": "contacts:read", - "handler": "ask_knowledge", - }, - { - "name": "start_workflow", - "description": "Start a workflow instance by workflow ID.", - "input_schema": { - "type": "object", - "properties": { - "workflow_id": {"type": "string", "description": "Workflow definition ID"}, - "context": {"type": "object", "description": "Initial context variables"}, - }, - "required": ["workflow_id"], - }, - "required_permission": "workflows:write", - "handler": "start_workflow", - }, - { - "name": "check_workflow_status", - "description": "Check the status of a workflow instance.", - "input_schema": { - "type": "object", - "properties": { - "instance_id": {"type": "string", "description": "Workflow instance ID"}, - }, - "required": ["instance_id"], - }, - "required_permission": "workflows:read", - "handler": "check_workflow_status", - }, - { - "name": "list_agents", - "description": "List available AI agents.", - "input_schema": { - "type": "object", - "properties": {}, - }, - "required_permission": "agents:read", - "handler": "list_agents", - }, - { - "name": "create_task", - "description": "Create a task (todo, follow-up, etc.).", - "input_schema": { - "type": "object", - "properties": { - "title": {"type": "string", "description": "Task title"}, - "description": {"type": "string", "description": "Task description"}, - "priority": {"type": "string", "description": "low|medium|high|urgent", "default": "medium"}, - "entity_type": {"type": "string", "description": "Linked entity type"}, - "entity_id": {"type": "string", "description": "Linked entity ID"}, - }, - "required": ["title"], - }, - "required_permission": "tasks:write", - "handler": "create_task", - }, -] - - -def get_mcp_tools() -> list[dict[str, Any]]: - """List all available MCP tools with their schemas.""" - return [ - { - "name": t["name"], - "description": t["description"], - "input_schema": t["input_schema"], - } - for t in MCP_TOOLS - ] - - -def get_mcp_tool(name: str) -> dict[str, Any] | None: - """Get a single MCP tool definition by name.""" - return next((t for t in MCP_TOOLS if t["name"] == name), None) - - -async def execute_mcp_tool( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - tool_name: str, - arguments: dict[str, Any], - user_permissions: dict[str, Any] | None = None, -) -> dict[str, Any]: - """Execute an MCP tool — thin wrapper over existing platform functions. - - MCP possesses no own rights. The existing auth/run-as context and - normal permission checks always apply. This function checks the - user's permissions before executing the tool. - """ - tool = get_mcp_tool(tool_name) - if tool is None: - return {"error": f"Unknown MCP tool: {tool_name}", "status": "not_found"} - - # Permission check — MCP has no own rights - required_perm = tool.get("required_permission") - if required_perm and user_permissions: - from app.core.permissions import check_permission - if not check_permission(user_permissions, required_perm): - return {"error": f"Permission denied: {required_perm}", "status": "forbidden"} - - handler = tool["handler"] - - try: - if handler == "search": - from app.ai.integration_tools import search_knowledge_tool - return await search_knowledge_tool( - db=db, tenant_id=tenant_id, user_id=user_id, - query=arguments.get("query", ""), - entity_type=arguments.get("entity_type"), - limit=arguments.get("limit", 10), - ) - - elif handler == "ask_knowledge": - from app.ai.integration_tools import ask_knowledge_tool - return await ask_knowledge_tool( - db=db, tenant_id=tenant_id, user_id=user_id, - query=arguments.get("query", ""), - source_types=arguments.get("source_types"), - max_results=arguments.get("max_results", 5), - ) - - elif handler == "start_workflow": - from app.ai.integration_tools import start_workflow_tool - return await start_workflow_tool( - db=db, tenant_id=tenant_id, user_id=user_id, - workflow_id=arguments.get("workflow_id", ""), - context=arguments.get("context"), - ) - - elif handler == "check_workflow_status": - from app.ai.integration_tools import check_workflow_status_tool - return await check_workflow_status_tool( - db=db, tenant_id=tenant_id, - instance_id=arguments.get("instance_id", ""), - ) - - elif handler == "list_agents": - # List available agents — thin wrapper - from app.plugins.builtins.automation.contracts import AutomationContract - contract = AutomationContract - list_fn = contract.get_function("list_agents") - if list_fn is None: - return {"error": "Agents not available", "status": "not_available"} - agents = await list_fn(db=db, tenant_id=tenant_id, user_id=user_id) - return {"agents": agents or [], "total": len(agents or [])} - - elif handler == "create_task": - from app.plugins.builtins.tasks.services import create_task - result = await create_task( - db=db, tenant_id=tenant_id, user_id=user_id, - data=arguments, - ) - return result or {"error": "Failed to create task"} - - else: - return {"error": f"Unknown handler: {handler}", "status": "not_implemented"} - - except Exception as e: - logger.warning("MCP tool '%s' failed: %s", tool_name, e) - return {"error": str(e), "status": "failed"} - - -__all__ = [ - "MCP_TOOLS", - "get_mcp_tools", - "get_mcp_tool", - "execute_mcp_tool", -] diff --git a/app/ai/onboarding.py b/app/ai/onboarding.py deleted file mode 100644 index 8149328..0000000 --- a/app/ai/onboarding.py +++ /dev/null @@ -1,154 +0,0 @@ -"""Feature onboarding — setup wizard backend (I-ONB). - -Provides API endpoints for the setup wizard that guides users through -configuring agents, workflows, knowledge, workstreams/miniapps, -and proactive collaboration. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -async def get_onboarding_status( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, -) -> dict[str, Any]: - """Get the onboarding progress for a user. - - Returns which steps are completed: agent_created, workflow_created, - knowledge_enabled, workstream_enabled, proactive_enabled. - """ - status: dict[str, Any] = { - "steps": { - "welcome": {"completed": True, "required": True}, - "create_agent": {"completed": False, "required": True}, - "create_workflow": {"completed": False, "required": False}, - "enable_knowledge": {"completed": False, "required": False}, - "enable_workstream": {"completed": False, "required": False}, - }, - "progress_pct": 20, # Welcome is done - } - - # Check if user has created an agent - try: - from app.models.workflow import AgentDefinition - from sqlalchemy import select, func - agent_count = await db.scalar( - select(func.count(AgentDefinition.id)).where( - AgentDefinition.tenant_id == tenant_id, - AgentDefinition.created_by == user_id, - ) - ) - if agent_count and agent_count > 0: - status["steps"]["create_agent"]["completed"] = True - status["progress_pct"] += 20 - except Exception as e: - logger.warning("Onboarding agent check failed: %s", e) - - # Check if user has created a workflow - try: - from app.models.workflow import WorkflowDefinition - from sqlalchemy import select, func - wf_count = await db.scalar( - select(func.count(WorkflowDefinition.id)).where( - WorkflowDefinition.tenant_id == tenant_id, - WorkflowDefinition.created_by == user_id, - ) - ) - if wf_count and wf_count > 0: - status["steps"]["create_workflow"]["completed"] = True - status["progress_pct"] += 20 - except Exception as e: - logger.warning("Onboarding workflow check failed: %s", e) - - # Check knowledge (wiki articles) - try: - from app.plugins.builtins.wiki.models import WikiArticle - from sqlalchemy import select, func - wiki_count = await db.scalar( - select(func.count(WikiArticle.id)).where( - WikiArticle.tenant_id == tenant_id, - WikiArticle.deleted_at.is_(None), - ) - ) - if wiki_count and wiki_count > 0: - status["steps"]["enable_knowledge"]["completed"] = True - status["progress_pct"] += 20 - except Exception as e: - logger.warning("Onboarding knowledge check failed: %s", e) - - # Check workstream (communication messages) - try: - from app.plugins.builtins.kommunikation.models import CommMessage - from sqlalchemy import select, func - msg_count = await db.scalar( - select(func.count(CommMessage.id)).where( - CommMessage.tenant_id == tenant_id, - ) - ) - if msg_count and msg_count > 0: - status["steps"]["enable_workstream"]["completed"] = True - status["progress_pct"] += 20 - except Exception as e: - logger.warning("Onboarding workstream check failed: %s", e) - - return status - - -def get_onboarding_guide() -> dict[str, Any]: - """Get the onboarding guide content for the setup wizard. - - Returns step-by-step instructions for each onboarding step. - """ - return { - "steps": [ - { - "id": "welcome", - "title": "Welcome to LeoCRM", - "description": "Get started with your AI-powered CRM platform.", - "icon": "sparkles", - }, - { - "id": "create_agent", - "title": "Create Your First Agent", - "description": "Set up an AI agent to help with email triage, contact enrichment, or follow-ups.", - "icon": "bot", - "action_url": "/agents/new", - }, - { - "id": "create_workflow", - "title": "Create Your First Workflow", - "description": "Automate repetitive tasks with workflows. Start with a template or build your own.", - "icon": "workflow", - "action_url": "/workflows/new", - }, - { - "id": "enable_knowledge", - "title": "Enable Knowledge Base", - "description": "Create wiki articles and let AI find answers from your company knowledge.", - "icon": "book-open", - "action_url": "/wiki", - }, - { - "id": "enable_workstream", - "title": "Enable Human-AI Workstream", - "description": "Connect humans, agents, and workflows in a unified communication stream.", - "icon": "message-square", - "action_url": "/workstream", - }, - ], - } - - -__all__ = [ - "get_onboarding_status", - "get_onboarding_guide", -] diff --git a/app/ai/proactive_feed.py b/app/ai/proactive_feed.py deleted file mode 100644 index d08ff9f..0000000 --- a/app/ai/proactive_feed.py +++ /dev/null @@ -1,236 +0,0 @@ -"""Proactive workstream feed — contextual suggestions and actions (I-WORK-PROACTIVE). - -UI-/Domain-Trigger erzeugen kontextuelle Vorschläge/Actions im Workstream -mit Priority, Dedupe, Cooldown und User-Einstellungen. Kein störendes -Popup-/Clippy-Verhalten. -""" - -from __future__ import annotations - -import logging -import uuid -from dataclasses import dataclass, field -from datetime import UTC, datetime, timedelta -from typing import Any, Literal - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -Priority = Literal["low", "medium", "high", "urgent"] - - -@dataclass -class ProactiveSuggestion: - """A proactive suggestion/action for the workstream feed.""" - - id: str = field(default_factory=lambda: str(uuid.uuid4())) - trigger: str = "" # What triggered this (e.g. "mail.received", "contact.created") - title: str = "" - description: str = "" - priority: Priority = "medium" - action_type: str = "" # suggestion, action_required, info - action_url: str = "" # Deep link to action - entity_type: str | None = None - entity_id: str | None = None - blocks: list[dict[str, Any]] = field(default_factory=list) - created_at: datetime = field(default_factory=lambda: datetime.now(UTC)) - expires_at: datetime | None = None - metadata: dict[str, Any] = field(default_factory=dict) - - def to_dict(self) -> dict[str, Any]: - return { - "id": self.id, - "trigger": self.trigger, - "title": self.title, - "description": self.description, - "priority": self.priority, - "action_type": self.action_type, - "action_url": self.action_url, - "entity_type": self.entity_type, - "entity_id": self.entity_id, - "blocks": self.blocks, - "created_at": self.created_at.isoformat(), - "expires_at": self.expires_at.isoformat() if self.expires_at else None, - "metadata": self.metadata, - } - - -# ─── Dedupe + Cooldown ─────────────────────────────────────────────────────── - -# In-memory dedupe cache (per-tenant). In production, use Redis. -_dedupe_cache: dict[str, dict[str, datetime]] = {} - -# Default cooldown per trigger type (seconds) -DEFAULT_COOLDOWNS: dict[str, int] = { - "mail.received": 300, # 5 min between suggestions for same mail - "contact.created": 600, # 10 min - "workflow.completed": 60, # 1 min - "agent.result": 120, # 2 min - "default": 300, # 5 min default -} - - -def get_cooldown(trigger: str) -> int: - """Get cooldown period for a trigger type.""" - return DEFAULT_COOLDOWNS.get(trigger, DEFAULT_COOLDOWNS["default"]) - - -def _dedupe_key(tenant_id: uuid.UUID, trigger: str, entity_id: str | None) -> str: - """Build a dedupe key for a suggestion.""" - return f"{tenant_id}:{trigger}:{entity_id or 'none'}" - - -def is_cooled_down(tenant_id: uuid.UUID, trigger: str, entity_id: str | None = None) -> bool: - """Check if a trigger is still in cooldown (should not produce new suggestions).""" - key = _dedupe_key(tenant_id, trigger, entity_id) - tenant_cache = _dedupe_cache.get(str(tenant_id), {}) - last_seen = tenant_cache.get(key) - if last_seen is None: - return False - cooldown = get_cooldown(trigger) - return datetime.now(UTC) - last_seen < timedelta(seconds=cooldown) - - -def mark_suggested(tenant_id: uuid.UUID, trigger: str, entity_id: str | None = None) -> None: - """Mark a trigger as having produced a suggestion (for cooldown tracking).""" - key = _dedupe_key(tenant_id, trigger, entity_id) - tenant_id_str = str(tenant_id) - if tenant_id_str not in _dedupe_cache: - _dedupe_cache[tenant_id_str] = {} - _dedupe_cache[tenant_id_str][key] = datetime.now(UTC) - - -# ─── Suggestion Generators ─────────────────────────────────────────────────── - -async def generate_suggestions( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - trigger: str, - payload: dict[str, Any], -) -> list[ProactiveSuggestion]: - """Generate proactive suggestions for a trigger event (I-WORK-PROACTIVE). - - Checks cooldown, generates suggestions, and marks them as suggested. - Returns a list of ProactiveSuggestion objects. - """ - entity_id = payload.get("entity_id") or payload.get("contact_id") or payload.get("message_id") - - # Check cooldown — don't spam - if is_cooled_down(tenant_id, trigger, entity_id): - return [] - - suggestions: list[ProactiveSuggestion] = [] - - # Generate based on trigger type - if trigger == "mail.received": - suggestions.append(ProactiveSuggestion( - trigger=trigger, - title="New email received", - description=f"You received a new email from {payload.get('sender', 'unknown')}", - priority="medium", - action_type="info", - action_url=f"/mail/messages/{entity_id}" if entity_id else "", - entity_type="mail", - entity_id=entity_id, - blocks=[], - )) - - elif trigger == "contact.created": - suggestions.append(ProactiveSuggestion( - trigger=trigger, - title="New contact created", - description=f"New contact: {payload.get('name', 'Unknown')}", - priority="low", - action_type="suggestion", - action_url=f"/contacts/{entity_id}" if entity_id else "", - entity_type="contact", - entity_id=entity_id, - )) - - elif trigger == "workflow.completed": - suggestions.append(ProactiveSuggestion( - trigger=trigger, - title="Workflow completed", - description=f"Workflow '{payload.get('workflow_name', 'Unknown')}' has been completed.", - priority="medium", - action_type="info", - action_url=f"/workflows/instances/{entity_id}" if entity_id else "", - entity_type="workflow_instance", - entity_id=entity_id, - )) - - elif trigger == "agent.result": - suggestions.append(ProactiveSuggestion( - trigger=trigger, - title="Agent completed task", - description=f"Agent finished: {payload.get('summary', 'Task completed')}", - priority="medium", - action_type="action_required", - action_url=f"/agents/runs/{entity_id}" if entity_id else "", - entity_type="agent_run", - entity_id=entity_id, - )) - - # Mark as suggested (cooldown) - if suggestions: - mark_suggested(tenant_id, trigger, entity_id) - - return suggestions - - -# ─── User Settings ─────────────────────────────────────────────────────────── - - -def get_user_proactive_settings(user_id: uuid.UUID) -> dict[str, Any]: - """Get proactive feed settings for a user. - - In production, this would load from DB/user preferences. - For now, returns defaults. - """ - return { - "enabled": True, - "min_priority": "low", # Don't show suggestions below this priority - "max_per_hour": 20, # Rate limit suggestions per hour - "triggers_enabled": { - "mail.received": True, - "contact.created": True, - "workflow.completed": True, - "agent.result": True, - }, - } - - -def filter_by_user_settings( - suggestions: list[ProactiveSuggestion], - settings: dict[str, Any], -) -> list[ProactiveSuggestion]: - """Filter suggestions by user settings.""" - if not settings.get("enabled", True): - return [] - - min_priority = settings.get("min_priority", "low") - priority_order = {"low": 0, "medium": 1, "high": 2, "urgent": 3} - min_level = priority_order.get(min_priority, 0) - - triggers_enabled = settings.get("triggers_enabled", {}) - - return [ - s for s in suggestions - if priority_order.get(s.priority, 0) >= min_level - and triggers_enabled.get(s.trigger, True) - ] - - -__all__ = [ - "ProactiveSuggestion", - "generate_suggestions", - "is_cooled_down", - "mark_suggested", - "get_cooldown", - "get_user_proactive_settings", - "filter_by_user_settings", - "DEFAULT_COOLDOWNS", -] diff --git a/app/ai/self_improvement.py b/app/ai/self_improvement.py deleted file mode 100644 index 1f83b78..0000000 --- a/app/ai/self_improvement.py +++ /dev/null @@ -1,659 +0,0 @@ -"""Controlled self-improvement system (Phase J). - -Implements the improvement loop: - Observe → Detect Patterns → Propose → Draft → Evaluate → - Human Approval → Activate → Measure → Keep/Rollback - -No autonomous production code changes. All improvements go through -versioned drafts, evaluation, and human approval. -""" - -from __future__ import annotations - -import logging -import uuid -from dataclasses import dataclass, field -from datetime import UTC, datetime, timedelta -from enum import Enum -from typing import Any, Literal - -from sqlalchemy import func, select -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -# ─── Enums ─────────────────────────────────────────────────────────────────── - - -class ProposalType(str, Enum): - AGENT = "agent" - SKILL = "skill" - TRIGGER = "trigger" - WORKFLOW = "workflow" - MINIAPP_TEMPLATE = "miniapp_template" - PLUGIN_PATCH = "plugin_patch" - - -class ProposalStatus(str, Enum): - DRAFT = "draft" - EVALUATING = "evaluating" - PENDING_APPROVAL = "pending_approval" - APPROVED = "approved" - REJECTED = "rejected" - ACTIVE = "active" - ROLLED_BACK = "rolled_back" - EXPIRED = "expired" - - -class SignalType(str, Enum): - AGENT_RUN = "agent_run" - WORKFLOW_RUN = "workflow_run" - PROACTIVE_SUGGESTION = "proactive_suggestion" - AUDIT_LOG = "audit_log" - ENTITY_HISTORY = "entity_history" - USER_CORRECTION = "user_correction" - HANDOFF = "handoff" - ERROR_RETRY = "error_retry" - - -# ─── J-SIGNAL: Improvement Signals ─────────────────────────────────────────── - - -@dataclass -class ImprovementSignal: - """A referenced signal from platform usage data (J-SIGNAL). - - Uses references/aggregates instead of full PII copies. - Data minimization/exposure-policy/retention apply. - """ - id: str = field(default_factory=lambda: str(uuid.uuid4())) - signal_type: SignalType = SignalType.AGENT_RUN - source_ref: str = "" # Reference to source (e.g. "agent_run:uuid") - tenant_id: str = "" - user_id: str | None = None - timestamp: datetime = field(default_factory=lambda: datetime.now(UTC)) - outcome: str = "" # success, failure, corrected, dismissed, accepted - metadata: dict[str, Any] = field(default_factory=dict) - - def to_dict(self) -> dict[str, Any]: - return { - "id": self.id, - "signal_type": self.signal_type.value, - "source_ref": self.source_ref, - "tenant_id": self.tenant_id, - "user_id": self.user_id, - "timestamp": self.timestamp.isoformat(), - "outcome": self.outcome, - "metadata": self.metadata, - } - - -async def collect_signals( - db: AsyncSession, - tenant_id: uuid.UUID, - days: int = 30, -) -> list[ImprovementSignal]: - """Collect improvement signals from platform usage data (J-SIGNAL). - - Aggregates signals from AgentRuns, WorkflowRuns, AuditLog, and - Proactive Suggestions. Uses references, not full PII copies. - """ - since = datetime.now(UTC) - timedelta(days=days) - signals: list[ImprovementSignal] = [] - - # Agent run signals - try: - from app.models.workflow import AgentRun - result = await db.execute( - select(AgentRun).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ).limit(200) - ) - for run in result.scalars().all(): - signals.append(ImprovementSignal( - signal_type=SignalType.AGENT_RUN, - source_ref=f"agent_run:{run.id}", - tenant_id=str(tenant_id), - user_id=str(run.user_id) if run.user_id else None, - timestamp=run.created_at or datetime.now(UTC), - outcome=run.status or "unknown", - metadata={"agent_id": str(run.agent_id) if run.agent_id else None, "cost_usd": float(run.total_cost_usd or 0)}, - )) - except Exception as e: - logger.warning("Signal collection (agent_runs) failed: %s", e) - - # Workflow instance signals - try: - from app.models.workflow import WorkflowInstance - result = await db.execute( - select(WorkflowInstance).where( - WorkflowInstance.tenant_id == tenant_id, - WorkflowInstance.created_at >= since, - ).limit(200) - ) - for inst in result.scalars().all(): - signals.append(ImprovementSignal( - signal_type=SignalType.WORKFLOW_RUN, - source_ref=f"workflow_instance:{inst.id}", - tenant_id=str(tenant_id), - timestamp=inst.created_at or datetime.now(UTC), - outcome=inst.status or "unknown", - metadata={"workflow_id": str(inst.workflow_id) if inst.workflow_id else None}, - )) - except Exception as e: - logger.warning("Signal collection (workflow_instances) failed: %s", e) - - # Audit log signals (user corrections) - try: - from app.models.audit import AuditLog - result = await db.execute( - select(AuditLog).where( - AuditLog.tenant_id == tenant_id, - AuditLog.created_at >= since, - AuditLog.action.like("%.correct%"), - ).limit(100) - ) - for entry in result.scalars().all(): - signals.append(ImprovementSignal( - signal_type=SignalType.USER_CORRECTION, - source_ref=f"audit:{entry.id}", - tenant_id=str(tenant_id), - user_id=str(entry.user_id) if entry.user_id else None, - timestamp=entry.created_at or datetime.now(UTC), - outcome="corrected", - metadata={"action": entry.action}, - )) - except Exception as e: - logger.warning("Signal collection (audit) failed: %s", e) - - return signals - - -# ─── J-PATTERN: Pattern/Bottleneck Detection ──────────────────────────────── - - -@dataclass -class DetectedPattern: - """A detected pattern or bottleneck from signals (J-PATTERN).""" - id: str = field(default_factory=lambda: str(uuid.uuid4())) - pattern_type: str = "" # repetitive_sequence, frequent_corrections, rejected_suggestions, error_retries, repetitive_handoffs - description: str = "" - confidence: float = 0.0 - occurrence_count: int = 0 - evidence_refs: list[str] = field(default_factory=list) - metadata: dict[str, Any] = field(default_factory=dict) - - def to_dict(self) -> dict[str, Any]: - return { - "id": self.id, - "pattern_type": self.pattern_type, - "description": self.description, - "confidence": self.confidence, - "occurrence_count": self.occurrence_count, - "evidence_refs": self.evidence_refs, - "metadata": self.metadata, - } - - -def detect_patterns(signals: list[ImprovementSignal]) -> list[DetectedPattern]: - """Detect patterns and bottlenecks from signals (J-PATTERN). - - Identifies: repetitive sequences, frequent corrections, - rejected suggestions, retries/errors, repetitive handoffs. - """ - patterns: list[DetectedPattern] = [] - - # Group signals by type and outcome - by_type: dict[str, list[ImprovementSignal]] = {} - for s in signals: - key = f"{s.signal_type.value}:{s.outcome}" - by_type.setdefault(key, []).append(s) - - # Detect frequent errors/retries - error_signals = [s for s in signals if s.outcome in ("stopped_error", "stopped_timeout", "failed")] - if len(error_signals) >= 3: - patterns.append(DetectedPattern( - pattern_type="error_retries", - description=f"{len(error_signals)} failed agent/workflow runs detected", - confidence=min(0.9, len(error_signals) / 20), - occurrence_count=len(error_signals), - evidence_refs=[s.source_ref for s in error_signals[:10]], - metadata={"avg_per_day": len(error_signals) / 30 if len(error_signals) > 0 else 0}, - )) - - # Detect frequent user corrections - correction_signals = [s for s in signals if s.signal_type == SignalType.USER_CORRECTION] - if len(correction_signals) >= 3: - patterns.append(DetectedPattern( - pattern_type="frequent_corrections", - description=f"{len(correction_signals)} user corrections detected — agents may need tuning", - confidence=min(0.85, len(correction_signals) / 15), - occurrence_count=len(correction_signals), - evidence_refs=[s.source_ref for s in correction_signals[:10]], - )) - - # Detect repetitive handoffs - handoff_signals = [s for s in signals if s.signal_type == SignalType.HANDOFF] - if len(handoff_signals) >= 3: - patterns.append(DetectedPattern( - pattern_type="repetitive_handoffs", - description=f"{len(handoff_signals)} handoffs detected — workflow may need automation", - confidence=min(0.8, len(handoff_signals) / 10), - occurrence_count=len(handoff_signals), - evidence_refs=[s.source_ref for s in handoff_signals[:10]], - )) - - # Detect dismissed proactive suggestions - dismissed = [s for s in signals if s.signal_type == SignalType.PROACTIVE_SUGGESTION and s.outcome == "dismissed"] - if len(dismissed) >= 5: - patterns.append(DetectedPattern( - pattern_type="rejected_suggestions", - description=f"{len(dismissed)} proactive suggestions dismissed — suggestions may be too frequent or irrelevant", - confidence=min(0.75, len(dismissed) / 20), - occurrence_count=len(dismissed), - evidence_refs=[s.source_ref for s in dismissed[:10]], - )) - - return patterns - - -# ─── J-PROP: ImprovementProposal ───────────────────────────────────────────── - - -@dataclass -class ImprovementProposal: - """An improvement proposal with evidence and status (J-PROP).""" - id: str = field(default_factory=lambda: str(uuid.uuid4())) - proposal_type: ProposalType = ProposalType.AGENT - title: str = "" - description: str = "" - rationale: str = "" - expected_benefit: str = "" - risk_assessment: str = "" - status: ProposalStatus = ProposalStatus.DRAFT - evidence_refs: list[str] = field(default_factory=list) - pattern_refs: list[str] = field(default_factory=list) - draft_config: dict[str, Any] = field(default_factory=dict) - evaluation_result: dict[str, Any] = field(default_factory=dict) - measurement_before: dict[str, Any] = field(default_factory=dict) - measurement_after: dict[str, Any] = field(default_factory=dict) - created_at: datetime = field(default_factory=lambda: datetime.now(UTC)) - updated_at: datetime = field(default_factory=lambda: datetime.now(UTC)) - approved_by: str | None = None - activated_at: datetime | None = None - - def to_dict(self) -> dict[str, Any]: - return { - "id": self.id, - "proposal_type": self.proposal_type.value, - "title": self.title, - "description": self.description, - "rationale": self.rationale, - "expected_benefit": self.expected_benefit, - "risk_assessment": self.risk_assessment, - "status": self.status.value, - "evidence_refs": self.evidence_refs, - "pattern_refs": self.pattern_refs, - "draft_config": self.draft_config, - "evaluation_result": self.evaluation_result, - "measurement_before": self.measurement_before, - "measurement_after": self.measurement_after, - "created_at": self.created_at.isoformat(), - "updated_at": self.updated_at.isoformat(), - "approved_by": self.approved_by, - "activated_at": self.activated_at.isoformat() if self.activated_at else None, - } - - -def create_proposal( - pattern: DetectedPattern, - proposal_type: ProposalType = ProposalType.AGENT, - title: str = "", - description: str = "", - draft_config: dict[str, Any] | None = None, -) -> ImprovementProposal: - """Create an improvement proposal from a detected pattern (J-PROP).""" - return ImprovementProposal( - proposal_type=proposal_type, - title=title or f"Improve: {pattern.pattern_type}", - description=description or pattern.description, - rationale=f"Based on {pattern.occurrence_count} occurrences with {pattern.confidence:.0%} confidence", - expected_benefit="Reduce manual effort and improve accuracy", - risk_assessment="Low — versioned draft with rollback capability", - evidence_refs=pattern.evidence_refs, - pattern_refs=[pattern.id], - draft_config=draft_config or {}, - ) - - -# ─── J-DRAFT: Versioned Draft ──────────────────────────────────────────────── - - -@dataclass -class VersionedDraft: - """A versioned draft of an agent/workflow/skill config (J-DRAFT).""" - id: str = field(default_factory=lambda: str(uuid.uuid4())) - proposal_id: str = "" - version: int = 1 - config: dict[str, Any] = field(default_factory=dict) - previous_version_id: str | None = None - created_at: datetime = field(default_factory=lambda: datetime.now(UTC)) - - def to_dict(self) -> dict[str, Any]: - return { - "id": self.id, - "proposal_id": self.proposal_id, - "version": self.version, - "config": self.config, - "previous_version_id": self.previous_version_id, - "created_at": self.created_at.isoformat(), - } - - -def create_draft(proposal: ImprovementProposal, previous_draft: VersionedDraft | None = None) -> VersionedDraft: - """Create a versioned draft from a proposal (J-DRAFT).""" - version = (previous_draft.version + 1) if previous_draft else 1 - return VersionedDraft( - proposal_id=proposal.id, - version=version, - config=proposal.draft_config, - previous_version_id=previous_draft.id if previous_draft else None, - ) - - -# ─── J-EVAL: Evaluation/Sandbox ────────────────────────────────────────────── - - -async def evaluate_proposal( - proposal: ImprovementProposal, - draft: VersionedDraft, - historical_signals: list[ImprovementSignal] | None = None, -) -> dict[str, Any]: - """Evaluate a proposal via dry-run/simulation (J-EVAL). - - No external side effects. Tests against historical/synthetic cases. - """ - result: dict[str, Any] = { - "proposal_id": proposal.id, - "draft_id": draft.id, - "evaluated_at": datetime.now(UTC).isoformat(), - "test_cases": 0, - "passed": 0, - "failed": 0, - "score": 0.0, - "recommendation": "", - "details": [], - } - - # Simulate against historical signals - test_signals = historical_signals or [] - result["test_cases"] = len(test_signals) - - for signal in test_signals: - # Simulate: would the new config have handled this better? - # This is a simplified evaluation — real implementation would - # replay the signal through the new config - if signal.outcome in ("stopped_error", "stopped_timeout", "failed"): - # Assume new config would fix 60% of errors - result["passed"] += 1 - else: - result["passed"] += 1 - - result["failed"] = result["test_cases"] - result["passed"] - result["score"] = (result["passed"] / result["test_cases"] * 100) if result["test_cases"] > 0 else 0.0 - - if result["score"] >= 80: - result["recommendation"] = "approve" - elif result["score"] >= 60: - result["recommendation"] = "approve_with_caution" - else: - result["recommendation"] = "reject" - - return result - - -# ─── J-APPROVAL: Human Approval ────────────────────────────────────────────── - - -async def request_approval( - db: AsyncSession, - tenant_id: uuid.UUID, - user_id: uuid.UUID, - proposal: ImprovementProposal, - evaluation: dict[str, Any], -) -> dict[str, Any]: - """Request human approval for a proposal (J-APPROVAL). - - Uses the central ApprovalRequest system. The approver sees - evidence, diff, tests, and expected impact. - """ - try: - from app.core.approval import create_approval_request - approval = await create_approval_request( - db=db, - tenant_id=tenant_id, - entity_type="improvement_proposal", - entity_id=uuid.UUID(proposal.id) if _is_valid_uuid(proposal.id) else uuid.uuid4(), - action=f"activate:{proposal.proposal_type.value}", - requested_by=user_id, - requested_by_type="system", - ) - proposal.status = ProposalStatus.PENDING_APPROVAL - proposal.updated_at = datetime.now(UTC) - return { - "approval_id": str(approval.id), - "proposal_id": proposal.id, - "status": "pending_approval", - "evaluation": evaluation, - } - except Exception as e: - logger.warning("Approval request failed: %s", e) - return {"error": str(e), "status": "failed"} - - -def _is_valid_uuid(s: str) -> bool: - try: - uuid.UUID(s) - return True - except (ValueError, AttributeError): - return False - - -# ─── J-ACTIVATE: Controlled Activate + Rollback ────────────────────────────── - - -async def activate_proposal( - db: AsyncSession, - tenant_id: uuid.UUID, - proposal: ImprovementProposal, - draft: VersionedDraft, -) -> dict[str, Any]: - """Atomically activate an approved proposal (J-ACTIVATE). - - Previous version remains rollback-capable. - """ - if proposal.status != ProposalStatus.APPROVED: - return {"error": "Proposal must be approved before activation", "status": "rejected"} - - try: - # Apply the draft config to the target system - # This would update the agent/workflow/skill definition - proposal.status = ProposalStatus.ACTIVE - proposal.activated_at = datetime.now(UTC) - proposal.updated_at = datetime.now(UTC) - - return { - "proposal_id": proposal.id, - "draft_id": draft.id, - "status": "active", - "activated_at": proposal.activated_at.isoformat(), - "rollback_available": True, - "previous_version_id": draft.previous_version_id, - } - except Exception as e: - logger.warning("Activation failed: %s", e) - return {"error": str(e), "status": "failed"} - - -async def rollback_proposal( - db: AsyncSession, - tenant_id: uuid.UUID, - proposal: ImprovementProposal, - previous_draft: VersionedDraft | None = None, -) -> dict[str, Any]: - """Rollback an active proposal to its previous version (J-ACTIVATE).""" - if proposal.status != ProposalStatus.ACTIVE: - return {"error": "Only active proposals can be rolled back", "status": "rejected"} - - try: - proposal.status = ProposalStatus.ROLLED_BACK - proposal.updated_at = datetime.now(UTC) - - return { - "proposal_id": proposal.id, - "status": "rolled_back", - "previous_version_id": previous_draft.id if previous_draft else None, - "rolled_back_at": datetime.now(UTC).isoformat(), - } - except Exception as e: - logger.warning("Rollback failed: %s", e) - return {"error": str(e), "status": "failed"} - - -# ─── J-MEASURE: Pre/Post Impact Measurement ───────────────────────────────── - - -async def measure_impact( - db: AsyncSession, - tenant_id: uuid.UUID, - proposal: ImprovementProposal, - days: int = 7, -) -> dict[str, Any]: - """Measure pre/post impact of an activated proposal (J-MEASURE). - - Compares time, errors, acceptance rate, cost, throughput, - and business outcome metrics. - """ - if not proposal.activated_at: - return {"error": "Proposal has not been activated", "status": "not_active"} - - since_activation = proposal.activated_at - before_start = since_activation - timedelta(days=days) - - measurement: dict[str, Any] = { - "proposal_id": proposal.id, - "measured_at": datetime.now(UTC).isoformat(), - "period_days": days, - "before": proposal.measurement_before, - "after": {}, - "delta": {}, - } - - # Collect post-activation metrics - try: - from app.models.workflow import AgentRun - - # Post-activation metrics - post_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since_activation, - ) - ) - post_errors = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since_activation, - AgentRun.status.in_(["stopped_error", "stopped_timeout"]), - ) - ) - post_cost = await db.scalar( - select(func.sum(AgentRun.total_cost_usd)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since_activation, - ) - ) - - measurement["after"] = { - "total_runs": post_runs or 0, - "errors": post_errors or 0, - "cost_usd": float(post_cost or 0), - "error_rate": (post_errors / post_runs * 100) if post_runs else 0.0, - } - - # Calculate delta - before = proposal.measurement_before - if before: - measurement["delta"] = { - "runs_change": (post_runs or 0) - before.get("total_runs", 0), - "errors_change": (post_errors or 0) - before.get("errors", 0), - "cost_change": float(post_cost or 0) - before.get("cost_usd", 0), - "error_rate_change": ((post_errors / post_runs * 100) if post_runs else 0) - before.get("error_rate", 0), - } - - except Exception as e: - measurement["error"] = str(e) - - return measurement - - -async def capture_baseline( - db: AsyncSession, - tenant_id: uuid.UUID, - days: int = 7, -) -> dict[str, Any]: - """Capture baseline metrics before activating a proposal (J-MEASURE).""" - since = datetime.now(UTC) - timedelta(days=days) - try: - from app.models.workflow import AgentRun - total_runs = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ) - ) - errors = await db.scalar( - select(func.count(AgentRun.id)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - AgentRun.status.in_(["stopped_error", "stopped_timeout"]), - ) - ) - cost = await db.scalar( - select(func.sum(AgentRun.total_cost_usd)).where( - AgentRun.tenant_id == tenant_id, - AgentRun.created_at >= since, - ) - ) - return { - "total_runs": total_runs or 0, - "errors": errors or 0, - "cost_usd": float(cost or 0), - "error_rate": (errors / total_runs * 100) if total_runs else 0.0, - "captured_at": datetime.now(UTC).isoformat(), - } - except Exception as e: - return {"error": str(e)} - - -__all__ = [ - "ProposalType", - "ProposalStatus", - "SignalType", - "ImprovementSignal", - "DetectedPattern", - "ImprovementProposal", - "VersionedDraft", - "collect_signals", - "detect_patterns", - "create_proposal", - "create_draft", - "evaluate_proposal", - "request_approval", - "activate_proposal", - "rollback_proposal", - "measure_impact", - "capture_baseline", -] diff --git a/app/ai/workstream_contract.py b/app/ai/workstream_contract.py deleted file mode 100644 index a480a93..0000000 --- a/app/ai/workstream_contract.py +++ /dev/null @@ -1,100 +0,0 @@ -"""Workstream contract — unified posting path for Human/System/Agent/Workflow (I-WORK-BASE, I-WORK-ACTOR, I-WORK-HANDOFF). - -Defines the contract for posting typed blocks to the central Communication -system. All actors use the same posting path with typed blocks. -""" - -from __future__ import annotations - -import logging -import uuid -from dataclasses import dataclass, field -from typing import Any, Literal - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - -ActorType = Literal["human", "system", "agent", "workflow"] -BlockType = Literal["text", "entity_card", "action_card", "evidence_card", "approval_card", "miniapp", "workflow_status", "workflow_handoff", "error"] - - -@dataclass -class WorkstreamBlock: - type: BlockType - content: str = "" - metadata: dict[str, Any] = field(default_factory=dict) - - def to_dict(self) -> dict[str, Any]: - return {"type": self.type, "content": self.content, "metadata": self.metadata} - - -@dataclass -class WorkstreamMessage: - actor_type: ActorType - actor_id: str | None = None - content: str = "" - blocks: list[WorkstreamBlock] = field(default_factory=list) - conversation_id: str | None = None - tenant_id: str | None = None - - def to_dict(self) -> dict[str, Any]: - return { - "actor_type": self.actor_type, - "actor_id": self.actor_id, - "content": self.content, - "blocks": [b.to_dict() for b in self.blocks], - "conversation_id": self.conversation_id, - "tenant_id": self.tenant_id, - } - - -async def post_to_workstream(db: AsyncSession, tenant_id: uuid.UUID, message: WorkstreamMessage) -> dict[str, Any] | None: - """Post a message to the central workstream (I-WORK-ACTOR).""" - try: - from app.plugins.builtins.kommunikation.contracts import KommunikationContract - contract = KommunikationContract - post_fn = contract.get_function("post_message") - if post_fn is None: - if message.actor_type == "human" and message.actor_id: - from app.core.notifications import post_system_message - await post_system_message(db, tenant_id, uuid.UUID(message.actor_id), "workstream", message.content[:200], message.content) - return None - return await post_fn(db=db, tenant_id=tenant_id, sender_id=uuid.UUID(message.actor_id) if message.actor_id else None, sender_type=message.actor_type, message_type=f"workstream_{message.blocks[0].type}" if message.blocks else "workstream_text", content=message.content, blocks=[b.to_dict() for b in message.blocks], conversation_id=message.conversation_id) - except Exception as e: - logger.warning("Failed to post to workstream: %s", e) - return None - - -async def create_handoff(db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, *, handoff_type: str, assignee_type: str = "user", assignee_id: str | None = None, entity_type: str | None = None, entity_id: str | None = None, description: str = "", agent_run_id: str | None = None, workflow_instance_id: str | None = None, conversation_id: str | None = None) -> dict[str, Any]: - """Create a Human<->Agent handoff (I-WORK-HANDOFF). Creates a Task with task_type='handoff'.""" - from app.plugins.builtins.tasks.services import create_task - task_data: dict[str, Any] = {"title": f"Handoff: {handoff_type}", "description": description, "task_type": "handoff", "assignee_type": assignee_type, "assignee_id": assignee_id, "entity_type": entity_type, "entity_id": entity_id, "status": "open", "priority": "medium"} - task = await create_task(db, tenant_id, user_id, task_data) - handoff_block = WorkstreamBlock(type="workflow_handoff", content=description, metadata={"handoff_type": handoff_type, "task_id": task.get("id") if task else None, "assignee_type": assignee_type, "assignee_id": assignee_id, "entity_type": entity_type, "entity_id": entity_id, "agent_run_id": agent_run_id, "workflow_instance_id": workflow_instance_id}) - message = WorkstreamMessage(actor_type="system", content=f"Handoff: {handoff_type} - {description}", blocks=[handoff_block], conversation_id=conversation_id, tenant_id=str(tenant_id)) - post_result = await post_to_workstream(db, tenant_id, message) - return {"task": task, "workstream_post": post_result, "handoff_type": handoff_type} - - -def build_entity_card(entity_type: str, entity_id: str, title: str = "", subtitle: str = "", url: str = "") -> WorkstreamBlock: - return WorkstreamBlock(type="entity_card", content=title, metadata={"entity_type": entity_type, "entity_id": entity_id, "title": title, "subtitle": subtitle, "url": url}) - - -def build_action_card(title: str, actions: list[dict[str, str]], description: str = "") -> WorkstreamBlock: - return WorkstreamBlock(type="action_card", content=title, metadata={"title": title, "description": description, "actions": actions}) - - -def build_evidence_card(source_type: str, source_id: str, title: str, snippet: str = "", url: str = "", confidence: float = 0.0) -> WorkstreamBlock: - return WorkstreamBlock(type="evidence_card", content=title, metadata={"source_type": source_type, "source_id": source_id, "title": title, "snippet": snippet[:200], "url": url, "confidence": confidence}) - - -def build_approval_card(approval_id: str, action: str, description: str = "") -> WorkstreamBlock: - return WorkstreamBlock(type="approval_card", content=f"Approval needed: {action}", metadata={"approval_id": approval_id, "action": action, "description": description}) - - -def build_miniapp_block(app_id: str, title: str = "", render_schema: dict[str, Any] | None = None) -> WorkstreamBlock: - return WorkstreamBlock(type="miniapp", content=title, metadata={"app_id": app_id, "title": title, "render_schema": render_schema or {}}) - - -__all__ = ["ActorType", "BlockType", "WorkstreamBlock", "WorkstreamMessage", "post_to_workstream", "create_handoff", "build_entity_card", "build_action_card", "build_evidence_card", "build_approval_card", "build_miniapp_block"] diff --git a/app/main.py b/app/main.py index 1a42ab4..c8d4e42 100644 --- a/app/main.py +++ b/app/main.py @@ -73,7 +73,6 @@ from app.routes import ( # noqa: E402 webhooks, workflows, workspaces, - platform, ) # ── Graceful shutdown signal ───────────────────────────────────────────────── @@ -581,7 +580,6 @@ def create_app() -> FastAPI: app.include_router(outbox.router) app.include_router(api_tokens.router) app.include_router(approvals.router) - app.include_router(platform.router) # ── Register plugin routes for all discovered plugins ── # Routes are registered at app creation time so OpenAPI docs are complete. diff --git a/app/routes/platform.py b/app/routes/platform.py deleted file mode 100644 index c6bfeb1..0000000 --- a/app/routes/platform.py +++ /dev/null @@ -1,161 +0,0 @@ -"""Platform routes — dashboard, onboarding, improvement, and DSGVO endpoints. - -Phase G-J platform-level API routes for platform dashboard, cost tracking, -usage analytics, onboarding status, improvement proposals/patterns, -and DSGVO data export. -""" - -from __future__ import annotations - -import uuid -from typing import Any - -from fastapi import APIRouter, Depends -from sqlalchemy.ext.asyncio import AsyncSession - -from app.core.db import get_db -from app.deps import get_current_user - -router = APIRouter(prefix="/api/v1", tags=["platform"]) - - -# ── Platform Dashboard ────────────────────────────────────────────────────── - -@router.get("/dashboard/platform") -async def get_platform_dashboard( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Get platform-level dashboard data: agents, workflows, search, knowledge, cost.""" - return { - "agents": {"active": 0, "total_runs": 0, "recent_runs_7d": 0}, - "workflows": {"active": 0, "running_instances": 0, "completed_instances": 0}, - "search": {"total_queries": 0, "avg_latency_ms": 0}, - "knowledge": {"wiki_articles": 0, "coverage": 0}, - "cost": {"total_cost_30d": 0, "budget_utilization": 0}, - "system_health": "healthy", - } - - -@router.get("/dashboard/cost") -async def get_cost_dashboard( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Get cost tracking dashboard data.""" - return { - "total_cost_30d": 0, - "budget_utilization": 0, - "by_service": {}, - "trend": [], - } - - -@router.get("/dashboard/usage") -async def get_usage_analytics( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Get usage analytics data.""" - return { - "active_users": 0, - "total_requests": 0, - "by_endpoint": {}, - "trend": [], - } - - -# ── Onboarding ────────────────────────────────────────────────────────────── - -@router.get("/onboarding/status") -async def get_onboarding_status( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Get onboarding status for the current tenant.""" - return { - "completed": False, - "steps": { - "welcome": True, - "first_agent": False, - "first_workflow": False, - "knowledge_workstream": False, - }, - } - - -@router.get("/onboarding/guide") -async def get_onboarding_guide( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Get onboarding guide content.""" - return { - "steps": [ - {"id": "welcome", "title": "Welcome", "description": "Get started with LeoCRM"}, - {"id": "first_agent", "title": "Create your first Agent", "description": "Set up an AI agent"}, - {"id": "first_workflow", "title": "Create your first Workflow", "description": "Automate a process"}, - {"id": "knowledge", "title": "Enable Knowledge & Workstream", "description": "Connect knowledge sources"}, - ], - } - - -# ── Improvement ───────────────────────────────────────────────────────────── - -@router.get("/improvement/proposals") -async def list_improvement_proposals( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> list[dict[str, Any]]: - """List improvement proposals (stub — returns empty list).""" - return [] - - -@router.get("/improvement/patterns") -async def list_improvement_patterns( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> list[dict[str, Any]]: - """List detected improvement patterns (stub — returns empty list).""" - return [] - - -# ── DSGVO ─────────────────────────────────────────────────────────────────── - -@router.get("/dsgvo/export/{user_id}") -async def export_user_data( - user_id: uuid.UUID, - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Export all data associated with a user (DSGVO/GDPR right to data portability).""" - return { - "user_id": str(user_id), - "exported_by": current_user.get("user_id"), - "data": { - "contacts": [], - "companies": [], - "emails": [], - "documents": [], - "calendar_events": [], - "tasks": [], - "audit_logs": [], - }, - } - - -@router.get("/dsgvo/compliance-export") -async def export_compliance_evidence( - db: AsyncSession = Depends(get_db), - current_user: dict[str, Any] = Depends(get_current_user), -) -> dict[str, Any]: - """Export compliance evidence for the current tenant.""" - return { - "tenant_id": current_user.get("tenant_id"), - "exported_by": current_user.get("user_id"), - "evidence": { - "audit_logs": [], - "consent_records": [], - "data_retention_policies": [], - }, - } diff --git a/app/workflows/workstream.py b/app/workflows/workstream.py deleted file mode 100644 index 55f7c56..0000000 --- a/app/workflows/workstream.py +++ /dev/null @@ -1,221 +0,0 @@ -"""Workflow workstream integration — posts workflow events to the central -Communication system (G-WORK). - -Replaces the old notification-based workflow messages with typed -CommMessage blocks: status, handoff, approval, action, and error. - -Used by the WorkflowEngine to post step transitions, approvals, -errors, and completions to the workstream. -""" - -from __future__ import annotations - -import logging -import uuid -from typing import Any - -from sqlalchemy.ext.asyncio import AsyncSession - -logger = logging.getLogger(__name__) - - -async def post_workflow_status( - db: AsyncSession, - tenant_id: uuid.UUID, - instance_id: uuid.UUID, - workflow_name: str, - status: str, - step_index: int | None = None, - step_name: str | None = None, - user_id: uuid.UUID | None = None, -) -> dict[str, Any] | None: - """Post a workflow status update to the workstream. - - Creates a CommMessage with a typed ``workflow_status`` block. - """ - try: - from app.plugins.builtins.kommunikation.contracts import KommunikationContract - contract = KommunikationContract - post_fn = contract.get_function("post_message") - if post_fn is None: - # Fallback to system notification - from app.core.notifications import post_system_message - if user_id: - await post_system_message( - db, tenant_id, user_id, "workflow_status", - f"Workflow: {workflow_name}", - f"Status: {status}" + (f" (Step: {step_name})" if step_name else ""), - ) - return None - - block = { - "type": "workflow_status", - "workflow_name": workflow_name, - "instance_id": str(instance_id), - "status": status, - "step_index": step_index, - "step_name": step_name, - } - return await post_fn( - db=db, - tenant_id=tenant_id, - sender_id=None, # System sender - sender_type="system", - message_type="workflow_status", - content=f"Workflow '{workflow_name}' → {status}", - blocks=[block], - ) - except Exception as e: - logger.warning("Failed to post workflow status to workstream: %s", e) - return None - - -async def post_workflow_handoff( - db: AsyncSession, - tenant_id: uuid.UUID, - instance_id: uuid.UUID, - workflow_name: str, - handoff_type: str, # review_needed, action_required, waiting_for_user - assignee_id: uuid.UUID | None = None, - assignee_type: str = "user", - description: str = "", - user_id: uuid.UUID | None = None, -) -> dict[str, Any] | None: - """Post a workflow handoff to the workstream. - - Creates a CommMessage with a typed ``workflow_handoff`` block. - The handoff indicates that the workflow is waiting for human input. - """ - try: - from app.plugins.builtins.kommunikation.contracts import KommunikationContract - contract = KommunikationContract - post_fn = contract.get_function("post_message") - if post_fn is None: - from app.core.notifications import post_system_message - if assignee_id: - await post_system_message( - db, tenant_id, assignee_id, "workflow_handoff", - f"Workflow Handoff: {workflow_name}", - f"{handoff_type}: {description}", - ) - return None - - block = { - "type": "workflow_handoff", - "workflow_name": workflow_name, - "instance_id": str(instance_id), - "handoff_type": handoff_type, - "assignee_id": str(assignee_id) if assignee_id else None, - "assignee_type": assignee_type, - "description": description, - } - return await post_fn( - db=db, - tenant_id=tenant_id, - sender_id=None, - sender_type="system", - message_type="workflow_handoff", - content=f"Workflow '{workflow_name}' → {handoff_type}", - blocks=[block], - ) - except Exception as e: - logger.warning("Failed to post workflow handoff to workstream: %s", e) - return None - - -async def post_workflow_error( - db: AsyncSession, - tenant_id: uuid.UUID, - instance_id: uuid.UUID, - workflow_name: str, - error: str, - step_index: int | None = None, - step_name: str | None = None, - user_id: uuid.UUID | None = None, -) -> dict[str, Any] | None: - """Post a workflow error to the workstream.""" - try: - from app.plugins.builtins.kommunikation.contracts import KommunikationContract - contract = KommunikationContract - post_fn = contract.get_function("post_message") - if post_fn is None: - from app.core.notifications import post_system_message - if user_id: - await post_system_message( - db, tenant_id, user_id, "workflow_error", - f"Workflow Error: {workflow_name}", - f"Error at step {step_name or step_index}: {error}", - ) - return None - - block = { - "type": "workflow_error", - "workflow_name": workflow_name, - "instance_id": str(instance_id), - "error": error, - "step_index": step_index, - "step_name": step_name, - } - return await post_fn( - db=db, - tenant_id=tenant_id, - sender_id=None, - sender_type="system", - message_type="workflow_error", - content=f"Workflow '{workflow_name}' → Error: {error}", - blocks=[block], - ) - except Exception as e: - logger.warning("Failed to post workflow error to workstream: %s", e) - return None - - -async def post_workflow_completed( - db: AsyncSession, - tenant_id: uuid.UUID, - instance_id: uuid.UUID, - workflow_name: str, - result: dict[str, Any] | None = None, - user_id: uuid.UUID | None = None, -) -> dict[str, Any] | None: - """Post a workflow completion to the workstream.""" - try: - from app.plugins.builtins.kommunikation.contracts import KommunikationContract - contract = KommunikationContract - post_fn = contract.get_function("post_message") - if post_fn is None: - from app.core.notifications import post_system_message - if user_id: - await post_system_message( - db, tenant_id, user_id, "workflow_completed", - f"Workflow Completed: {workflow_name}", - f"Workflow '{workflow_name}' has been completed successfully.", - ) - return None - - block = { - "type": "workflow_completed", - "workflow_name": workflow_name, - "instance_id": str(instance_id), - "result": result or {}, - } - return await post_fn( - db=db, - tenant_id=tenant_id, - sender_id=None, - sender_type="system", - message_type="workflow_completed", - content=f"Workflow '{workflow_name}' → Completed", - blocks=[block], - ) - except Exception as e: - logger.warning("Failed to post workflow completion to workstream: %s", e) - return None - - -__all__ = [ - "post_workflow_status", - "post_workflow_handoff", - "post_workflow_error", - "post_workflow_completed", -] diff --git a/frontend/src/api/improvement.ts b/frontend/src/api/improvement.ts deleted file mode 100644 index 4297b3d..0000000 --- a/frontend/src/api/improvement.ts +++ /dev/null @@ -1,159 +0,0 @@ -/** - * Controlled Self-Improvement API client (Phase J). - * - * All requests use the shared `apiClient` (`baseURL: '/api/v1'`) and target - * the self-improvement routes under `/improvement/...`. - */ - -import { apiDelete, apiGet, apiPost, apiPut } from './client'; - -// ─── Types ─── - -export type ProposalType = - | 'agent' - | 'skill' - | 'trigger' - | 'workflow' - | 'miniapp_template' - | 'plugin_patch'; - -export type ProposalStatus = - | 'draft' - | 'evaluating' - | 'pending_approval' - | 'approved' - | 'rejected' - | 'active' - | 'rolled_back' - | 'expired'; - -export type SignalType = - | 'agent_run' - | 'workflow_run' - | 'proactive_suggestion' - | 'audit_log' - | 'entity_history' - | 'user_correction' - | 'handoff' - | 'error_retry'; - -export interface ImprovementSignal { - id: string; - signal_type: SignalType; - source_ref: string; - tenant_id: string; - user_id?: string | null; - timestamp: string; - outcome: string; - metadata: Record; -} - -export interface DetectedPattern { - id: string; - pattern_type: string; - description: string; - confidence: number; - occurrence_count: number; - evidence_refs: string[]; - metadata: Record; -} - -export interface EvaluationResult { - proposal_id?: string; - draft_id?: string; - evaluated_at?: string; - test_cases: number; - passed: number; - failed: number; - score: number; - recommendation: string; - details?: unknown[]; - [key: string]: unknown; -} - -export interface ImprovementProposal { - id: string; - proposal_type: ProposalType; - title: string; - description: string; - rationale: string; - expected_benefit: string; - risk_assessment: string; - status: ProposalStatus; - evidence_refs: string[]; - pattern_refs: string[]; - draft_config: Record; - evaluation_result: EvaluationResult; - measurement_before: Record; - measurement_after: Record; - created_at: string; - updated_at: string; - approved_by?: string | null; - activated_at?: string | null; -} - -export interface VersionedDraft { - id: string; - proposal_id: string; - version: number; - config: Record; - previous_version_id?: string | null; - created_at: string; -} - -export interface ProposalActionResponse { - proposal_id: string; - status: string; - error?: string; - approval_id?: string; - activated_at?: string; - rolled_back_at?: string; - rollback_available?: boolean; - previous_version_id?: string | null; - evaluation?: EvaluationResult; - [key: string]: unknown; -} - -// ─── Proposals ─── - -export const fetchProposals = (status?: ProposalStatus) => - apiGet('/improvement/proposals', { - params: status ? { status } : {}, - }); - -export const fetchProposal = (id: string) => - apiGet(`/improvement/proposals/${id}`); - -export const approveProposal = (id: string) => - apiPost(`/improvement/proposals/${id}/approve`); - -export const rejectProposal = (id: string) => - apiPost(`/improvement/proposals/${id}/reject`); - -export const rollbackProposal = (id: string) => - apiPost(`/improvement/proposals/${id}/rollback`); - -export const activateProposal = (id: string) => - apiPost(`/improvement/proposals/${id}/activate`); - -export const deleteProposal = (id: string) => - apiDelete<{ status: string }>(`/improvement/proposals/${id}`); - -// ─── Patterns ─── - -export const fetchPatterns = () => apiGet('/improvement/patterns'); - -// ─── Signals ─── - -export const fetchSignals = (days?: number) => - apiGet('/improvement/signals', { - params: days ? { days } : {}, - }); - -// ─── Drafts ─── - -export const fetchDraft = (proposalId: string) => - apiGet(`/improvement/proposals/${proposalId}/draft`); - -export const updateDraft = (proposalId: string, config: Record) => - apiPut(`/improvement/proposals/${proposalId}/draft`, { config }); diff --git a/frontend/src/api/miniapps.ts b/frontend/src/api/miniapps.ts deleted file mode 100644 index 503563f..0000000 --- a/frontend/src/api/miniapps.ts +++ /dev/null @@ -1,87 +0,0 @@ -/** - * MiniApp manifest types + workstream block types (Phase I.2 Workstream & MiniApps). - * - * Mirrors backend contracts: - * - app/plugins/manifest.py → MiniAppContribution - * - app/ai/workstream_contract.py → WorkstreamBlock / WorkstreamMessage - * - app/ai/proactive_feed.py → ProactiveSuggestion - */ - -// ── MiniApp manifest contribution (backend: MiniAppContribution) ── - -export interface MiniAppContribution { - app_id: string; - name: string; - icon: string; - description: string; - render_schema: Record; -} - -/** - * Plugin manifest with optional `miniapps` field. - * Only the fields relevant to MiniApps are declared here; the full manifest - * is typed in store/pluginStore.ts (PluginUiManifest). - */ -export interface PluginManifestWithMiniApps { - name: string; - display_name: string; - version: string; - miniapps?: MiniAppContribution[]; -} - -// ── Workstream block types (backend: WorkstreamBlock) ── - -export type WorkstreamBlockType = - | 'text' - | 'entity_card' - | 'action_card' - | 'evidence_card' - | 'approval_card' - | 'miniapp' - | 'workflow_status' - | 'workflow_handoff' - | 'error'; - -export interface WorkstreamBlock { - type: WorkstreamBlockType; - content: string; - metadata: Record; -} - -export type WorkstreamActorType = 'human' | 'system' | 'agent' | 'workflow'; - -export interface WorkstreamMessage { - actor_type: WorkstreamActorType; - actor_id?: string | null; - content: string; - blocks: WorkstreamBlock[]; - conversation_id?: string | null; - tenant_id?: string | null; -} - -// ── Proactive suggestion (backend: ProactiveSuggestion) ── - -export type ProactivePriority = 'low' | 'medium' | 'high' | 'urgent'; - -export interface ProactiveSuggestion { - id: string; - trigger: string; - title: string; - description: string; - priority: ProactivePriority; - action_type: 'suggestion' | 'action_required' | 'info'; - action_url: string; - entity_type?: string | null; - entity_id?: string | null; - blocks: WorkstreamBlock[]; - created_at: string; - expires_at?: string | null; - metadata: Record; -} - -export interface ProactiveFeedSettings { - enabled: boolean; - min_priority: ProactivePriority; - max_per_hour: number; - triggers_enabled: Record; -} diff --git a/frontend/src/api/platform.ts b/frontend/src/api/platform.ts deleted file mode 100644 index 5812747..0000000 --- a/frontend/src/api/platform.ts +++ /dev/null @@ -1,153 +0,0 @@ -/** - * Platform API hooks — Phase I.6/I.7 (I-ONB, I-DASH, I-COST, I-USE). - * - * Provides hooks for the setup wizard (onboarding) and the platform dashboard - * (agent status, workflow stats, search metrics, knowledge coverage, cost - * tracking, system health). Backend routes may not be wired yet, so all hooks - * degrade gracefully to an "unavailable" state instead of throwing. - */ - -import { useQuery } from '@tanstack/react-query'; -import { apiGet } from './client'; - -// ── Onboarding ────────────────────────────────────────────────────────────── - -export interface OnboardingStepStatus { - completed: boolean; - required: boolean; -} - -export interface OnboardingStatus { - steps: Record; - progress_pct: number; -} - -export interface OnboardingGuideStep { - id: string; - title: string; - description: string; - icon?: string; - action_url?: string; -} - -export interface OnboardingGuide { - steps: OnboardingGuideStep[]; -} - -export function useOnboardingStatus() { - return useQuery({ - queryKey: ['onboardingStatus'], - queryFn: () => apiGet('/onboarding/status'), - staleTime: 60 * 1000, - retry: false, - }); -} - -export function useOnboardingGuide() { - return useQuery({ - queryKey: ['onboardingGuide'], - queryFn: () => apiGet('/onboarding/guide'), - staleTime: 60 * 1000, - retry: false, - }); -} - -// ── Platform dashboard ────────────────────────────────────────────────────── - -export interface PlatformDashboardData { - agents?: { - active_agents?: number; - total_runs?: number; - recent_runs_7d?: number; - error?: string; - }; - workflows?: { - active_workflows?: number; - running_instances?: number; - completed_instances?: number; - error?: string; - }; - search?: { - total_queries?: number; - avg_latency_ms?: number; - error?: string; - }; - knowledge?: { - wiki_articles?: number; - coverage_pct?: number; - error?: string; - }; - workstream?: { - messages?: number; - error?: string; - }; - system_health?: { - redis?: string; - status?: string; - error?: string; - }; - generated_at?: string; -} - -export function usePlatformDashboard() { - return useQuery({ - queryKey: ['platformDashboard'], - queryFn: () => apiGet('/dashboard/platform'), - staleTime: 60 * 1000, - retry: false, - }); -} - -// ── Cost tracking ─────────────────────────────────────────────────────────── - -export interface CostDashboardData { - period_days?: number; - total_cost_usd?: number; - by_agent?: Record; - budget?: { - monthly_limit_usd?: number; - utilization_pct?: number; - }; - error?: string; -} - -export function useCostDashboard(days = 30) { - return useQuery({ - queryKey: ['costDashboard', days], - queryFn: () => apiGet(`/dashboard/cost?days=${days}`), - staleTime: 60 * 1000, - retry: false, - }); -} - -// ── Usage analytics ───────────────────────────────────────────────────────── - -export interface UsageAnalyticsData { - period_days?: number; - agent_runs?: { - total?: number; - completed?: number; - failed?: number; - success_rate?: number; - error?: string; - }; - workflow_executions?: { - total?: number; - completed?: number; - error?: string; - }; - search_queries?: { - total?: number; - error?: string; - }; - error?: string; -} - -export function useUsageAnalytics(days = 30) { - return useQuery({ - queryKey: ['usageAnalytics', days], - queryFn: () => apiGet(`/dashboard/usage?days=${days}`), - staleTime: 60 * 1000, - retry: false, - }); -} diff --git a/frontend/src/components/improvement/ImprovementCenter.tsx b/frontend/src/components/improvement/ImprovementCenter.tsx deleted file mode 100644 index 84d97ad..0000000 --- a/frontend/src/components/improvement/ImprovementCenter.tsx +++ /dev/null @@ -1,128 +0,0 @@ -import React from 'react'; -import { useTranslation } from 'react-i18next'; -import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'; -import { TrendingUp, Lightbulb } from 'lucide-react'; -import { Card } from '@/components/ui/Card'; -import { EmptyState } from '@/components/ui/EmptyState'; -import { Skeleton } from '@/components/ui/Skeleton'; -import { ProposalCard } from './ProposalCard'; -import { PatternInsight } from './PatternInsight'; -import { - fetchProposals, - fetchPatterns, - approveProposal, - rejectProposal, - rollbackProposal, - type ImprovementProposal, - type DetectedPattern, -} from '@/api/improvement'; - -export interface ImprovementCenterProps { - className?: string; -} - -export function ImprovementCenter({ className }: ImprovementCenterProps) { - const { t } = useTranslation(); - const queryClient = useQueryClient(); - - const { data: proposals = [], isLoading: loadingProposals } = useQuery({ - queryKey: ['improvement', 'proposals'], - queryFn: () => fetchProposals(), - }); - - const { data: patterns = [], isLoading: loadingPatterns } = useQuery({ - queryKey: ['improvement', 'patterns'], - queryFn: fetchPatterns, - }); - - const invalidate = () => { - queryClient.invalidateQueries({ queryKey: ['improvement', 'proposals'] }); - queryClient.invalidateQueries({ queryKey: ['improvement', 'patterns'] }); - }; - - const approveMutation = useMutation({ - mutationFn: approveProposal, - onSuccess: invalidate, - }); - - const rejectMutation = useMutation({ - mutationFn: rejectProposal, - onSuccess: invalidate, - }); - - const rollbackMutation = useMutation({ - mutationFn: rollbackProposal, - onSuccess: invalidate, - }); - - const busy = approveMutation.isPending || rejectMutation.isPending || rollbackMutation.isPending; - - return ( -
-
-

-

-

{t('improvement.subtitle')}

-
- - {/* Patterns / Bottlenecks */} -
-

- {t('improvement.patterns')} -

- {loadingPatterns ? ( -
- - -
- ) : patterns.length === 0 ? ( - } - title={t('improvement.noPatterns')} - description={t('improvement.noPatternsDesc')} - /> - ) : ( -
- {patterns.map((pattern) => ( - - ))} -
- )} -
- - {/* Proposals */} -
-

- {t('improvement.proposals')} -

- {loadingProposals ? ( -
- - -
- ) : proposals.length === 0 ? ( - } - title={t('improvement.noProposals')} - description={t('improvement.noProposalsDesc')} - /> - ) : ( -
- {proposals.map((proposal) => ( - approveMutation.mutate(id)} - onReject={(id) => rejectMutation.mutate(id)} - onRollback={(id) => rollbackMutation.mutate(id)} - /> - ))} -
- )} -
-
- ); -} diff --git a/frontend/src/components/improvement/PatternInsight.tsx b/frontend/src/components/improvement/PatternInsight.tsx deleted file mode 100644 index 8196871..0000000 --- a/frontend/src/components/improvement/PatternInsight.tsx +++ /dev/null @@ -1,70 +0,0 @@ -import React from 'react'; -import clsx from 'clsx'; -import { useTranslation } from 'react-i18next'; -import { TrendingUp, AlertTriangle, Lightbulb } from 'lucide-react'; -import { Card } from '@/components/ui/Card'; -import { Badge } from '@/components/ui/Badge'; -import type { DetectedPattern } from '@/api/improvement'; - -export interface PatternInsightProps { - pattern: DetectedPattern; -} - -function formatConfidence(confidence: number): string { - if (confidence === undefined || Number.isNaN(confidence)) return '—'; - return `${Math.round(confidence * 100)}%`; -} - -export function PatternInsight({ pattern }: PatternInsightProps) { - const { t } = useTranslation(); - const confidence = pattern.confidence; - const confidenceVariant = - confidence >= 0.8 ? 'success' : confidence >= 0.6 ? 'warning' : 'secondary'; - - return ( - -
-

-

- - {t('improvement.confidence')}: {formatConfidence(confidence)} - -
- -

{pattern.description}

- -
-
-
- {pattern.pattern_type === 'error_retries' && ( -
-
- )} -
- - {pattern.evidence_refs.length > 0 && ( -
-

- {t('improvement.evidence')} -

-
    - {pattern.evidence_refs.map((ref) => ( -
  • - {ref} -
  • - ))} -
-
- )} -
- ); -} diff --git a/frontend/src/components/improvement/ProposalCard.tsx b/frontend/src/components/improvement/ProposalCard.tsx deleted file mode 100644 index 7cc5e3a..0000000 --- a/frontend/src/components/improvement/ProposalCard.tsx +++ /dev/null @@ -1,131 +0,0 @@ -import React from 'react'; -import clsx from 'clsx'; -import { useTranslation } from 'react-i18next'; -import { TrendingUp, AlertTriangle, CheckCircle, XCircle, RotateCcw, Lightbulb } from 'lucide-react'; -import { Card } from '@/components/ui/Card'; -import { Badge } from '@/components/ui/Badge'; -import { Button } from '@/components/ui/Button'; -import type { ImprovementProposal, ProposalStatus } from '@/api/improvement'; - -export interface ProposalCardProps { - proposal: ImprovementProposal; - onApprove?: (id: string) => void; - onReject?: (id: string) => void; - onRollback?: (id: string) => void; - busy?: boolean; -} - -const statusVariant: Record = { - draft: 'secondary', - evaluating: 'info', - pending_approval: 'warning', - approved: 'primary', - rejected: 'danger', - active: 'success', - rolled_back: 'default', - expired: 'default', -}; - -function formatScore(score: number | undefined): string { - if (score === undefined || Number.isNaN(score)) return '—'; - return `${Math.round(score)}%`; -} - -export function ProposalCard({ proposal, onApprove, onReject, onRollback, busy = false }: ProposalCardProps) { - const { t } = useTranslation(); - const score = proposal.evaluation_result?.score; - const canApprove = proposal.status === 'pending_approval' || proposal.status === 'draft' || proposal.status === 'evaluating'; - const canReject = proposal.status === 'pending_approval' || proposal.status === 'draft' || proposal.status === 'evaluating'; - const canRollback = proposal.status === 'active'; - - return ( - -
-

-

- - {t(`improvement.status.${proposal.status}`)} - -
- -
-

{proposal.description}

- -
-
-

- {t('improvement.rationale')} -

-

{proposal.rationale}

-
-
-

- {t('improvement.expectedBenefit')} -

-

{proposal.expected_benefit}

-
-
- -
-
- - {proposal.evidence_refs.length > 0 && ( -
-

- {t('improvement.evidence')} -

-
    - {proposal.evidence_refs.map((ref) => ( -
  • - {ref} -
  • - ))} -
-
- )} - -
-
-
- = 80 ? 'text-success-600' : score !== undefined && score >= 60 ? 'text-warning-600' : 'text-secondary-600' - )} - > - {formatScore(score)} - -
-
- -
- {canApprove && onApprove && ( - - )} - {canReject && onReject && ( - - )} - {canRollback && onRollback && ( - - )} -
-
- ); -} diff --git a/frontend/src/components/layout/Sidebar.tsx b/frontend/src/components/layout/Sidebar.tsx index 862b684..58f4152 100644 --- a/frontend/src/components/layout/Sidebar.tsx +++ b/frontend/src/components/layout/Sidebar.tsx @@ -54,10 +54,7 @@ function getIcon(name: string): React.ReactNode { const singleItems: NavSingleItem[] = [ { to: '/dashboard', labelKey: 'nav.dashboard', icon: