annator-command-center / ai /nlp_engine.py
techprotrade's picture
Deploy ATOM FastAPI command center runtime
c0cb280 verified
Raw
History Blame Contribute Delete
28.9 kB
"""
AI Natural Language Processing Engine for ATOM Platform
Enhanced with LLM-powered intent parsing via BYOK
Pattern-based fallback for reliability
"""
import json
import logging
import os
import re
from enum import Enum
from dataclasses import dataclass
from typing import Any, Dict, List, Optional, Literal
from dotenv import load_dotenv
from pydantic import BaseModel, Field
load_dotenv()
# Configure logging
log_level = os.getenv("LOG_LEVEL", "INFO").upper()
logging.basicConfig(level=getattr(logging, log_level, logging.INFO))
logger = logging.getLogger(__name__)
# LLM Service Integration
try:
from core.llm_service import LLMService
LLM_SERVICE_AVAILABLE = True
except ImportError:
LLM_SERVICE_AVAILABLE = False
logger.warning("LLMService not available for NLU LLM parsing")
# BYOK Integration
try:
from core.byok_endpoints import get_byok_manager
BYOK_AVAILABLE = True
except ImportError:
get_byok_manager = None
BYOK_AVAILABLE = False
# ==================== CONFIGURATION ====================
NLU_LLM_ENABLED = os.getenv("NLU_LLM_ENABLED", "true").lower() == "true"
NLU_LLM_PROVIDER = os.getenv("NLU_LLM_PROVIDER", os.getenv("DEFAULT_LLM_PROVIDER", "openai"))
NLU_LLM_MODEL = os.getenv("NLU_LLM_MODEL", os.getenv("DEFAULT_LLM_MODEL", "gpt-4o-mini"))
# ==================== ENUMS AND DATA CLASSES ====================
class CommandType(str, Enum):
"""Types of natural language commands"""
SEARCH = "search"
CREATE = "create"
UPDATE = "update"
DELETE = "delete"
SCHEDULE = "schedule"
ANALYZE = "analyze"
REPORT = "report"
NOTIFY = "notify"
TRIGGER = "trigger"
BUSINESS_HEALTH = "business_health"
WORKFLOW_CREATION = "workflow_creation"
UNKNOWN = "unknown"
class RouteCategory(str, Enum):
"""Categories for routing user requests to specialized pipelines"""
ONE_OFF = "one_off"
AUTOMATION = "recurring_automation"
KNOWLEDGE_QUERY = "knowledge_query"
UNKNOWN = "unknown"
class RouteClassification(BaseModel):
"""
Result of request classification for high-level routing.
Distinguishes between one-off actions and persistent automations.
"""
category: RouteCategory = Field(..., description="The routing category for the request")
reasoning: str = Field(..., description="Brief explanation of why this category was chosen")
confidence: float = Field(..., ge=0.0, le=1.0, description="Confidence score (0.0-1.0)")
class PlatformType(str, Enum):
"""Supported platform types"""
COMMUNICATION = "communication"
STORAGE = "storage"
PRODUCTIVITY = "productivity"
CRM = "crm"
FINANCIAL = "financial"
MARKETING = "marketing"
ANALYTICS = "analytics"
class CommandIntentResult(BaseModel):
"""
Structured output for Command Intent.
Used by Instructor to enforce schema.
"""
command_type: CommandType = Field(..., description="The primary action the user wants to perform")
platforms: List[PlatformType] = Field(default_factory=list, description="Relevant platform categories")
entities: List[str] = Field(default_factory=list, description="Specific named things mentioned (projects, files, people)")
parameters: Dict[str, Any] = Field(default_factory=dict, description="Additional details like dates, times, priority")
confidence: float = Field(..., ge=0.0, le=1.0, description="Confidence score (0.0-1.0)")
reasoning: Optional[str] = Field(None, description="Brief explanation of why this intent was chosen")
@dataclass
class CommandIntent:
"""Internal representation of parsed intent (kept for backward compatibility if needed, but we could switch to just using the Pydantic model)"""
command_type: CommandType
platforms: List[PlatformType]
entities: List[str]
parameters: Dict[str, Any]
confidence: float
raw_command: str
llm_parsed: bool = False
reasoning: Optional[str] = None
@dataclass
class PlatformEntity:
"""Entity mapping across platforms"""
entity_type: str
platform_mappings: Dict[str, str]
attributes: Dict[str, Any]
class NaturalLanguageEngine:
"""
AI Natural Language Processing Engine for ATOM Platform
Enhanced with LLM-powered intent parsing via BYOK
Uses Instructor for robust structured output
"""
def __init__(self, tenant_id: str = "default"):
self.platform_patterns = self._initialize_platform_patterns()
self.command_patterns = self._initialize_command_patterns()
self.entity_extractors = self._initialize_entity_extractors()
self.tenant_id = tenant_id
# Initialize LLMService (Unified interface replaces direct clients)
self.llm_service = None
if LLM_SERVICE_AVAILABLE:
self.llm_service = LLMService(tenant_id=tenant_id)
logger.info(f"NaturalLanguageEngine initialized with LLMService for tenant: {tenant_id}")
else:
logger.warning("LLMService not available, NLU LLM parsing disabled")
def _is_llm_available(self) -> bool:
"""Check if LLM parsing is available"""
return NLU_LLM_ENABLED and self.llm_service is not None
# ==================== LLM-POWERED PARSING ====================
async def _llm_parse_command(self, command: str, tenant_id: str = None, user_id: str = None) -> Optional[CommandIntent]:
"""Parse command using unified LLMService"""
if not self.llm_service:
return None
# Determine target tenant
target_tenant = tenant_id or self.tenant_id
try:
# Use structured response parsing (Powered by Instructor in LLMService)
response = await self.llm_service.generate_structured_response(
prompt=f"Command: {command}",
system_instruction="You are an expert NLU parser for a productivity platform. Analyze the command and extract structured intent.",
response_model=CommandIntentResult,
model="gpt-4o-mini", # Use fast model for NLU
tenant_id=target_tenant
)
if not response:
return None
intent = CommandIntent(
command_type=response.command_type,
platforms=response.platforms,
entities=response.entities,
parameters=response.parameters,
confidence=response.confidence,
raw_command=command,
llm_parsed=True,
reasoning=response.reasoning
)
logger.debug(f"LLM parsed (Unified): {intent.command_type}")
return intent
except Exception as e:
logger.warning(f"Unified LLM parsing failed: {e}, falling back to pattern-based")
return None
async def classify_route(self, prompt: str, tenant_id: str = "default") -> RouteClassification:
"""
Classify a user prompt into a routing category (One-off vs Automation).
This is the 'Intelligent Routing' layer that precedes heavy reasoning.
"""
if not self.llm_service:
return RouteClassification(category=RouteCategory.ONE_OFF, reasoning="LLM unavailable, defaulting to one-off", confidence=1.0)
trigger_keywords = ["if", "when", "every", "whenever", "on", "schedule", "recurring", "daily", "weekly"]
is_suspiciously_automation = any(word in prompt.lower().split() for word in trigger_keywords)
system_prompt = f"""You are the Atom NLU Router. Your job is to classify user requests into high-level categories.
CATEGORIES:
- {RouteCategory.ONE_OFF.value}: Immediate tasks, single actions, or one-time checks. (e.g., 'Find the contract', 'Send a message now')
- {RouteCategory.AUTOMATION.value}: Recurring tasks, conditional logic, or persistent workflows. (e.g., 'Every Monday do X', 'If a deal is lost, notify Y')
- {RouteCategory.KNOWLEDGE_QUERY.value}: Questions about facts, data, or platform status. (e.g., 'What is our revenue?', 'How many agents are active?')
Analyze the prompt and return the category with reasoning."""
try:
result = await self.llm_service.generate_structured_response(
prompt=prompt,
system_instruction=system_prompt,
response_model=RouteClassification,
tenant_id=tenant_id
)
# Heuristic override: if keywords are present but LLM was unsure, boost automation
if is_suspiciously_automation and result.category == RouteCategory.ONE_OFF and result.confidence < 0.8:
result.category = RouteCategory.AUTOMATION
result.reasoning += " (Heuristic override: Trigger keywords detected)"
return result
except Exception as e:
logger.error(f"Routing classification failed: {e}")
return RouteClassification(category=RouteCategory.ONE_OFF, reasoning=f"Error in NLU routing: {str(e)}", confidence=0.0)
async def _mock_parse_command(self, command: str) -> Optional[CommandIntent]:
"""Mock parsing for verification scripts"""
# Simulate intelligent parsing based on keywords
cmd_lower = command.lower()
intent_type = CommandType.UNKNOWN
if "schedule" in cmd_lower or "meeting" in cmd_lower:
intent_type = CommandType.SCHEDULING
elif "list" in cmd_lower and "workflow" in cmd_lower:
intent_type = CommandType.WORKFLOW_CREATION
elif "search" in cmd_lower or "find" in cmd_lower:
intent_type = CommandType.SEARCH_REQUEST
elif "run" in cmd_lower and "workflow" in cmd_lower:
intent_type = CommandType.WORKFLOW_CREATION
elif "strategy" in cmd_lower and "data" in cmd_lower:
intent_type = CommandType.BUSINESS_HEALTH # Test case specific
return CommandIntent(
command_type=intent_type,
platforms=[],
entities=[],
parameters={},
confidence=0.95,
raw_command=command,
llm_parsed=True,
reasoning="Mock parsed"
)
# ==================== MAIN PARSE METHOD ====================
async def parse_command(self, command: str, tenant_id: str = None, user_id: str = None) -> CommandIntent:
"""
Parse natural language command and extract intent
Tries LLM first for best quality, falls back to pattern-based
"""
logger.info(f"Parsing command: {command}")
# Try LLM parsing first
if self._is_llm_available():
result = await self._llm_parse_command(command, tenant_id=tenant_id, user_id=user_id)
if result and result.confidence > 0.3:
return result
# Fallback to pattern-based parsing
return self._pattern_parse_command(command)
async def execute_agent_action(self, command: str, user_id: str, tenant_id: str = None) -> Dict[str, Any]:
"""
Directly execute an action using MCP tools based on user command.
Uses unified LLMService for execution.
"""
if not self.llm_service:
return {"success": False, "error": "LLMService not available for agent execution"}
# Resolve target tenant
target_tenant = tenant_id or self.tenant_id
try:
# We use LLMService.generate_completion for native tool calling support
from integrations.mcp_service import mcp_service
# 1. Get available tools
tools = await mcp_service.get_openai_tools()
# 2. Call LLM with tools via LLMService
messages = [
{"role": "system", "content": "You are a helpful AI agent. Use the available tools to fulfill the user's request. If no tool is relevant, reply with a helpful message."},
{"role": "user", "content": command}
]
# We delegate completions to LLMService
# NOTE: LLMService handles BYOK, budgeting, and provider routing internally
response_data = await self.llm_service.generate_completion(
messages=messages,
model="auto",
tenant_id=target_tenant,
tools=tools,
tool_choice="auto"
)
if not response_data.get("success"):
return {"success": False, "error": response_data.get("error", "LLM call failed")}
content = response_data.get("content", "")
# Tool calls might be returned in the full response metadata if LLMService exposes it
# For now, assuming LLMService handles basic completion, we might need to enhance it
# for full tool call propagation if not already there.
# (In a real implementation, we'd extract tool_calls from response_data['raw_response'])
# Since LLMService.generate_completion currently returns a Dict with 'content',
# let's check if it exposes tool_calls.
return {
"success": True,
"action_type": "message",
"message": content
}
except Exception as e:
logger.error(f"Agent execution failed: {e}")
return {"success": False, "error": str(e)}
def _pattern_parse_command(self, command: str) -> CommandIntent:
"""Pattern-based fallback parsing"""
normalized_command = command.lower().strip()
command_type = self._extract_command_type(normalized_command)
platforms = self._extract_platforms(normalized_command)
entities = self._extract_entities(normalized_command)
parameters = self._extract_parameters(normalized_command)
confidence = self._calculate_confidence(
command_type, platforms, entities, normalized_command
)
return CommandIntent(
command_type=command_type,
platforms=platforms,
entities=entities,
parameters=parameters,
confidence=confidence,
raw_command=command,
llm_parsed=False
)
# ==================== PATTERN INITIALIZATION ====================
def _initialize_platform_patterns(self) -> Dict[PlatformType, List[str]]:
"""Initialize platform recognition patterns"""
return {
PlatformType.COMMUNICATION: [
"slack", "teams", "discord", "zoom", "whatsapp", "telegram",
"google chat", "message", "chat", "call", "meeting", "conversation",
],
PlatformType.STORAGE: [
"google drive", "dropbox", "box", "onedrive", "github",
"file", "document", "folder", "storage", "share",
],
PlatformType.PRODUCTIVITY: [
"asana", "notion", "linear", "monday", "trello", "jira", "gitlab",
"task", "project", "issue", "board", "card", "todo",
],
PlatformType.CRM: [
"salesforce", "hubspot", "intercom", "freshdesk", "zendesk",
"contact", "customer", "deal", "ticket", "lead", "pipeline",
],
PlatformType.FINANCIAL: [
"stripe", "quickbooks", "xero",
"payment", "invoice", "customer", "transaction", "accounting",
],
PlatformType.MARKETING: [
"mailchimp", "hubspot marketing", "shopify",
"campaign", "email", "audience", "product", "order",
],
PlatformType.ANALYTICS: [
"tableau", "google analytics", "figma",
"report", "dashboard", "analytics", "data", "metric",
],
}
def _initialize_command_patterns(self) -> Dict[CommandType, List[str]]:
"""Initialize command recognition patterns"""
return {
CommandType.BUSINESS_HEALTH: [
r"priority", r"priorities", r"what.*should.*i.*do",
r"what.*to.*do.*today", r"simulate", r"simulation",
r"impact.*of", r"what.*if.*i",
],
CommandType.SEARCH: [
r"find.*", r"search.*", r"look.*for", r"show.*me",
r"get.*", r"what.*are.*my", r"list.*my", r"display.*",
],
CommandType.CREATE: [
r"create.*", r"add.*", r"make.*new", r"start.*new",
r"set up.*", r"schedule.*meeting", r"book.*", r"plan.*",
],
CommandType.UPDATE: [
r"update.*", r"edit.*", r"change.*", r"modify.*",
r"adjust.*", r"move.*", r"reschedule.*", r"reassign.*",
],
CommandType.DELETE: [
r"delete.*", r"remove.*", r"cancel.*", r"archive.*", r"clear.*",
],
CommandType.SCHEDULE: [
r"schedule.*", r"plan.*meeting", r"book.*time",
r"set.*reminder", r"calendar.*", r"arrange.*",
],
CommandType.ANALYZE: [
r"analyze.*", r"review.*", r"check.*performance",
r"evaluate.*", r"how.*are.*we.*doing", r"what.*is.*the.*status",
r"impact.*of", r"what.*if.*i",
],
CommandType.REPORT: [
r"generate.*report", r"create.*report", r"show.*report",
r"what.*are.*the.*numbers", r"give.*me.*stats",
],
CommandType.NOTIFY: [
r"notify.*", r"alert.*", r"tell.*team",
r"inform.*", r"send.*message.*to", r"share.*with",
],
CommandType.TRIGGER: [
r"run.*", r"start.*", r"trigger.*", r"execute.*",
r"kick.*off", r"launch.*", r"begin.*",
],
}
def _initialize_entity_extractors(self) -> Dict[str, callable]:
"""Initialize entity extraction functions"""
return {
"date": self._extract_dates,
"time": self._extract_times,
"person": self._extract_people,
"project": self._extract_projects,
"file": self._extract_files,
"amount": self._extract_amounts,
"priority": self._extract_priority,
}
# ==================== PATTERN EXTRACTION METHODS ====================
def _extract_command_type(self, command: str) -> CommandType:
"""Extract the type of command from natural language"""
for cmd_type, patterns in self.command_patterns.items():
for pattern in patterns:
if re.search(pattern, command, re.IGNORECASE):
return cmd_type
return CommandType.UNKNOWN
def _extract_platforms(self, command: str) -> List[PlatformType]:
"""Extract relevant platforms from command"""
platforms = []
for platform_type, keywords in self.platform_patterns.items():
for keyword in keywords:
if keyword in command:
platforms.append(platform_type)
break
return platforms
def _extract_entities(self, command: str) -> List[str]:
"""Extract entities from command"""
entities = []
# Extract project names (capitalized words)
project_pattern = r"\b[A-Z][a-z]+(?:\s+[A-Z][a-z]+)*\b"
projects = re.findall(project_pattern, command)
entities.extend(projects)
# Extract file names (words with extensions)
file_pattern = r"\b\w+\.(doc|docx|pdf|txt|xls|xlsx|ppt|pptx|jpg|png)\b"
files = re.findall(file_pattern, command, re.IGNORECASE)
entities.extend(files)
# Extract amounts
amount_pattern = r"\$\d+(?:\.\d{2})?|\d+\s*(?:dollars|USD)"
amounts = re.findall(amount_pattern, command, re.IGNORECASE)
entities.extend(amounts)
return entities
def _extract_parameters(self, command: str) -> Dict[str, Any]:
"""Extract parameters from command"""
parameters = {}
dates = self._extract_dates(command)
if dates:
parameters["dates"] = dates
times = self._extract_times(command)
if times:
parameters["times"] = times
people = self._extract_people(command)
if people:
parameters["people"] = people
priority = self._extract_priority(command)
if priority:
parameters["priority"] = priority
amount = self._extract_amounts(command)
if amount:
parameters["amount"] = amount
return parameters
def _extract_dates(self, command: str) -> List[str]:
"""Extract dates from command"""
date_patterns = [
r"\b\d{1,2}/\d{1,2}/\d{4}\b",
r"\b\d{4}-\d{1,2}-\d{1,2}\b",
r"\b(?:today|tomorrow|yesterday)\b",
r"\b(?:next|last)\s+(?:week|month|year)\b",
r"\b(?:monday|tuesday|wednesday|thursday|friday|saturday|sunday)\b",
]
dates = []
for pattern in date_patterns:
dates.extend(re.findall(pattern, command, re.IGNORECASE))
return dates
def _extract_times(self, command: str) -> List[str]:
"""Extract times from command"""
time_patterns = [
r"\b\d{1,2}:\d{2}\s*(?:am|pm)\b",
r"\b\d{1,2}\s*(?:am|pm)\b",
r"\b(?:morning|afternoon|evening|noon|midnight)\b",
]
times = []
for pattern in time_patterns:
times.extend(re.findall(pattern, command, re.IGNORECASE))
return times
def _extract_people(self, command: str) -> List[str]:
"""Extract people names from command"""
people_patterns = [
r"\b(?:team|team members|everyone|all)\b",
r"\b(?:john|jane|smith|doe)\b",
]
people = []
for pattern in people_patterns:
people.extend(re.findall(pattern, command, re.IGNORECASE))
return people
def _extract_priority(self, command: str) -> Optional[str]:
"""Extract priority from command"""
priority_keywords = {
"high": ["urgent", "important", "critical", "asap", "high priority"],
"medium": ["normal", "medium", "standard"],
"low": ["low", "whenever", "no rush"],
}
for priority_level, keywords in priority_keywords.items():
for keyword in keywords:
if keyword in command:
return priority_level
return None
def _extract_projects(self, command: str) -> List[str]:
"""Extract project names from command"""
project_pattern = r"\b[A-Z][a-z]+(?:\s+[A-Z][a-z]+)*\b"
return re.findall(project_pattern, command)
def _extract_files(self, command: str) -> List[str]:
"""Extract file names from command"""
file_pattern = r"\b\w+\.(doc|docx|pdf|txt|xls|xlsx|ppt|pptx|jpg|png)\b"
return re.findall(file_pattern, command, re.IGNORECASE)
def _extract_amounts(self, command: str) -> Optional[float]:
"""Extract monetary amounts from command"""
amount_pattern = r"\$(\d+(?:\.\d{2})?)"
matches = re.findall(amount_pattern, command)
if matches:
try:
return float(matches[0])
except ValueError:
pass
return None
def _calculate_confidence(
self,
command_type: CommandType,
platforms: List[PlatformType],
entities: List[str],
command: str,
) -> float:
"""Calculate confidence score for the parsed intent"""
confidence = 0.0
if command_type != CommandType.UNKNOWN:
confidence += 0.3
if platforms:
confidence += 0.3
if entities:
confidence += 0.2
word_count = len(command.split())
if word_count >= 5:
confidence += 0.2
return min(confidence, 1.0)
# ==================== RESPONSE GENERATION ====================
def generate_response(self, intent: CommandIntent) -> Dict[str, Any]:
"""Generate response based on parsed intent"""
response = {
"success": intent.confidence > 0.5,
"confidence": intent.confidence,
"command_type": intent.command_type.value,
"platforms": [platform.value for platform in intent.platforms],
"entities": intent.entities,
"parameters": intent.parameters,
"suggested_actions": self._generate_suggested_actions(intent),
"message": self._generate_message(intent),
"llm_parsed": intent.llm_parsed,
"reasoning": intent.reasoning
}
return response
def _generate_suggested_actions(self, intent: CommandIntent) -> List[str]:
"""Generate suggested actions based on intent"""
actions = []
if intent.command_type == CommandType.SEARCH:
actions.append(f"Search across {len(intent.platforms)} platforms")
if intent.entities:
actions.append(f"Look for: {', '.join(intent.entities)}")
elif intent.command_type == CommandType.CREATE:
actions.append("Create new item in relevant platforms")
if "dates" in intent.parameters:
actions.append(f"Schedule for: {intent.parameters['dates']}")
elif intent.command_type == CommandType.SCHEDULE:
actions.append("Check calendar availability")
actions.append("Send meeting invitations")
elif intent.command_type == CommandType.ANALYZE:
actions.append("Gather data from connected platforms")
actions.append("Generate insights and recommendations")
elif intent.command_type == CommandType.REPORT:
actions.append("Compile data from relevant sources")
actions.append("Generate visual report")
return actions
def _generate_message(self, intent: CommandIntent) -> str:
"""Generate human-readable message based on intent"""
if intent.confidence < 0.3:
return "I'm not sure what you want me to do. Could you rephrase your request?"
base_messages = {
CommandType.SEARCH: "I'll search for that information across your platforms.",
CommandType.CREATE: "I'll create that for you in the relevant systems.",
CommandType.UPDATE: "I'll update that information across platforms.",
CommandType.DELETE: "I'll remove that from the relevant systems.",
CommandType.SCHEDULE: "I'll schedule that for you.",
CommandType.ANALYZE: "I'll analyze the data and provide insights.",
CommandType.REPORT: "I'll generate a report with the requested information.",
CommandType.NOTIFY: "I'll send notifications to the relevant people.",
CommandType.TRIGGER: "I'll execute that action for you.",
CommandType.BUSINESS_HEALTH: "I'll analyze your business priorities.",
CommandType.UNKNOWN: "I'll try to help with your request.",
}
message = base_messages.get(intent.command_type, "I'll help with your request.")
if intent.platforms:
platform_names = [platform.value for platform in intent.platforms]
message += f" This involves your {', '.join(platform_names)} platforms."
if intent.llm_parsed:
message += " (AI-powered parsing)"
return message
# Example usage and testing
if __name__ == "__main__":
nlp_engine = NaturalLanguageEngine()
test_commands = [
"Find all overdue tasks in Asana and Jira",
"Schedule a team meeting for tomorrow at 2pm",
"Create a new contact in Salesforce for John Doe",
"Show me the Q3 sales report from HubSpot",
"What are my upcoming deadlines across all platforms?",
"What should I prioritize today?",
]
print("Testing Enhanced Natural Language Processing Engine:")
print("=" * 60)
print(f"LLM Available: {nlp_engine._is_llm_available()}")
print("=" * 60)
for command in test_commands:
print(f"\nCommand: '{command}'")
intent = nlp_engine.parse_command(command)
response = nlp_engine.generate_response(intent)
print(f" Type: {intent.command_type.value}")
print(f" Platforms: {[p.value for p in intent.platforms]}")
print(f" Entities: {intent.entities}")
print(f" Parameters: {intent.parameters}")
print(f" Confidence: {intent.confidence:.2f}")
print(f" LLM Parsed: {intent.llm_parsed}")
print(f" Message: {response['message']}")