annator-command-center / integrations /atom_workflow_automation_service.py
techprotrade's picture
Deploy ATOM FastAPI command center runtime (part 5)
90c6b42 verified
Raw
History Blame Contribute Delete
108 kB
"""
ATOM Enterprise Workflow Automation Service
Comprehensive workflow automation integrating all enterprise services with intelligent automation
"""
import asyncio
from collections import Counter, defaultdict
from dataclasses import asdict, dataclass
from datetime import datetime, timedelta, timezone
from enum import Enum
import hashlib
import json
import logging
import os
import time
from typing import Any, Callable, Dict, List, Optional, Tuple, Union
import aiohttp
import httpx
import numpy as np
import pandas as pd
from core.circuit_breaker import circuit_breaker
from core.rate_limiter import rate_limiter, should_retry, calculate_backoff
from core.audit_logger import log_integration_call, log_integration_error, log_integration_attempt, log_integration_complete
from fastapi import HTTPException
# Configure logging
logger = logging.getLogger(__name__)
# Import existing ATOM services (all optional)
try:
from ai_enhanced_service import (
AIModelType,
AIRequest,
AIResponse,
AIServiceType,
AITaskType,
ai_enhanced_service,
)
except ImportError:
logger.debug("ai_enhanced_service not available")
ai_enhanced_service = None
AIModelType = None
AIRequest = None
AIResponse = None
AIServiceType = None
AITaskType = None
try:
from atom_ai_integration import atom_ai_integration
except ImportError:
logger.debug("atom_ai_integration not available")
atom_ai_integration = None
try:
from atom_discord_integration import atom_discord_integration
except ImportError:
logger.debug("atom_discord_integration not available")
atom_discord_integration = None
try:
from atom_enterprise_security_service import (
AuditEventType,
ComplianceReport,
ComplianceStandard,
SecurityAudit,
SecurityLevel,
SecurityPolicy,
ThreatDetection,
ThreatType,
atom_enterprise_security_service,
)
except ImportError:
logger.debug("atom_enterprise_security_service not available")
atom_enterprise_security_service = None
AuditEventType = None
ComplianceReport = None
ComplianceStandard = None
SecurityAudit = None
SecurityLevel = None
SecurityPolicy = None
ThreatDetection = None
ThreatType = None
try:
from atom_enterprise_unified_service import (
AutomationTriggerType,
ComplianceAutomation,
ComplianceWorkflowType,
EnterpriseServiceType,
EnterpriseWorkflow,
SecurityWorkflowAction,
WorkflowSecurityLevel,
atom_enterprise_unified_service,
)
except ImportError:
logger.debug("atom_enterprise_unified_service not available")
atom_enterprise_unified_service = None
AutomationTriggerType = None
ComplianceAutomation = None
ComplianceWorkflowType = None
EnterpriseServiceType = None
EnterpriseWorkflow = None
SecurityWorkflowAction = None
WorkflowSecurityLevel = None
try:
from atom_google_chat_integration import atom_google_chat_integration
except ImportError:
logger.debug("atom_google_chat_integration not available")
atom_google_chat_integration = None
try:
from atom_ingestion_pipeline import AtomIngestionPipeline
except ImportError:
logger.debug("AtomIngestionPipeline not available")
AtomIngestionPipeline = None
try:
from atom_memory_service import AtomMemoryService
except ImportError:
logger.debug("AtomMemoryService not available")
AtomMemoryService = None
try:
from atom_search_service import AtomSearchService
except ImportError:
logger.debug("AtomSearchService not available")
AtomSearchService = None
try:
from atom_slack_integration import atom_slack_integration
except ImportError:
logger.debug("atom_slack_integration not available")
atom_slack_integration = None
try:
from atom_teams_integration import atom_teams_integration
except ImportError:
logger.debug("atom_teams_integration not available")
atom_teams_integration = None
try:
from atom_workflow_service import (
AtomWorkflowService,
Workflow,
WorkflowAction,
WorkflowStatus,
WorkflowStep,
WorkflowTrigger,
)
except ImportError:
logger.debug("atom_workflow_service not available")
AtomWorkflowService = None
Workflow = None
WorkflowAction = None
WorkflowStatus = None
WorkflowStep = None
WorkflowTrigger = None
class WorkflowAutomationType(Enum):
"""Workflow automation types"""
SECURITY = "security"
COMPLIANCE = "compliance"
GOVERNANCE = "governance"
MONITORING = "monitoring"
AUDITING = "auditing"
INCIDENT_RESPONSE = "incident_response"
RISK_MANAGEMENT = "risk_management"
DATA_PROTECTION = "data_protection"
ACCESS_CONTROL = "access_control"
USER_MANAGEMENT = "user_management"
RESOURCE_MANAGEMENT = "resource_management"
NOTIFICATION = "notification"
REPORTING = "reporting"
INTEGRATION = "integration"
class AutomationConditionType(Enum):
"""Automation condition types"""
EVENT_TRIGGERED = "event_triggered"
SCHEDULED = "scheduled"
THRESHOLD_EXCEEDED = "threshold_exceeded"
ANOMALY_DETECTED = "anomaly_detected"
MANUAL = "manual"
WEBHOOK = "webhook"
API_CALLED = "api_called"
SYSTEM_EVENT = "system_event"
USER_ACTION = "user_action"
DATA_CHANGED = "data_changed"
SECURITY_ALERT = "security_alert"
COMPLIANCE_VIOLATION = "compliance_violation"
class AutomationActionType(Enum):
"""Automation action types"""
NOTIFICATION = "notification"
WORKFLOW_EXECUTION = "workflow_execution"
SECURITY_ENFORCEMENT = "security_enforcement"
COMPLIANCE_CHECK = "compliance_check"
DATA_PROCESSING = "data_processing"
USER_ACTION = "user_action"
SYSTEM_CONFIG = "system_config"
API_CALL = "api_call"
EMAIL_SEND = "email_send"
MESSAGE_SEND = "message_send"
FILE_OPERATION = "file_operation"
DATABASE_OPERATION = "database_operation"
LOGGING = "logging"
AUDITING = "auditing"
REPORTING = "reporting"
REMEDIATION = "remediation"
class AutomationPriority(Enum):
"""Automation priority levels"""
CRITICAL = "critical"
HIGH = "high"
MEDIUM = "medium"
LOW = "low"
INFO = "info"
class AutomationStatus(Enum):
"""Automation status"""
ACTIVE = "active"
INACTIVE = "inactive"
PAUSED = "paused"
SUSPENDED = "suspended"
ERROR = "error"
COMPLETED = "completed"
RUNNING = "running"
PENDING = "pending"
FAILED = "failed"
@dataclass
class WorkflowAutomation:
"""Workflow automation data model"""
automation_id: str
name: str
description: str
automation_type: WorkflowAutomationType
priority: AutomationPriority
status: AutomationStatus
conditions: List[Dict[str, Any]]
actions: List[Dict[str, Any]]
schedule: Optional[str]
created_at: datetime
updated_at: datetime
created_by: str
last_executed: Optional[datetime]
execution_count: int
success_count: int
failure_count: int
timeout: int
retry_policy: Dict[str, Any]
notification_rules: List[Dict[str, Any]]
metadata: Dict[str, Any]
audit_trail: List[Dict[str, Any]]
@dataclass
class AutomationExecution:
"""Automation execution data model"""
execution_id: str
automation_id: str
triggered_by: str
trigger_context: Dict[str, Any]
status: AutomationStatus
started_at: datetime
completed_at: Optional[datetime]
execution_time: float
result: Dict[str, Any]
error: Optional[str]
actions_executed: List[Dict[str, Any]]
notifications_sent: List[Dict[str, Any]]
compliance_checks: List[Dict[str, Any]]
security_checks: List[Dict[str, Any]]
metadata: Dict[str, Any]
# Auth Type: Internal
class AtomWorkflowAutomationService:
"""Enterprise workflow automation service with comprehensive integration"""
def __init__(self, tenant_id: str = "default", config: Dict[str, Any] = None):
if config is None:
config = {}
self.config = config
self.db = config.get('database')
self.cache = config.get('cache')
# Enterprise services
self.security_service = config.get('security_service') or atom_enterprise_security_service
self.unified_service = config.get('unified_service') or atom_enterprise_unified_service
self.workflow_service = config.get('workflow_service')
self.ai_service = config.get('ai_service') or ai_enhanced_service
self.ai_integration = config.get('ai_integration') or atom_ai_integration
# Platform integrations
self.platform_integrations = {
'slack': atom_slack_integration,
'teams': atom_teams_integration,
'google_chat': atom_google_chat_integration,
'discord': atom_discord_integration
}
# Automation state
self.is_initialized = False
self.automations: Dict[str, WorkflowAutomation] = {}
self.executions: Dict[str, AutomationExecution] = {}
self.scheduled_automations: Dict[str, Dict[str, Any]] = {}
self.active_triggers: Dict[str, Dict[str, Any]] = {}
self.automation_templates: Dict[str, Dict[str, Any]] = {}
# Automation metrics
self.automation_metrics = {
'total_automations': 0,
'active_automations': 0,
'executed_today': 0,
'executed_this_week': 0,
'executed_this_month': 0,
'success_rate': 0.0,
'average_execution_time': 0.0,
'automations_by_type': defaultdict(int),
'automations_by_priority': defaultdict(int),
'executions_by_status': defaultdict(int),
'error_rate': 0.0,
'time_saved_hours': 0.0,
'cost_savings': 0.0
}
# Automation scheduling
self.scheduler_running = False
self.scheduler_task = None
self.trigger_listeners = {}
# HTTP sessions for API calls
self.http_sessions = {}
logger.info("Workflow Automation Service initialized")
async def initialize(self) -> bool:
"""Initialize workflow automation service"""
try:
if not all([self.security_service, self.unified_service, self.ai_service]):
logger.error("Required services not available for workflow automation service")
return False
# Initialize automation templates
await self._initialize_automation_templates()
# Load existing automations
await self._load_automations()
# Initialize automation scheduling
await self._initialize_automation_scheduling()
# Initialize trigger listeners
await self._initialize_trigger_listeners()
# Initialize integration endpoints
await self._initialize_integration_endpoints()
# Start automation monitoring
await self._start_automation_monitoring()
self.is_initialized = True
logger.info("Workflow Automation Service initialized successfully")
return True
except Exception as e:
logger.error(f"Error initializing workflow automation service: {e}")
return False
async def create_automation(self, automation_data: Dict[str, Any], user_id: str) -> Dict[str, Any]:
"""Create workflow automation"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "initialize", locals())
try:
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
automation_id = f"auto_{int(time.time())}_{hashlib.md5(automation_data['name'].encode()).hexdigest()[:8]}"
# Validate automation data
validation_result = await self._validate_automation_data(automation_data)
if not validation_result['valid']:
return {
'ok': False,
'error': f"Automation validation failed: {validation_result['errors']}"
}
# Create automation
automation = WorkflowAutomation(
automation_id=automation_id,
name=automation_data['name'],
description=automation_data['description'],
automation_type=WorkflowAutomationType(automation_data['automation_type']),
priority=AutomationPriority(automation_data['priority']),
status=AutomationStatus.ACTIVE,
conditions=automation_data['conditions'],
actions=automation_data['actions'],
schedule=automation_data.get('schedule'),
created_at=datetime.utcnow(),
updated_at=datetime.utcnow(),
created_by=user_id,
last_executed=None,
execution_count=0,
success_count=0,
failure_count=0,
timeout=automation_data.get('timeout', 3600),
retry_policy=automation_data.get('retry_policy', {
'max_retries': 3,
'backoff': 'exponential',
'max_delay': 3600
}),
notification_rules=automation_data.get('notification_rules', []),
metadata=automation_data.get('metadata', {}),
audit_trail=[]
)
# Setup automation triggers
await self._setup_automation_triggers(automation)
# Store automation
self.automations[automation_id] = automation
# Store in database
if self.db:
await self.db.store_workflow_automation(asdict(automation))
# Update metrics
self.automation_metrics['total_automations'] += 1
self.automation_metrics['active_automations'] += 1
self.automation_metrics['automations_by_type'][automation.automation_type.value] += 1
self.automation_metrics['automations_by_priority'][automation.priority.value] += 1
# Log creation
await self._log_automation_event(
automation_id=automation_id,
event_type='automation_created',
user_id=user_id,
details={
'automation_name': automation.name,
'automation_type': automation.automation_type.value,
'priority': automation.priority.value
}
)
return {
'ok': True,
'automation_id': automation_id,
'automation': asdict(automation),
'message': "Workflow automation created successfully"
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
logger.error(f"Error creating workflow automation: {e}")
return {'ok': False, 'error': str(e)}
async def execute_automation(self, automation_id: str, trigger_context: Dict[str, Any], triggered_by: str) -> Dict[str, Any]:
"""Execute workflow automation"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "create_automation", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
automation = self.automations.get(automation_id)
if not automation:
return {'ok': False, 'error': 'Automation not found'}
if automation.status != AutomationStatus.ACTIVE:
return {'ok': False, 'error': 'Automation is not active'}
# Create execution record
execution_id = f"exec_{int(time.time())}_{hashlib.md5(automation_id.encode()).hexdigest()[:8]}"
execution = AutomationExecution(
execution_id=execution_id,
automation_id=automation_id,
triggered_by=triggered_by,
trigger_context=trigger_context,
status=AutomationStatus.RUNNING,
started_at=datetime.utcnow(),
completed_at=None,
execution_time=0.0,
result={},
error=None,
actions_executed=[],
notifications_sent=[],
compliance_checks=[],
security_checks=[],
metadata={'trigger_context': trigger_context}
)
self.executions[execution_id] = execution
# Pre-execution security checks
security_check = await self._pre_execution_security_check(automation, trigger_context)
if not security_check['passed']:
execution.status = AutomationStatus.FAILED
execution.error = f"Security check failed: {security_check['reason']}"
execution.security_checks.append(security_check)
return {
'ok': False,
'error': execution.error,
'security_violation': security_check
}
# Pre-execution compliance checks
compliance_check = await self._pre_execution_compliance_check(automation, trigger_context)
if not compliance_check['passed']:
execution.status = AutomationStatus.FAILED
execution.error = f"Compliance check failed: {compliance_check['reason']}"
execution.compliance_checks.append(compliance_check)
return {
'ok': False,
'error': execution.error,
'compliance_violation': compliance_check
}
# ========================================================================
# NEW: Maturity-Based Trigger Interception for Agent Actions
# ========================================================================
# Pre-check all actions for agent triggers that require maturity checks
for action in automation.actions:
if action.get('type') == 'workflow_execution' or action.get('type') == 'agent_trigger':
agent_id = action.get('config', {}).get('agent_id')
if agent_id:
from core.trigger_interceptor import TriggerInterceptor, TriggerSource
interceptor = TriggerInterceptor(self.db, self.workspace_id)
trigger_context = {
"action_type": action.get('type'),
"action_config": action.get('config'),
"automation_id": automation_id,
"trigger_context": trigger_context
}
decision = await interceptor.intercept_trigger(
agent_id=agent_id,
trigger_source=TriggerSource.WORKFLOW_ENGINE,
trigger_context=trigger_context
)
# Log routing decision
logger.info(
f"Workflow automation routing decision for agent {agent_id}: "
f"{decision.routing_decision.value} (maturity: {decision.agent_maturity}, "
f"confidence: {decision.confidence_score:.2f})"
)
# Handle blocked triggers
if not decision.execute:
execution.status = AutomationStatus.FAILED
execution.error = (
f"Agent action blocked by maturity guard: {decision.reason}"
)
execution.metadata['maturity_check'] = {
'agent_id': agent_id,
'blocked': True,
'reason': decision.reason,
'routing_decision': decision.routing_decision.value
}
self.db.commit()
logger.warning(
f"Workflow automation {automation_id} action blocked: {decision.reason}"
)
return {
'ok': False,
'error': execution.error,
'maturity_check': execution.metadata['maturity_check']
}
except ValueError as e:
# Agent not found or other error
logger.warning(
f"Could not check maturity for agent {agent_id} in automation: {e}"
)
# Continue with execution for backward compatibility
# ========================================================================
# Execute automation actions
execution_results = []
for action in automation.actions:
try:
action_result = await self._execute_automation_action(action, trigger_context, execution)
execution_results.append(action_result)
execution.actions_executed.append({
'action': action,
'result': action_result,
'timestamp': datetime.utcnow().isoformat()
})
# Check if execution should stop
if action_result.get('stop_execution', False):
break
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing automation action: {e}")
execution_results.append({
'success': False,
'error': str(e),
'action': action
})
# Post-execution checks
post_security_check = await self._post_execution_security_check(automation, execution_results)
post_compliance_check = await self._post_execution_compliance_check(automation, execution_results)
# Calculate execution result
successful_actions = [r for r in execution_results if r.get('success', False)]
success_rate = len(successful_actions) / len(execution_results) if execution_results else 0
# Update execution
execution.status = AutomationStatus.COMPLETED if success_rate >= 0.8 else AutomationStatus.FAILED
execution.completed_at = datetime.utcnow()
execution.execution_time = (execution.completed_at - execution.started_at).total_seconds()
execution.result = {
'success_rate': success_rate,
'total_actions': len(execution_results),
'successful_actions': len(successful_actions),
'failed_actions': len(execution_results) - len(successful_actions)
}
execution.security_checks.append(post_security_check)
execution.compliance_checks.append(post_compliance_check)
# Update automation metrics
automation.execution_count += 1
automation.last_executed = execution.completed_at
if execution.status == AutomationStatus.COMPLETED:
automation.success_count += 1
else:
automation.failure_count += 1
# Send notifications
await self._send_automation_notifications(automation, execution)
# Update metrics
await self._update_automation_metrics(automation, execution)
# Store execution in database
if self.db:
await self.db.store_automation_execution(asdict(execution))
return {
'ok': True,
'execution_id': execution_id,
'automation_id': automation_id,
'status': execution.status.value,
'execution_time': execution.execution_time,
'result': execution.result,
'actions_executed': len(execution_results),
'successful_actions': len(successful_actions),
'message': "Automation executed successfully"
}
except Exception as e:
logger.error(f"Error executing workflow automation: {e}")
return {'ok': False, 'error': str(e)}
async def create_security_automation(self, security_event: Dict[str, Any], automation_config: Dict[str, Any]) -> Dict[str, Any]:
"""Create automation from security event"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "execute_automation", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
# Determine automation type based on security event
threat_type = security_event.get('threat_type', 'unknown')
severity = security_event.get('severity', 'medium')
automation_data = {
'name': f"Security Response: {threat_type}",
'description': f"Automated response for {threat_type} security events",
'automation_type': WorkflowAutomationType.SECURITY,
'priority': severity,
'conditions': [
{
'type': AutomationConditionType.SECURITY_ALERT.value,
'threat_type': threat_type,
'severity_level': severity,
'source_ip': security_event.get('source_ip'),
'user_id': security_event.get('user_id')
}
],
'actions': automation_config.get('actions', [
{
'type': AutomationActionType.SECURITY_ENFORCEMENT.value,
'config': {
'action': 'block_ip',
'duration': 3600,
'reason': f'Security threat detected: {threat_type}'
}
},
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['security_team', 'management'],
'message': f"Security threat {threat_type} detected with severity {severity}",
'urgency': severity
}
}
]),
'schedule': None,
'timeout': 600,
'retry_policy': {
'max_retries': 2,
'backoff': 'exponential'
},
'notification_rules': [
{
'condition': 'always',
'channels': ['security_team'],
'urgency': severity
}
],
'metadata': {
'security_event': security_event,
'threat_type': threat_type,
'severity': severity
}
}
# Create automation
result = await self.create_automation(automation_data, 'security_system')
if result.get('ok'):
# Execute automation immediately
execution_result = await self.execute_automation(
automation_id=result['automation_id'],
trigger_context={'security_event': security_event},
triggered_by='security_event'
)
result['execution_result'] = execution_result
return result
except Exception as e:
logger.error(f"Error creating security automation: {e}")
return {'ok': False, 'error': str(e)}
async def create_compliance_automation(self, compliance_violation: Dict[str, Any], automation_config: Dict[str, Any]) -> Dict[str, Any]:
"""Create automation from compliance violation"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "create_security_automation", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
# Determine automation type based on compliance violation
standard = compliance_violation.get('standard', 'unknown')
violation_type = compliance_violation.get('violation_type', 'unknown')
severity = compliance_violation.get('severity', 'medium')
automation_data = {
'name': f"Compliance Response: {standard}_{violation_type}",
'description': f"Automated response for {standard} compliance violations",
'automation_type': WorkflowAutomationType.COMPLIANCE,
'priority': severity,
'conditions': [
{
'type': AutomationConditionType.COMPLIANCE_VIOLATION.value,
'standard': standard,
'violation_type': violation_type,
'severity_level': severity,
'affected_resources': compliance_violation.get('affected_resources', [])
}
],
'actions': automation_config.get('actions', [
{
'type': AutomationActionType.COMPLIANCE_CHECK.value,
'config': {
'action': 'remediate',
'standard': standard,
'violation_type': violation_type
}
},
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['compliance_officer', 'management'],
'message': f"Compliance violation {violation_type} detected for {standard}",
'urgency': severity
}
}
]),
'schedule': None,
'timeout': 1800,
'retry_policy': {
'max_retries': 3,
'backoff': 'linear'
},
'notification_rules': [
{
'condition': 'always',
'channels': ['compliance_officer'],
'urgency': severity
}
],
'metadata': {
'compliance_violation': compliance_violation,
'standard': standard,
'violation_type': violation_type,
'severity': severity
}
}
# Create automation
result = await self.create_automation(automation_data, 'compliance_system')
if result.get('ok'):
# Execute automation immediately
execution_result = await self.execute_automation(
automation_id=result['automation_id'],
trigger_context={'compliance_violation': compliance_violation},
triggered_by='compliance_violation'
)
result['execution_result'] = execution_result
return result
except Exception as e:
logger.error(f"Error creating compliance automation: {e}")
return {'ok': False, 'error': str(e)}
async def create_integration_automation(self, platform: str, integration_config: Dict[str, Any]) -> Dict[str, Any]:
"""Create automation for platform integration"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "create_compliance_automation", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
# Validate platform
if platform not in self.platform_integrations:
return {'ok': False, 'error': f'Unsupported platform: {platform}'}
automation_data = {
'name': f"Integration: {platform}",
'description': f"Automation for {platform} platform integration",
'automation_type': WorkflowAutomationType.INTEGRATION,
'priority': AutomationPriority.MEDIUM,
'conditions': [
{
'type': AutomationConditionType.EVENT_TRIGGERED.value,
'platform': platform,
'events': integration_config.get('events', ['message_received', 'user_joined'])
}
],
'actions': integration_config.get('actions', [
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['platform_admin'],
'message': f"Integration event from {platform}",
'urgency': 'low'
}
}
]),
'schedule': None,
'timeout': 300,
'retry_policy': {
'max_retries': 2,
'backoff': 'exponential'
},
'notification_rules': [
{
'condition': 'on_error',
'channels': ['platform_admin'],
'urgency': 'medium'
}
],
'metadata': {
'platform': platform,
'integration_config': integration_config
}
}
# Create automation
result = await self.create_automation(automation_data, 'integration_system')
# Setup platform-specific trigger listeners
if result.get('ok'):
await self._setup_platform_triggers(platform, result['automation_id'], integration_config)
return result
except Exception as e:
logger.error(f"Error creating integration automation: {e}")
return {'ok': False, 'error': str(e)}
async def get_automations(self, filters: Dict[str, Any] = None) -> List[Dict[str, Any]]:
"""Get workflow automations with filters"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "create_integration_automation", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
filters = filters or {}
automations = []
for automation in self.automations.values():
# Apply filters
if filters.get('automation_type') and automation.automation_type.value != filters['automation_type']:
continue
if filters.get('priority') and automation.priority.value != filters['priority']:
continue
if filters.get('status') and automation.status.value != filters['status']:
continue
if filters.get('created_by') and automation.created_by != filters['created_by']:
continue
# Include automation details
automation_details = {
'automation_id': automation.automation_id,
'name': automation.name,
'description': automation.description,
'automation_type': automation.automation_type.value,
'priority': automation.priority.value,
'status': automation.status.value,
'conditions': automation.conditions,
'actions': automation.actions,
'schedule': automation.schedule,
'created_at': automation.created_at.isoformat(),
'updated_at': automation.updated_at.isoformat(),
'created_by': automation.created_by,
'last_executed': automation.last_executed.isoformat() if automation.last_executed else None,
'execution_count': automation.execution_count,
'success_count': automation.success_count,
'failure_count': automation.failure_count,
'success_rate': automation.success_count / automation.execution_count if automation.execution_count > 0 else 0.0,
'timeout': automation.timeout,
'retry_policy': automation.retry_policy,
'notification_rules': automation.notification_rules,
'metadata': automation.metadata
}
automations.append(automation_details)
return automations
except Exception as e:
logger.error(f"Error getting automations: {e}")
return []
async def get_automation_executions(self, automation_id: str = None, filters: Dict[str, Any] = None) -> List[Dict[str, Any]]:
"""Get automation executions"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "get_automations", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
filters = filters or {}
executions = []
for execution in self.executions.values():
# Filter by automation_id if specified
if automation_id and execution.automation_id != automation_id:
continue
# Apply additional filters
if filters.get('status') and execution.status.value != filters['status']:
continue
if filters.get('triggered_by') and execution.triggered_by != filters['triggered_by']:
continue
if filters.get('date_from') and execution.started_at.date() < filters['date_from']:
continue
if filters.get('date_to') and execution.started_at.date() > filters['date_to']:
continue
# Include execution details
execution_details = {
'execution_id': execution.execution_id,
'automation_id': execution.automation_id,
'triggered_by': execution.triggered_by,
'trigger_context': execution.trigger_context,
'status': execution.status.value,
'started_at': execution.started_at.isoformat(),
'completed_at': execution.completed_at.isoformat() if execution.completed_at else None,
'execution_time': execution.execution_time,
'result': execution.result,
'error': execution.error,
'actions_executed': len(execution.actions_executed),
'notifications_sent': len(execution.notifications_sent),
'compliance_checks': len(execution.compliance_checks),
'security_checks': len(execution.security_checks),
'metadata': execution.metadata
}
executions.append(execution_details)
# Sort by started_at descending
executions.sort(key=lambda x: x['started_at'], reverse=True)
return executions
except Exception as e:
logger.error(f"Error getting automation executions: {e}")
return []
async def get_automation_metrics(self) -> Dict[str, Any]:
"""Get automation metrics"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "get_automation_executions", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
try:
return {
'total_automations': self.automation_metrics['total_automations'],
'active_automations': self.automation_metrics['active_automations'],
'executed_today': self.automation_metrics['executed_today'],
'executed_this_week': self.automation_metrics['executed_this_week'],
'executed_this_month': self.automation_metrics['executed_this_month'],
'success_rate': self.automation_metrics['success_rate'],
'average_execution_time': self.automation_metrics['average_execution_time'],
'error_rate': self.automation_metrics['error_rate'],
'time_saved_hours': self.automation_metrics['time_saved_hours'],
'cost_savings': self.automation_metrics['cost_savings'],
'automations_by_type': dict(self.automation_metrics['automations_by_type']),
'automations_by_priority': dict(self.automation_metrics['automations_by_priority']),
'executions_by_status': dict(self.automation_metrics['executions_by_status']),
'scheduled_automations': len(self.scheduled_automations),
'active_triggers': len(self.active_triggers),
'automation_templates': len(self.automation_templates)
}
except Exception as e:
logger.error(f"Error getting automation metrics: {e}")
return {}
# Private methods
async def _validate_automation_data(self, automation_data: Dict[str, Any]) -> Dict[str, Any]:
"""Validate automation data"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "get_automation_metrics", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
validation_result = {'valid': True, 'errors': [], 'warnings': []}
required_fields = ['name', 'description', 'automation_type', 'priority', 'conditions', 'actions']
for field in required_fields:
if field not in automation_data:
validation_result['valid'] = False
validation_result['errors'].append(f"Required field missing: {field}")
# Validate conditions
if 'conditions' in automation_data:
for condition in automation_data['conditions']:
if 'type' not in condition:
validation_result['valid'] = False
validation_result['errors'].append("Condition missing type")
# Validate actions
if 'actions' in automation_data:
for action in automation_data['actions']:
if 'type' not in action:
validation_result['valid'] = False
validation_result['errors'].append("Action missing type")
return validation_result
async def _setup_automation_triggers(self, automation: WorkflowAutomation):
"""Setup automation triggers"""
for condition in automation.conditions:
condition_type = condition.get('type')
if condition_type == AutomationConditionType.SCHEDULED.value:
# Schedule automation
await self._schedule_automation(automation, condition)
elif condition_type == AutomationConditionType.EVENT_TRIGGERED.value:
# Setup event trigger
await self._setup_event_trigger(automation, condition)
elif condition_type == AutomationConditionType.THRESHOLD_EXCEEDED.value:
# Setup threshold trigger
await self._setup_threshold_trigger(automation, condition)
elif condition_type == AutomationConditionType.ANOMALY_DETECTED.value:
# Setup anomaly trigger
await self._setup_anomaly_trigger(automation, condition)
elif condition_type == AutomationConditionType.SECURITY_ALERT.value:
# Setup security trigger
await self._setup_security_trigger(automation, condition)
elif condition_type == AutomationConditionType.COMPLIANCE_VIOLATION.value:
# Setup compliance trigger
await self._setup_compliance_trigger(automation, condition)
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up automation triggers: {e}")
async def _execute_automation_action(self, action: Dict[str, Any], trigger_context: Dict[str, Any], execution: AutomationExecution) -> Dict[str, Any]:
"""Execute automation action"""
action_type = action.get('type')
action_config = action.get('config', {})
# Execute based on action type
if action_type == AutomationActionType.NOTIFICATION.value:
return await self._execute_notification_action(action_config, trigger_context)
elif action_type == AutomationActionType.WORKFLOW_EXECUTION.value:
return await self._execute_workflow_action(action_config, trigger_context)
elif action_type == AutomationActionType.SECURITY_ENFORCEMENT.value:
return await self._execute_security_enforcement_action(action_config, trigger_context)
elif action_type == AutomationActionType.COMPLIANCE_CHECK.value:
return await self._execute_compliance_check_action(action_config, trigger_context)
elif action_type == AutomationActionType.DATA_PROCESSING.value:
return await self._execute_data_processing_action(action_config, trigger_context)
elif action_type == AutomationActionType.API_CALL.value:
return await self._execute_api_call_action(action_config, trigger_context)
elif action_type == AutomationActionType.EMAIL_SEND.value:
return await self._execute_email_action(action_config, trigger_context)
elif action_type == AutomationActionType.MESSAGE_SEND.value:
return await self._execute_message_action(action_config, trigger_context)
elif action_type == AutomationActionType.LOGGING.value:
return await self._execute_logging_action(action_config, trigger_context)
elif action_type == AutomationActionType.AUDITING.value:
return await self._execute_auditing_action(action_config, trigger_context)
elif action_type == AutomationActionType.REPORTING.value:
return await self._execute_reporting_action(action_config, trigger_context)
elif action_type == AutomationActionType.REMEDIATION.value:
return await self._execute_remediation_action(action_config, trigger_context)
else:
return {
'success': False,
'error': f"Unsupported action type: {action_type}"
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing automation action: {e}")
return {
'success': False,
'error': str(e)
}
# Action execution methods
async def _execute_notification_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute notification action"""
channels = config.get('channels', [])
message = config.get('message', 'Automation triggered')
urgency = config.get('urgency', 'medium')
# Send notifications to different channels
notification_results = []
for channel in channels:
if channel == 'security_team':
# Send to security team
await self._notify_security_team(message, urgency, trigger_context)
elif channel == 'compliance_officer':
# Send to compliance officer
await self._notify_compliance_officer(message, urgency, trigger_context)
elif channel == 'management':
# Send to management
await self._notify_management(message, urgency, trigger_context)
elif channel == 'slack':
# Send to Slack
await self._notify_slack(message, urgency, trigger_context)
elif channel == 'teams':
# Send to Teams
await self._notify_teams(message, urgency, trigger_context)
elif channel == 'email':
# Send email
await self._notify_email(message, urgency, trigger_context)
notification_results.append({
'channel': channel,
'success': True,
'message': message
})
return {
'success': True,
'notification_results': notification_results,
'channels': channels,
'message': message,
'urgency': urgency
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing notification action: {e}")
return {
'success': False,
'error': str(e)
}
async def _execute_workflow_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute workflow action"""
workflow_id = config.get('workflow_id')
workflow_data = config.get('workflow_data', {})
if not workflow_id:
return {
'success': False,
'error': 'workflow_id is required for workflow action'
}
# Execute workflow using unified service
if self.unified_service:
result = await self.unified_service.execute_enterprise_workflow(
workflow_id=workflow_id,
trigger_context={
'automation_trigger': trigger_context,
'workflow_data': workflow_data
},
user_id='automation_system'
)
return {
'success': result.get('ok', False),
'result': result,
'workflow_id': workflow_id
}
else:
return {
'success': False,
'error': 'Unified service not available'
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing workflow action: {e}")
return {
'success': False,
'error': str(e)
}
async def _execute_security_enforcement_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute security enforcement action"""
enforcement_action = config.get('action')
target = config.get('target')
reason = config.get('reason', 'Security policy violation')
if not enforcement_action:
return {
'success': False,
'error': 'action is required for security enforcement'
}
# Execute using security service
if self.security_service:
if enforcement_action == 'block_ip':
await self.security_service._block_ip(target, config.get('duration', 3600))
elif enforcement_action == 'lock_user':
await self.security_service._lock_user_account(target)
elif enforcement_action == 'terminate_session':
await self.security_service._terminate_session(target)
elif enforcement_action == 'quarantine':
await self.security_service._quarantine_resource(target)
return {
'success': True,
'enforcement_action': enforcement_action,
'target': target,
'reason': reason
}
else:
return {
'success': False,
'error': 'Security service not available'
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing security enforcement action: {e}")
return {
'success': False,
'error': str(e)
}
async def _execute_compliance_check_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute compliance check action"""
standard = config.get('standard')
check_type = config.get('check_type', 'automated')
if not standard:
return {
'success': False,
'error': 'standard is required for compliance check'
}
# Execute compliance check using security service
if self.security_service:
compliance_report = await self.security_service.check_compliance(
ComplianceStandard(standard),
trigger_context.get('period', 'immediate')
)
return {
'success': compliance_report is not None,
'compliance_report': asdict(compliance_report) if compliance_report else None,
'standard': standard,
'check_type': check_type
}
else:
return {
'success': False,
'error': 'Security service not available'
}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error executing compliance check action: {e}")
return {
'success': False,
'error': str(e)
}
# Additional action execution methods would be implemented here
async def _execute_data_processing_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute data processing action"""
return {'success': True, 'message': 'Data processing action executed'}
async def _execute_api_call_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute API call action"""
return {'success': True, 'message': 'API call action executed'}
async def _execute_email_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute email action"""
return {'success': True, 'message': 'Email action executed'}
async def _execute_message_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute message action"""
return {'success': True, 'message': 'Message action executed'}
async def _execute_logging_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute logging action"""
return {'success': True, 'message': 'Logging action executed'}
async def _execute_auditing_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute auditing action"""
return {'success': True, 'message': 'Auditing action executed'}
async def _execute_reporting_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute reporting action"""
return {'success': True, 'message': 'Reporting action executed'}
async def _execute_remediation_action(self, config: Dict[str, Any], trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Execute remediation action"""
return {'success': True, 'message': 'Remediation action executed'}
# Security and compliance checks
async def _pre_execution_security_check(self, automation: WorkflowAutomation, trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Pre-execution security check"""
# Check security level
if automation.automation_type == WorkflowAutomationType.SECURITY:
security_level = WorkflowSecurityLevel.RESTRICTED
else:
security_level = WorkflowSecurityLevel.INTERNAL
# Validate trigger context
if not trigger_context.get('authorized', True):
return {
'passed': False,
'reason': 'Trigger context not authorized'
}
return {'passed': True}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in pre-execution security check: {e}")
return {
'passed': False,
'reason': str(e)
}
async def _pre_execution_compliance_check(self, automation: WorkflowAutomation, trigger_context: Dict[str, Any]) -> Dict[str, Any]:
"""Pre-execution compliance check"""
# Check compliance requirements
if automation.automation_type == WorkflowAutomationType.COMPLIANCE:
return {
'passed': True,
'compliance_level': 'compliant'
}
return {'passed': True}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in pre-execution compliance check: {e}")
return {
'passed': False,
'reason': str(e)
}
async def _post_execution_security_check(self, automation: WorkflowAutomation, execution_results: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Post-execution security check"""
# Validate execution results
for result in execution_results:
if not result.get('success', False):
logger.warning(f"Security action failed: {result}")
return {'passed': True}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in post-execution security check: {e}")
return {
'passed': False,
'reason': str(e)
}
async def _post_execution_compliance_check(self, automation: WorkflowAutomation, execution_results: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Post-execution compliance check"""
# Validate compliance
return {'passed': True}
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in post-execution compliance check: {e}")
return {
'passed': False,
'reason': str(e)
}
# Notification methods
async def _notify_security_team(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify security team"""
# Mock implementation
logger.info(f"Security Team Notification: {message} (Urgency: {urgency})")
async def _notify_compliance_officer(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify compliance officer"""
# Mock implementation
logger.info(f"Compliance Officer Notification: {message} (Urgency: {urgency})")
async def _notify_management(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify management"""
# Mock implementation
logger.info(f"Management Notification: {message} (Urgency: {urgency})")
async def _notify_slack(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify Slack"""
# Mock implementation
logger.info(f"Slack Notification: {message} (Urgency: {urgency})")
async def _notify_teams(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify Teams"""
# Mock implementation
logger.info(f"Teams Notification: {message} (Urgency: {urgency})")
async def _notify_email(self, message: str, urgency: str, context: Dict[str, Any]):
"""Notify via email"""
# Mock implementation
logger.info(f"Email Notification: {message} (Urgency: {urgency})")
# Additional private methods - Full implementations
async def _initialize_automation_templates(self):
"""Initialize automation templates with default templates"""
# Default automation templates for common use cases
default_templates = {
'security_alert_response': {
'name': 'Security Alert Response',
'description': 'Automatically respond to security alerts based on severity',
'type': WorkflowAutomationType.SECURITY.value,
'conditions': [
{
'type': AutomationConditionType.SECURITY_ALERT.value,
'severity': ['high', 'critical']
}
],
'actions': [
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['security_team'],
'urgency': 'high'
}
},
{
'type': AutomationActionType.WORKFLOW_EXECUTION.value,
'config': {
'workflow_id': 'security_incident_response'
}
}
],
'priority': AutomationPriority.HIGH.value,
'enabled': True
},
'compliance_violation_handling': {
'name': 'Compliance Violation Handling',
'description': 'Handle compliance violations automatically',
'type': WorkflowAutomationType.COMPLIANCEANCE.value,
'conditions': [
{
'type': AutomationConditionType.COMPLIANCE_VIOLATION.value,
'standards': ['SOC2', 'HIPAA', 'GDPR']
}
],
'actions': [
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['compliance_team'],
'urgency': 'critical'
}
},
{
'type': AutomationActionType.AUDITING.value,
'config': {
'audit_type': 'compliance_violation'
}
}
],
'priority': AutomationPriority.CRITICAL.value,
'enabled': True
},
'daily_security_scan': {
'name': 'Daily Security Scan',
'description': 'Run daily security scans',
'type': WorkflowAutomationType.SECURITY.value,
'conditions': [
{
'type': AutomationConditionType.SCHEDULED.value,
'schedule': '0 2 * * *' # 2 AM daily
}
],
'actions': [
{
'type': AutomationActionType.WORKFLOW_EXECUTION.value,
'config': {
'workflow_id': 'security_scan_workflow'
}
}
],
'priority': AutomationPriority.MEDIUM.value,
'enabled': False # Disabled by default
},
'user_access_review': {
'name': 'User Access Review',
'description': 'Review user access permissions periodically',
'type': WorkflowAutomationType.ACCESS_CONTROL.value,
'conditions': [
{
'type': AutomationConditionType.SCHEDULED.value,
'schedule': '0 9 * * 1' # 9 AM every Monday
}
],
'actions': [
{
'type': AutomationActionType.NOTIFICATION.value,
'config': {
'channels': ['admin_team'],
'urgency': 'medium'
}
},
{
'type': AutomationActionType.REPORTING.value,
'config': {
'report_type': 'user_access_report'
}
}
],
'priority': AutomationPriority.MEDIUM.value,
'enabled': False
}
}
# Load templates from database if available, otherwise use defaults
if self.db:
try:
from sqlalchemy import text
result = self.db.execute(text("SELECT data FROM automation_templates WHERE active = :active"), {"active": True})
for row in result:
template_data = json.loads(row[0]) if isinstance(row[0], str) else row[0]
if 'template_id' in template_data:
self.automation_templates[template_data['template_id']] = template_data
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.warning(f"Could not load templates from database: {e}")
# Add default templates
self.automation_templates.update(default_templates)
logger.info(f"Initialized {len(self.automation_templates)} automation templates")
return True
except Exception as e:
logger.error(f"Error initializing automation templates: {e}")
return False
async def _load_automations(self):
"""Load existing automations from database"""
if not self.db:
logger.warning("No database connection, skipping automation load")
return False
from sqlalchemy import text
# Load active automations
result = self.db.execute(text("""
SELECT automation_id, name, description, type, conditions, actions,
priority, status, enabled, created_by, created_at, updated_at,
schedule, next_run, last_run, execution_count, success_count,
failure_count, last_execution_status, metadata
FROM workflow_automations
WHERE status IN (:active, :paused)
ORDER BY created_at DESC
"""), {
"active": AutomationStatus.ACTIVE.value,
"paused": AutomationStatus.PAUSED.value
})
for row in result:
automation = WorkflowAutomation(
automation_id=row[0],
name=row[1],
description=row[2],
type=row[3],
conditions=json.loads(row[4]) if row[4] else [],
actions=json.loads(row[5]) if row[5] else [],
priority=row[6],
status=AutomationStatus(row[7]),
enabled=row[8],
created_by=row[9],
created_at=datetime.fromisoformat(row[10]) if row[10] else datetime.utcnow(),
updated_at=datetime.fromisoformat(row[11]) if row[11] else datetime.utcnow(),
schedule=row[12],
next_run=datetime.fromisoformat(row[13]) if row[13] else None,
last_run=datetime.fromisoformat(row[14]) if row[14] else None,
execution_count=row[15] or 0,
success_count=row[16] or 0,
failure_count=row[17] or 0,
last_execution_status=row[18],
metadata=json.loads(row[19]) if row[19] else {}
)
self.automations[automation.automation_id] = automation
# Schedule automation if it has a schedule and is enabled
if automation.enabled and automation.schedule and automation.next_run:
await self._schedule_automation(automation, {'type': AutomationConditionType.SCHEDULED.value})
logger.info(f"Loaded {len(self.automations)} automations from database")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error loading automations: {e}")
return False
async def _initialize_automation_scheduling(self):
"""Initialize automation scheduling system"""
if self.scheduler_running:
logger.warning("Scheduler already running")
return True
# Start the scheduler task
self.scheduler_task = asyncio.create_task(self._scheduler_loop())
self.scheduler_running = True
logger.info("Automation scheduling initialized")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error initializing automation scheduling: {e}")
return False
async def _scheduler_loop(self):
"""Background scheduler loop"""
while self.scheduler_running:
now = datetime.utcnow()
# Check automations that need to run
for automation_id, automation in self.automations.items():
if automation.enabled and automation.next_run:
if automation.next_run <= now:
logger.info(f"Running scheduled automation: {automation_id}")
await self.execute_automation(
automation_id=automation_id,
trigger_context={'trigger_type': 'scheduled'}
)
# Sleep for a short interval before checking again
await asyncio.sleep(60) # Check every minute
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in scheduler loop: {e}")
await asyncio.sleep(60) # Wait before retrying
async def _initialize_trigger_listeners(self):
"""Initialize trigger listeners for event-based automations"""
# Register event listeners for different trigger types
event_types = [
AutomationConditionType.EVENT_TRIGGERED.value,
AutomationConditionType.SECURITY_ALERT.value,
AutomationConditionType.COMPLIANCE_VIOLATION.value,
AutomationConditionType.THRESHOLD_EXCEEDED.value,
AutomationConditionType.ANOMALY_DETECTED.value,
AutomationConditionType.SYSTEM_EVENT.value,
AutomationConditionType.USER_ACTION.value,
AutomationConditionType.DATA_CHANGED.value
]
for event_type in event_types:
self.trigger_listeners[event_type] = {
'automations': [],
'callback': self._handle_event_trigger
}
# Find automations with event triggers and register them
for automation_id, automation in self.automations.items():
for condition in automation.conditions:
if condition['type'] in event_types:
if automation_id not in self.trigger_listeners[condition['type']]['automations']:
self.trigger_listeners[condition['type']]['automations'].append(automation_id)
logger.info(f"Initialized trigger listeners for {len(self.trigger_listeners)} event types")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error initializing trigger listeners: {e}")
return False
async def _handle_event_trigger(self, event_type: str, event_data: Dict[str, Any]):
"""Handle an event trigger"""
if event_type not in self.trigger_listeners:
logger.warning(f"Unknown event type: {event_type}")
return
listener = self.trigger_listeners[event_type]
automation_ids = listener['automations']
for automation_id in automation_ids:
if automation_id in self.automations:
automation = self.automations[automation_id]
if automation.enabled:
await self.execute_automation(
automation_id=automation_id,
trigger_context=event_data
)
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error handling event trigger: {e}")
async def _initialize_integration_endpoints(self):
"""Initialize integration endpoints for platform-specific automations"""
# Validate platform integrations
for platform_name, integration in self.platform_integrations.items():
if integration:
try:
# Test the integration
if hasattr(integration, 'test_connection'):
is_connected = await integration.test_connection()
logger.info(f"Platform {platform_name} integration: {'connected' if is_connected else 'disconnected'}")
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.warning(f"Could not validate {platform_name} integration: {e}")
logger.info("Integration endpoints initialized")
return True
except Exception as e:
logger.error(f"Error initializing integration endpoints: {e}")
return False
async def _start_automation_monitoring(self):
"""Start background automation monitoring"""
# Start monitoring task
asyncio.create_task(self._monitoring_loop())
logger.info("Automation monitoring started")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error starting automation monitoring: {e}")
return False
async def _monitoring_loop(self):
"""Background monitoring loop for automation health"""
while True:
# Update metrics
self.automation_metrics['total_automations'] = len(self.automations)
self.automation_metrics['active_automations'] = sum(
1 for auto in self.automations.values() if auto.enabled and auto.status == AutomationStatus.ACTIVE
)
# Check for failed automations
for automation_id, automation in self.automations.items():
if automation.last_execution_status == 'failed':
# Check if failure rate is high
if automation.execution_count > 0:
failure_rate = automation.failure_count / automation.execution_count
if failure_rate > 0.5: # More than 50% failure rate
logger.warning(f"Automation {automation_id} has high failure rate: {failure_rate:.2%}")
# Sleep for 5 minutes between checks
await asyncio.sleep(300)
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error in monitoring loop: {e}")
await asyncio.sleep(300)
async def _schedule_automation(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Schedule automation based on condition"""
if condition['type'] == AutomationConditionType.SCHEDULED.value:
schedule = automation.schedule or condition.get('schedule')
if schedule:
# Calculate next run time based on cron schedule
# This is a simplified implementation - use a proper cron library in production
from datetime import timedelta
# For now, just schedule for next day at same time
if automation.next_run:
next_run = automation.next_run + timedelta(days=1)
else:
next_run = datetime.utcnow() + timedelta(days=1)
automation.next_run = next_run
self.scheduled_automations[automation.automation_id] = {
'schedule': schedule,
'next_run': next_run.isoformat()
}
logger.info(f"Scheduled automation {automation.automation_id} for {next_run}")
return True
return False
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error scheduling automation: {e}")
return False
async def _setup_event_trigger(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Setup event-based trigger"""
event_type = condition.get('event_type', condition.get('type'))
if not event_type:
logger.warning(f"No event type specified for automation {automation.automation_id}")
return False
# Register automation for event type
if event_type not in self.trigger_listeners:
self.trigger_listeners[event_type] = {
'automations': [],
'callback': self._handle_event_trigger
}
if automation.automation_id not in self.trigger_listeners[event_type]['automations']:
self.trigger_listeners[event_type]['automations'].append(automation.automation_id)
self.active_triggers[automation.automation_id] = {
'type': 'event',
'event_type': event_type,
'condition': condition,
'enabled': automation.enabled
}
logger.info(f"Setup event trigger for automation {automation.automation_id}: {event_type}")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up event trigger: {e}")
return False
async def _setup_threshold_trigger(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Setup threshold-based trigger"""
metric = condition.get('metric')
threshold = condition.get('threshold')
operator = condition.get('operator', 'gt') # gt, lt, gte, lte, eq
if not metric or threshold is None:
logger.warning(f"Invalid threshold condition for automation {automation.automation_id}")
return False
self.active_triggers[automation.automation_id] = {
'type': 'threshold',
'metric': metric,
'threshold': threshold,
'operator': operator,
'condition': condition,
'enabled': automation.enabled
}
logger.info(f"Setup threshold trigger for automation {automation.automation_id}: {metric} {operator} {threshold}")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up threshold trigger: {e}")
return False
async def _setup_anomaly_trigger(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Setup anomaly detection trigger"""
metric = condition.get('metric')
sensitivity = condition.get('sensitivity', 'medium') # low, medium, high
if not metric:
logger.warning(f"Invalid anomaly condition for automation {automation.automation_id}")
return False
self.active_triggers[automation.automation_id] = {
'type': 'anomaly',
'metric': metric,
'sensitivity': sensitivity,
'condition': condition,
'enabled': automation.enabled
}
logger.info(f"Setup anomaly trigger for automation {automation.automation_id}: {metric} (sensitivity: {sensitivity})")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up anomaly trigger: {e}")
return False
async def _setup_security_trigger(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Setup security event trigger"""
threat_type = condition.get('threat_type')
severity = condition.get('severity', 'medium') # low, medium, high, critical
self.active_triggers[automation.automation_id] = {
'type': 'security',
'threat_type': threat_type,
'severity': severity,
'condition': condition,
'enabled': automation.enabled
}
# Register with security service if available
if self.security_service and hasattr(self.security_service, 'register_security_trigger'):
await self.security_service.register_security_trigger(
automation_id=automation.automation_id,
threat_type=threat_type,
severity=severity,
callback=lambda event: self.execute_automation(automation.automation_id, event)
)
logger.info(f"Setup security trigger for automation {automation.automation_id}: {threat_type} (severity: {severity})")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up security trigger: {e}")
return False
async def _setup_compliance_trigger(self, automation: WorkflowAutomation, condition: Dict[str, Any]):
"""Setup compliance violation trigger"""
standard = condition.get('standard') # SOC2, HIPAA, GDPR, etc.
violation_type = condition.get('violation_type')
self.active_triggers[automation.automation_id] = {
'type': 'compliance',
'standard': standard,
'violation_type': violation_type,
'condition': condition,
'enabled': automation.enabled
}
# Register with unified service if available
if self.unified_service and hasattr(self.unified_service, 'register_compliance_trigger'):
await self.unified_service.register_compliance_trigger(
automation_id=automation.automation_id,
standard=standard,
violation_type=violation_type,
callback=lambda event: self.execute_automation(automation.automation_id, event)
)
logger.info(f"Setup compliance trigger for automation {automation.automation_id}: {standard} - {violation_type}")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up compliance trigger: {e}")
return False
async def _setup_platform_triggers(self, platform: str, automation_id: str, config: Dict[str, Any]):
"""Setup platform-specific triggers"""
if platform not in self.platform_integrations:
logger.warning(f"Unknown platform: {platform}")
return False
integration = self.platform_integrations[platform]
if not integration:
logger.warning(f"Platform {platform} integration not available")
return False
# Setup platform-specific triggers based on config
trigger_type = config.get('trigger_type')
if trigger_type == 'webhook':
# Register webhook with platform
if hasattr(integration, 'register_webhook'):
webhook_url = config.get('webhook_url')
events = config.get('events', [])
await integration.register_webhook(webhook_url, events)
logger.info(f"Registered webhook for {platform}: {webhook_url}")
elif trigger_type == 'polling':
# Setup polling interval
if hasattr(integration, 'start_polling'):
interval = config.get('polling_interval', 300) # 5 minutes default
await integration.start_polling(automation_id, interval)
logger.info(f"Started polling for {platform} with interval {interval}s")
elif trigger_type == 'event_subscription':
# Subscribe to platform events
if hasattr(integration, 'subscribe_to_events'):
events = config.get('events', [])
await integration.subscribe_to_events(automation_id, events)
logger.info(f"Subscribed to events for {platform}: {events}")
logger.info(f"Setup platform triggers for {platform}: {automation_id}")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error setting up platform triggers: {e}")
return False
async def _send_automation_notifications(self, automation: WorkflowAutomation, execution: AutomationExecution):
"""Send notifications based on automation execution"""
# Get notification rules from automation metadata
notification_rules = automation.metadata.get('notification_rules', [])
if not notification_rules:
# Default notification behavior
if execution.status == AutomationStatus.FAILED:
await self._notify_via_slack(
message=f"Automation {automation.name} failed: {execution.error}",
urgency='high'
)
return
# Process each notification rule
for rule in notification_rules:
should_notify = False
# Check if rule matches execution status
if rule.get('status') == execution.status.value:
should_notify = True
# Check if rule is for errors
if rule.get('on_error') and execution.error:
should_notify = True
if should_notify:
channels = rule.get('channels', [])
message = rule.get('message', f"Automation {automation.name} executed with status: {execution.status.value}")
urgency = rule.get('urgency', 'medium')
# Send to each channel
for channel in channels:
if channel.startswith('slack:'):
await self._notify_via_slack(message, urgency)
elif channel.startswith('email:'):
await self._notify_via_email(message, urgency)
elif channel.startswith('teams:'):
await self._notify_via_teams(message, urgency)
logger.info(f"Sent notifications for automation {automation.automation_id}")
return True
except Exception as e:
logger.error(f"Operation failed: {e}")
log_integration_complete(audit_ctx, error=e)
return {'ok': False, 'error': str(e)}
logger.error(f"Error sending automation notifications: {e}")
return False
async def _update_automation_metrics(self, automation: WorkflowAutomation, execution: AutomationExecution):
"""Update automation metrics"""
self.automation_metrics['executed_today'] += 1
self.automation_metrics['executed_this_week'] += 1
self.automation_metrics['executed_this_month'] += 1
# Update success rate
total_executions = sum(self.automation_metrics['executions_by_status'].values())
if total_executions > 0:
successful_executions = self.automation_metrics['executions_by_status'].get('completed', 0)
self.automation_metrics['success_rate'] = successful_executions / total_executions
# Update average execution time
if execution.execution_time > 0:
self.automation_metrics['average_execution_time'] = (
(self.automation_metrics['average_execution_time'] * (total_executions - 1) + execution.execution_time)
/ total_executions
)
async def _log_automation_event(self, automation_id: str, event_type: str, user_id: str, details: Dict[str, Any]):
"""Log automation event"""
if self.security_service:
await self.security_service.audit_event({
'event_type': event_type,
'user_id': user_id,
'resource': 'automation',
'action': 'log',
'result': 'success',
'metadata': {
'automation_id': automation_id,
'details': details
}
})
async def get_service_info(self) -> Dict[str, Any]:
"""Get workflow automation service information"""
return {
"name": "Workflow Automation Service",
"version": "6.0.0",
"description": "Comprehensive workflow automation integrating all enterprise services",
"features": [
"multi_platform_integration",
"security_automations",
"compliance_automations",
"trigger_based_execution",
"scheduled_executions",
"ai_enhanced_automations",
"error_handling",
"monitoring",
"analytics"
],
"supported_automation_types": [t.value for t in WorkflowAutomationType],
"supported_action_types": [t.value for t in AutomationActionType],
"supported_condition_types": [t.value for t in AutomationConditionType],
"supported_priorities": [t.value for t in AutomationPriority],
"status": "ACTIVE"
}
async def close(self):
"""Close workflow automation service"""
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "get_service_info", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)
# Stop scheduler
if self.scheduler_task:
self.scheduler_task.cancel()
# Close HTTP sessions
for session in self.http_sessions.values():
await session.close()
logger.info("Workflow Automation Service closed")
# Global workflow automation service instance
# Initialize with None values - will be configured when dependencies are available
atom_workflow_automation_service = AtomWorkflowAutomationService({
'database': None, # Would be actual database connection
'cache': None, # Would be actual cache client
'security_service': atom_enterprise_security_service if 'atom_enterprise_security_service' in globals() else None,
'unified_service': atom_enterprise_unified_service if 'atom_enterprise_unified_service' in globals() else None,
'workflow_service': None, # Would be actual workflow service
'ai_service': ai_enhanced_service if 'ai_enhanced_service' in globals() else None,
'ai_integration': atom_ai_integration if 'atom_ai_integration' in globals() else None
})
except Exception as e:
logger.warning(f"Could not initialize global workflow automation service: {e}")
atom_workflow_automation_service = None
# Start audit logging
audit_ctx = log_integration_attempt("atom_workflow_automation", "close", locals())
# Check circuit breaker
if not await circuit_breaker.is_enabled("atom_workflow_automation"):
logger.warning(f"Circuit breaker is open for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Circuit breaker open"))
raise HTTPException(
status_code=503,
detail=f"Atom_workflow_automation integration temporarily disabled"
)
# Check rate limiter
is_limited, remaining = await rate_limiter.is_rate_limited("atom_workflow_automation")
if is_limited:
logger.warning(f"Rate limit exceeded for atom_workflow_automation")
log_integration_complete(audit_ctx, error=Exception("Rate limit exceeded"))
raise HTTPException(
status_code=429,
detail=f"Rate limit exceeded for atom_workflow_automation"
)