Spaces:
Sleeping
Sleeping
| """ | |
| ATOM Zoom Integration | |
| Advanced Zoom platform integration with enterprise features and 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 | |
| # Import existing ATOM services | |
| try: | |
| from ai_enhanced_service import ( | |
| AIModelType, | |
| AIRequest, | |
| AIResponse, | |
| AIServiceType, | |
| AITaskType, | |
| ai_enhanced_service, | |
| ) | |
| from atom_ai_integration import atom_ai_integration | |
| from atom_discord_integration import atom_discord_integration | |
| from atom_enterprise_security_service import atom_enterprise_security_service | |
| from atom_enterprise_unified_service import atom_enterprise_unified_service | |
| from atom_google_chat_integration import atom_google_chat_integration | |
| from atom_ingestion_pipeline import AtomIngestionPipeline | |
| from atom_memory_service import AtomMemoryService | |
| from atom_search_service import AtomSearchService | |
| from atom_slack_integration import atom_slack_integration | |
| from atom_teams_integration import atom_teams_integration | |
| from atom_telegram_integration import atom_telegram_integration | |
| from atom_whatsapp_integration import atom_whatsapp_integration | |
| from atom_workflow_automation_service import ( | |
| AutomationPriority, | |
| AutomationStatus, | |
| atom_workflow_automation_service, | |
| ) | |
| from atom_workflow_service import AtomWorkflowService | |
| from atom_zoom_integration import atom_zoom_integration | |
| from integrations.atom_ingestion_pipeline import RecordType, atom_ingestion_pipeline | |
| except ImportError as e: | |
| logging.warning(f"Enterprise services not available: {e}") | |
| # Configure logging | |
| logger = logging.getLogger(__name__) | |
| class ZoomEventType(Enum): | |
| """Zoom event types""" | |
| MEETING_STARTED = "meeting.started" | |
| MEETING_ENDED = "meeting.ended" | |
| MEETING_PARTICIPANT_JOINED = "meeting.participant_joined" | |
| MEETING_PARTICIPANT_LEFT = "meeting.participant_left" | |
| RECORDING_COMPLETED = "recording.completed" | |
| WEBINAR_STARTED = "webinar.started" | |
| WEBINAR_ENDED = "webinar.ended" | |
| USER_CREATED = "user.created" | |
| USER_UPDATED = "user.updated" | |
| USER_DELETED = "user.deleted" | |
| class ZoomMeetingType(Enum): | |
| """Zoom meeting types""" | |
| INSTANT = "instant" | |
| SCHEDULED = "scheduled" | |
| RECURRING = "recurring" | |
| RECURRING_WITH_FIXED_TIME = "recurring_with_fixed_time" | |
| class ZoomUserType(Enum): | |
| """Zoom user types""" | |
| BASIC = "basic" | |
| LICENSED = "licensed" | |
| ON_PREM = "on_prem" | |
| OFFLINE_ON_PREM = "offline_on_prem" | |
| class ZoomCommandType(Enum): | |
| """Zoom command types""" | |
| START = "start" | |
| SCHEDULE = "schedule" | |
| JOIN = "join" | |
| LEAVE = "leave" | |
| RECORD = "record" | |
| SHARE = "share" | |
| CHAT = "chat" | |
| STATUS = "status" | |
| ANALYTICS = "analytics" | |
| AUTOMATE = "automate" | |
| SECURITY = "security" | |
| COMPLIANCE = "compliance" | |
| class ZoomUser: | |
| """Zoom user data model""" | |
| user_id: str | |
| email: str | |
| first_name: str | |
| last_name: str | |
| display_name: str | |
| user_type: ZoomUserType | |
| role: str | |
| timezone: str | |
| is_active: bool | |
| permissions: List[str] | |
| security_level: str | |
| created_at: datetime | |
| last_active: datetime | |
| metadata: Dict[str, Any] | |
| class ZoomMeeting: | |
| """Zoom meeting data model""" | |
| meeting_id: str | |
| topic: str | |
| meeting_type: ZoomMeetingType | |
| host_id: str | |
| start_time: datetime | |
| duration: int | |
| timezone: str | |
| agenda: str | |
| participants: List[str] | |
| is_recorded: bool | |
| password: Optional[str] | |
| waiting_room: bool | |
| security_level: str | |
| created_at: datetime | |
| status: str | |
| metadata: Dict[str, Any] | |
| class ZoomEvent: | |
| """Zoom event data model""" | |
| event_id: str | |
| event_type: ZoomEventType | |
| meeting_id: Optional[str] | |
| user_id: Optional[str] | |
| timestamp: datetime | |
| data: Dict[str, Any] | |
| security_flags: Dict[str, Any] | |
| compliance_flags: Dict[str, Any] | |
| metadata: Dict[str, Any] | |
| class AtomZoomIntegration: | |
| """Advanced Zoom integration with enterprise features""" | |
| def __init__(self, config: Dict[str, Any]): | |
| self.config = config | |
| self.db = config.get('database') | |
| self.cache = config.get('cache') | |
| # Zoom configuration | |
| self.zoom_config = { | |
| 'api_key': config.get('api_key') or os.getenv('ZOOM_API_KEY'), | |
| 'api_secret': config.get('api_secret') or os.getenv('ZOOM_API_SECRET'), | |
| 'webhook_secret': config.get('webhook_secret') or os.getenv('ZOOM_WEBHOOK_SECRET'), | |
| 'webhook_url': config.get('webhook_url') or os.getenv('ZOOM_WEBHOOK_URL'), | |
| 'admin_emails': config.get('admin_emails', []), | |
| # 'account_id': config.get('account_id') or os.getenv('ZOOM_ACCOUNT_ID'), # Deprecated: Using User OAuth | |
| 'client_id': config.get('client_id') or os.getenv('ZOOM_CLIENT_ID'), | |
| 'client_secret': config.get('client_secret') or os.getenv('ZOOM_CLIENT_SECRET'), | |
| 'redirect_uri': config.get('redirect_uri') or os.getenv('ZOOM_REDIRECT_URI'), | |
| 'max_participants': config.get('max_participants', 1000), | |
| 'enable_enterprise_features': config.get('enable_enterprise_features', True), | |
| 'security_level': config.get('security_level', 'standard'), | |
| 'compliance_standards': config.get('compliance_standards', ['SOC2', 'ISO27001', 'HIPAA']), | |
| 'api_base_url': config.get('api_base_url', 'https://api.zoom.us/v2'), | |
| 'oauth_token_url': config.get('oauth_token_url', 'https://zoom.us/oauth/token') | |
| } | |
| # Integration state | |
| self.is_initialized = False | |
| self.active_users: Dict[str, ZoomUser] = {} | |
| self.active_meetings: Dict[str, ZoomMeeting] = {} | |
| self.meeting_history: Dict[str, List[ZoomEvent]] = {} | |
| self.webhook_handlers: Dict[str, Callable] = {} | |
| self.command_handlers: Dict[str, Callable] = {} | |
| # Enterprise integration (use safe defaults if optional services didn't import) | |
| self.enterprise_security = config.get('security_service') or globals().get('atom_enterprise_security_service') | |
| self.enterprise_automation = config.get('automation_service') or globals().get('atom_workflow_automation_service') | |
| self.ai_service = config.get('ai_service') or globals().get('ai_enhanced_service') | |
| # Analytics and monitoring | |
| self.analytics_metrics = { | |
| 'total_meetings': 0, | |
| 'total_participants': 0, | |
| 'total_users': 0, | |
| 'total_recording_hours': 0, | |
| 'meetings_today': 0, | |
| 'meetings_this_week': 0, | |
| 'meetings_this_month': 0, | |
| 'commands_executed': 0, | |
| 'automations_triggered': 0, | |
| 'security_incidents': 0, | |
| 'compliance_checks': 0, | |
| 'ai_requests': 0, | |
| 'meeting_types': defaultdict(int), | |
| 'user_types': defaultdict(int), | |
| 'active_meetings': 0, | |
| 'active_users': 0 | |
| } | |
| # Security and compliance | |
| self.security_policies = {} | |
| self.compliance_rules = {} | |
| self.automation_triggers = {} | |
| self.meeting_filters = [] | |
| # Performance metrics | |
| self.performance_metrics = { | |
| 'meeting_start_time': 0.0, | |
| 'participant_join_time': 0.0, | |
| 'command_response_time': 0.0, | |
| 'ai_processing_time': 0.0, | |
| 'security_check_time': 0.0, | |
| 'compliance_check_time': 0.0, | |
| 'automation_execution_time': 0.0, | |
| 'webhook_response_time': 0.0 | |
| } | |
| # HTTP session for API calls | |
| self.http_session = httpx.AsyncClient(timeout=30.0) | |
| self.oauth_token = None | |
| self.oauth_token_expires = None | |
| logger.info("Zoom Integration initialized") | |
| async def initialize(self) -> bool: | |
| """Initialize Zoom integration""" | |
| try: | |
| if not (self.zoom_config['api_key'] and self.zoom_config['api_secret']) and \ | |
| not (self.zoom_config['client_id'] and self.zoom_config['client_secret']): | |
| logger.error("Zoom API credentials not provided") | |
| return False | |
| # Test API connection | |
| await self._test_api_connection() | |
| # Setup webhook | |
| if self.zoom_config['webhook_url']: | |
| await self._setup_webhook() | |
| # Setup webhook handlers | |
| await self._setup_webhook_handlers() | |
| # Setup enterprise features | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._setup_enterprise_features() | |
| # Setup security and compliance | |
| await self._setup_security_and_compliance() | |
| # Setup automation | |
| await self._setup_automation() | |
| # Setup monitoring | |
| await self._setup_monitoring() | |
| # Load existing data | |
| await self._load_existing_data() | |
| self.is_initialized = True | |
| logger.info("Zoom Integration initialized successfully") | |
| return True | |
| except Exception as e: | |
| logger.error(f"Error initializing Zoom integration: {e}") | |
| return False | |
| async def get_intelligent_workspaces(self, user_id: str) -> List[Dict[str, Any]]: | |
| """Get intelligent workspaces for user""" | |
| try: | |
| workspaces = [] | |
| # Get user's meetings as workspaces | |
| user_meetings = [m for m in self.active_meetings.values() if m.host_id == user_id] | |
| for meeting in user_meetings: | |
| workspace = { | |
| 'id': meeting.meeting_id, | |
| 'name': meeting.topic, | |
| 'type': 'zoom_meeting', | |
| 'participant_count': len(meeting.participants), | |
| 'description': meeting.agenda, | |
| 'start_time': meeting.start_time.isoformat(), | |
| 'duration': meeting.duration, | |
| 'status': meeting.status, | |
| 'security_level': meeting.security_level, | |
| 'is_recorded': meeting.is_recorded, | |
| 'is_active': meeting.status == 'started', | |
| 'permissions': { | |
| 'can_join': user_id in meeting.participants or user_id == meeting.host_id, | |
| 'can_start': user_id == meeting.host_id, | |
| 'can_record': user_id == meeting.host_id, | |
| 'can_manage': user_id == meeting.host_id | |
| }, | |
| 'platform': 'zoom' | |
| } | |
| workspaces.append(workspace) | |
| # Sort by start time | |
| workspaces.sort(key=lambda x: x['start_time'], reverse=True) | |
| return workspaces | |
| except Exception as e: | |
| logger.error(f"Error getting intelligent workspaces: {e}") | |
| return [] | |
| async def get_intelligent_channels(self, workspace_id: str, user_id: str) -> List[Dict[str, Any]]: | |
| """Get intelligent channels for workspace""" | |
| try: | |
| channels = [] | |
| meeting = self.active_meetings.get(workspace_id) | |
| if meeting: | |
| # For Zoom, workspace is meeting itself | |
| channel = { | |
| 'id': meeting.meeting_id, | |
| 'name': meeting.topic, | |
| 'type': 'zoom_meeting', | |
| 'participant_count': len(meeting.participants), | |
| 'description': meeting.agenda, | |
| 'start_time': meeting.start_time.isoformat(), | |
| 'duration': meeting.duration, | |
| 'status': meeting.status, | |
| 'security_level': meeting.security_level, | |
| 'is_recorded': meeting.is_recorded, | |
| 'is_active': meeting.status == 'started', | |
| 'permissions': { | |
| 'can_join': user_id in meeting.participants or user_id == meeting.host_id, | |
| 'can_start': user_id == meeting.host_id, | |
| 'can_record': user_id == meeting.host_id, | |
| 'can_manage': user_id == meeting.host_id | |
| }, | |
| 'participants': meeting.participants, | |
| 'host_id': meeting.host_id, | |
| 'meeting_type': meeting.meeting_type.value, | |
| 'platform': 'zoom' | |
| } | |
| channels.append(channel) | |
| return channels | |
| except Exception as e: | |
| logger.error(f"Error getting intelligent channels: {e}") | |
| return [] | |
| async def send_intelligent_message(self, channel_id: str, message: str, metadata: Dict[str, Any] = None) -> Dict[str, Any]: | |
| """Send intelligent message""" | |
| try: | |
| # Send message to Zoom meeting chat | |
| meeting = self.active_meetings.get(channel_id) | |
| if not meeting: | |
| return { | |
| 'success': False, | |
| 'error': 'Meeting not found', | |
| 'platform': 'zoom' | |
| } | |
| # Send chat message via Zoom API | |
| result = await self._send_chat_message(channel_id, message) | |
| if result['success']: | |
| # Log message | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._log_message_event('chat_message_sent', channel_id, { | |
| 'message': message, | |
| 'timestamp': datetime.utcnow().isoformat(), | |
| 'metadata': metadata or {} | |
| }) | |
| return result | |
| except Exception as e: | |
| logger.error(f"Error sending intelligent message: {e}") | |
| return { | |
| 'success': False, | |
| 'error': str(e), | |
| 'platform': 'zoom' | |
| } | |
| async def perform_intelligent_search(self, query: str, user_id: str, workspace_id: str = None) -> List[Dict[str, Any]]: | |
| """Perform intelligent search""" | |
| try: | |
| search_results = [] | |
| # Search in meetings | |
| for meeting_id, meeting in self.active_meetings.items(): | |
| if workspace_id and meeting_id != workspace_id: | |
| continue | |
| if query.lower() in meeting.topic.lower() or \ | |
| (meeting.agenda and query.lower() in meeting.agenda.lower()): | |
| result = { | |
| 'id': meeting.meeting_id, | |
| 'type': 'zoom_meeting', | |
| 'title': meeting.topic, | |
| 'snippet': meeting.agenda[:100] + "..." if meeting.agenda and len(meeting.agenda) > 100 else meeting.agenda or "No agenda", | |
| 'content': f"Meeting: {meeting.topic}\n{meeting.agenda or ''}", | |
| 'channel_id': meeting.meeting_id, | |
| 'host_id': meeting.host_id, | |
| 'timestamp': meeting.start_time.isoformat(), | |
| 'meeting_type': meeting.meeting_type.value, | |
| 'platform': 'zoom', | |
| 'relevance_score': self._calculate_relevance_score(query, meeting.topic + (meeting.agenda or '')) | |
| } | |
| search_results.append(result) | |
| # Sort by relevance score | |
| search_results.sort(key=lambda x: x['relevance_score'], reverse=True) | |
| # AI-enhanced search | |
| if self.ai_service: | |
| ai_results = await self._perform_ai_search(query, workspace_id) | |
| search_results.extend(ai_results) | |
| return search_results[:20] # Return top 20 results | |
| except Exception as e: | |
| logger.error(f"Error performing intelligent search: {e}") | |
| return [] | |
| async def get_user_conversation_history(self, user_id: str, channel_id: str, limit: int = 50) -> List[Dict[str, Any]]: | |
| """Get user conversation history""" | |
| try: | |
| conversation_history = [] | |
| events = self.meeting_history.get(channel_id, []) | |
| # Filter by user and limit | |
| user_events = [e for e in events if e.user_id == user_id][-limit:] | |
| for event in user_events: | |
| history_item = { | |
| 'id': event.event_id, | |
| 'content': f"{event.event_type.value}: {event.data}", | |
| 'event_type': event.event_type.value, | |
| 'timestamp': event.timestamp.isoformat(), | |
| 'channel_id': event.meeting_id or '', | |
| 'platform': 'zoom', | |
| 'metadata': event.metadata | |
| } | |
| conversation_history.append(history_item) | |
| return conversation_history | |
| except Exception as e: | |
| logger.error(f"Error getting user conversation history: {e}") | |
| return [] | |
| async def get_service_status(self) -> Dict[str, Any]: | |
| """Get Zoom service status""" | |
| try: | |
| return { | |
| 'platform': 'zoom', | |
| 'status': 'active' if self.is_initialized else 'inactive', | |
| 'webhook_url': self.zoom_config['webhook_url'], | |
| 'enterprise_features': self.zoom_config['enable_enterprise_features'], | |
| 'security_level': self.zoom_config['security_level'], | |
| 'compliance_standards': self.zoom_config['compliance_standards'], | |
| 'total_meetings': self.analytics_metrics['total_meetings'], | |
| 'total_users': self.analytics_metrics['total_users'], | |
| 'total_participants': self.analytics_metrics['total_participants'], | |
| 'total_recording_hours': self.analytics_metrics['total_recording_hours'], | |
| 'commands_executed': self.analytics_metrics['commands_executed'], | |
| 'automations_triggered': self.analytics_metrics['automations_triggered'], | |
| 'security_incidents': self.analytics_metrics['security_incidents'], | |
| 'compliance_checks': self.analytics_metrics['compliance_checks'], | |
| 'ai_requests': self.analytics_metrics['ai_requests'], | |
| 'active_meetings': self.analytics_metrics['active_meetings'], | |
| 'active_users': self.analytics_metrics['active_users'], | |
| 'performance_metrics': self.performance_metrics, | |
| 'uptime': time.time() - (self._start_time if hasattr(self, '_start_time') else time.time()) | |
| } | |
| except Exception as e: | |
| logger.error(f"Error getting service status: {e}") | |
| return {'error': str(e), 'platform': 'zoom'} | |
| # Private helper methods | |
| async def _test_api_connection(self): | |
| """Test Zoom API connection""" | |
| try: | |
| # Get OAuth token | |
| await self._get_oauth_token() | |
| # Test API call | |
| if self.oauth_token: | |
| api_url = f"{self.zoom_config['api_base_url']}/users/me" | |
| headers = {'Authorization': f'Bearer {self.oauth_token}'} | |
| response = await self.http_session.get(api_url, headers=headers) | |
| if response.status_code == 200: | |
| logger.info("Zoom API connection verified") | |
| else: | |
| logger.error(f"Zoom API connection failed: {response.status_code}") | |
| else: | |
| logger.warning("Skipping API test - No OAuth token available (User OAuth flow required)") | |
| except Exception as e: | |
| logger.error(f"Error testing API connection: {e}") | |
| async def _get_oauth_token(self): | |
| """Get OAuth token""" | |
| try: | |
| if self.oauth_token and self.oauth_token_expires and datetime.utcnow() < self.oauth_token_expires: | |
| return | |
| # Deprecated: Server-to-Server OAuth (Client Credentials) | |
| # This flow requires ZOOM_ACCOUNT_ID which is no longer used. | |
| # We now use User Managed OAuth (Authorization Code) which requires user interaction. | |
| # This method now just logs a warning if no token is present, as tokens should be | |
| # injected or retrieved via the ZoomAuthHandler after user login. | |
| logger.info("User Managed OAuth configured. Token should be provided via user authentication flow.") | |
| # In a real implementation, we would fetch the stored user token from DB here | |
| # For now, we avoid breaking the initialization by not attempting the client_credentials flow | |
| except Exception as e: | |
| logger.error(f"Error getting OAuth token: {e}") | |
| async def _setup_webhook(self): | |
| """Setup webhook""" | |
| try: | |
| # Setup webhook via Zoom API | |
| await self._get_oauth_token() | |
| webhook_url = f"{self.zoom_config['api_base_url']}/webhooks" | |
| headers = {'Authorization': f'Bearer {self.oauth_token}'} | |
| webhook_data = { | |
| 'event_url': self.zoom_config['webhook_url'], | |
| 'event_types': [ | |
| 'meeting.started', 'meeting.ended', | |
| 'meeting.participant_joined', 'meeting.participant_left', | |
| 'recording.completed' | |
| ], | |
| 'authorization_token': self.zoom_config['webhook_secret'] | |
| } | |
| response = await self.http_session.post(webhook_url, json=webhook_data, headers=headers) | |
| if response.status_code == 201: | |
| logger.info(f"Zoom webhook setup complete: {self.zoom_config['webhook_url']}") | |
| else: | |
| logger.error(f"Zoom webhook setup failed: {response.status_code}") | |
| except Exception as e: | |
| logger.error(f"Error setting up webhook: {e}") | |
| async def _setup_webhook_handlers(self): | |
| """Setup webhook handlers""" | |
| try: | |
| # Define webhook handlers | |
| async def handle_meeting_started(event_data: Dict[str, Any]): | |
| await self._handle_meeting_started(event_data) | |
| async def handle_meeting_ended(event_data: Dict[str, Any]): | |
| await self._handle_meeting_ended(event_data) | |
| async def handle_participant_joined(event_data: Dict[str, Any]): | |
| await self._handle_participant_joined(event_data) | |
| async def handle_participant_left(event_data: Dict[str, Any]): | |
| await self._handle_participant_left(event_data) | |
| async def handle_recording_completed(event_data: Dict[str, Any]): | |
| await self._handle_recording_completed(event_data) | |
| # Register handlers | |
| self.webhook_handlers = { | |
| ZoomEventType.MEETING_STARTED: handle_meeting_started, | |
| ZoomEventType.MEETING_ENDED: handle_meeting_ended, | |
| ZoomEventType.MEETING_PARTICIPANT_JOINED: handle_participant_joined, | |
| ZoomEventType.MEETING_PARTICIPANT_LEFT: handle_participant_left, | |
| ZoomEventType.RECORDING_COMPLETED: handle_recording_completed | |
| } | |
| logger.info("Webhook handlers setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up webhook handlers: {e}") | |
| async def _setup_enterprise_features(self): | |
| """Setup enterprise features""" | |
| try: | |
| if not self.enterprise_security or not self.enterprise_automation: | |
| logger.warning("Enterprise services not available") | |
| return | |
| # Setup security policies | |
| await self._setup_security_policies() | |
| # Setup compliance rules | |
| await self._setup_compliance_rules() | |
| # Setup automation triggers | |
| await self._setup_automation_triggers() | |
| logger.info("Enterprise features setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up enterprise features: {e}") | |
| async def _setup_security_policies(self): | |
| """Setup security policies""" | |
| try: | |
| self.security_policies = { | |
| 'meeting_access_control': { | |
| 'enabled': True, | |
| 'require_password': True, | |
| 'waiting_room': True, | |
| 'action': 'restrict' | |
| }, | |
| 'participant_authentication': { | |
| 'enabled': True, | |
| 'verified_domains': [], | |
| 'blocked_users': [], | |
| 'action': 'restrict' | |
| }, | |
| 'meeting_security': { | |
| 'enabled': True, | |
| 'encryption_required': True, | |
| 'recording_consent': True, | |
| 'action': 'monitor' | |
| } | |
| } | |
| logger.info("Security policies setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up security policies: {e}") | |
| async def _setup_compliance_rules(self): | |
| """Setup compliance rules""" | |
| try: | |
| self.compliance_rules = { | |
| 'recording_compliance': { | |
| 'enabled': True, | |
| 'consent_required': True, | |
| 'retention_period': 365, | |
| 'action': 'comply' | |
| }, | |
| 'meeting_compliance': { | |
| 'enabled': True, | |
| 'participant_logging': True, | |
| 'duration_logging': True, | |
| 'action': 'log' | |
| }, | |
| 'data_protection': { | |
| 'enabled': True, | |
| 'data_encryption': True, | |
| 'access_logging': True, | |
| 'action': 'protect' | |
| } | |
| } | |
| logger.info("Compliance rules setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up compliance rules: {e}") | |
| async def _setup_automation_triggers(self): | |
| """Setup automation triggers""" | |
| try: | |
| self.automation_triggers = { | |
| 'meeting_started': { | |
| 'enabled': True, | |
| 'conditions': ['host_role', 'meeting_type', 'participant_count'], | |
| 'actions': ['send_notification', 'start_recording', 'ai_analysis'] | |
| }, | |
| 'meeting_ended': { | |
| 'enabled': True, | |
| 'conditions': ['duration', 'participant_count', 'recording_status'], | |
| 'actions': ['generate_report', 'archive_data', 'compliance_check'] | |
| }, | |
| 'participant_joined': { | |
| 'enabled': True, | |
| 'conditions': ['user_verification', 'meeting_permissions'], | |
| 'actions': ['welcome_message', 'role_assignment', 'security_check'] | |
| } | |
| } | |
| logger.info("Automation triggers setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up automation triggers: {e}") | |
| async def _setup_automation(self): | |
| """Setup automation""" | |
| try: | |
| if not self.enterprise_automation: | |
| logger.warning("Automation service not available") | |
| return | |
| # Create Zoom-specific automation | |
| zoom_automation_data = { | |
| 'name': 'Zoom Meeting Automation', | |
| 'description': 'Automation for Zoom meeting events', | |
| 'automation_type': 'integration', | |
| 'priority': 'medium', | |
| 'conditions': [ | |
| { | |
| 'type': 'event_triggered', | |
| 'platform': 'zoom', | |
| 'events': ['meeting_started', 'meeting_ended', 'participant_joined', 'recording_completed'] | |
| } | |
| ], | |
| 'actions': [ | |
| { | |
| 'type': 'notification', | |
| 'config': { | |
| 'channels': ['platform_admin'], | |
| 'message': 'Zoom meeting event occurred', | |
| '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': 'zoom', | |
| 'integration_version': '1.0.0' | |
| } | |
| } | |
| result = await self.enterprise_automation.create_integration_automation('zoom', zoom_automation_data) | |
| if result.get('ok'): | |
| logger.info("Zoom automation created successfully") | |
| else: | |
| logger.error(f"Failed to create Zoom automation: {result.get('error')}") | |
| except Exception as e: | |
| logger.error(f"Error setting up automation: {e}") | |
| async def _setup_security_and_compliance(self): | |
| """Setup security and compliance monitoring""" | |
| try: | |
| # Setup monitoring for security events | |
| if self.zoom_config['enable_enterprise_features']: | |
| # Security monitoring | |
| await self._setup_security_monitoring() | |
| # Compliance monitoring | |
| await self._setup_compliance_monitoring() | |
| logger.info("Security and compliance setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up security and compliance: {e}") | |
| async def _setup_security_monitoring(self): | |
| """Setup security monitoring""" | |
| try: | |
| self.security_monitoring = { | |
| 'meeting_anomaly_detection': { | |
| 'enabled': True, | |
| 'threshold': 0.8, | |
| 'action': 'alert' | |
| }, | |
| 'participant_behavior_analysis': { | |
| 'enabled': True, | |
| 'baseline_period': 30, | |
| 'action': 'monitor' | |
| }, | |
| 'meeting_security_monitoring': { | |
| 'enabled': True, | |
| 'security_score_threshold': 0.7, | |
| 'action': 'flag' | |
| } | |
| } | |
| logger.info("Security monitoring setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up security monitoring: {e}") | |
| async def _setup_compliance_monitoring(self): | |
| """Setup compliance monitoring""" | |
| try: | |
| self.compliance_monitoring = { | |
| 'meeting_compliance_checking': { | |
| 'enabled': True, | |
| 'check_frequency': 'real_time', | |
| 'action': 'flag' | |
| }, | |
| 'participant_activity_auditing': { | |
| 'enabled': True, | |
| 'audit_retention': 365, | |
| 'action': 'log' | |
| }, | |
| 'recording_compliance_management': { | |
| 'enabled': True, | |
| 'retention_policy': 'standard', | |
| 'action': 'manage' | |
| } | |
| } | |
| logger.info("Compliance monitoring setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up compliance monitoring: {e}") | |
| async def _setup_monitoring(self): | |
| """Setup monitoring""" | |
| try: | |
| self.performance_metrics = { | |
| 'meeting_start_time': 0.0, | |
| 'participant_join_time': 0.0, | |
| 'command_response_time': 0.0, | |
| 'ai_processing_time': 0.0, | |
| 'security_check_time': 0.0, | |
| 'compliance_check_time': 0.0, | |
| 'automation_execution_time': 0.0, | |
| 'webhook_response_time': 0.0 | |
| } | |
| self._start_time = time.time() | |
| logger.info("Monitoring setup complete") | |
| except Exception as e: | |
| logger.error(f"Error setting up monitoring: {e}") | |
| async def _load_existing_data(self): | |
| """Load existing data""" | |
| try: | |
| # Mock implementation - would load from database | |
| logger.info("Existing data loaded") | |
| except Exception as e: | |
| logger.error(f"Error loading existing data: {e}") | |
| async def _send_chat_message(self, meeting_id: str, message: str) -> Dict[str, Any]: | |
| """Send chat message to meeting""" | |
| try: | |
| await self._get_oauth_token() | |
| api_url = f"{self.zoom_config['api_base_url']}/meetings/{meeting_id}/chat" | |
| headers = {'Authorization': f'Bearer {self.oauth_token}'} | |
| chat_data = { | |
| 'to_all': True, | |
| 'message': message | |
| } | |
| response = await self.http_session.post(api_url, json=chat_data, headers=headers) | |
| if response.status_code == 201: | |
| return { | |
| 'success': True, | |
| 'message_id': response.json().get('message_id'), | |
| 'timestamp': datetime.utcnow().isoformat() | |
| } | |
| else: | |
| return { | |
| 'success': False, | |
| 'error': response.json().get('message', 'Unknown error') | |
| } | |
| except Exception as e: | |
| logger.error(f"Error sending chat message: {e}") | |
| return { | |
| 'success': False, | |
| 'error': str(e) | |
| } | |
| def _calculate_relevance_score(self, query: str, content: str) -> float: | |
| """Calculate relevance score for search""" | |
| try: | |
| query_words = query.lower().split() | |
| content_words = content.lower().split() | |
| matches = 0 | |
| for word in query_words: | |
| if word in content_words: | |
| matches += 1 | |
| return matches / len(query_words) if query_words else 0.0 | |
| except Exception as e: | |
| logger.error(f"Error calculating relevance score: {e}") | |
| return 0.0 | |
| async def _perform_ai_search(self, query: str, workspace_id: str = None) -> List[Dict[str, Any]]: | |
| """Perform AI-enhanced search""" | |
| try: | |
| if not self.ai_service: | |
| return [] | |
| # Create AI request | |
| ai_request = AIRequest( | |
| request_id=f"zoom_search_{int(time.time())}", | |
| task_type=AITaskType.SEARCH_QUERY, | |
| model_type=AIModelType.GPT_4, | |
| service_type=AIServiceType.OPENAI, | |
| input_data={ | |
| 'query': query, | |
| 'context': 'zoom_search', | |
| 'workspace_id': workspace_id, | |
| 'meetings': [asdict(m) for m in self.active_meetings.values()] | |
| }, | |
| context={ | |
| 'platform': 'zoom', | |
| 'workspace_id': workspace_id | |
| }, | |
| platform='zoom' | |
| ) | |
| # Process AI request | |
| ai_response = await self.ai_service.process_ai_request(ai_request) | |
| if ai_response.ok and ai_response.output_data: | |
| return ai_response.output_data.get('results', []) | |
| return [] | |
| except Exception as e: | |
| logger.error(f"Error performing AI search: {e}") | |
| return [] | |
| async def _log_message_event(self, event_type: str, channel_id: str, data: Dict[str, Any]): | |
| """Log message event""" | |
| try: | |
| if self.enterprise_security: | |
| await self.enterprise_security.audit_event({ | |
| 'event_type': event_type, | |
| 'user_id': data.get('user_id'), | |
| 'resource': 'zoom_meeting', | |
| 'action': event_type, | |
| 'result': 'success', | |
| 'ip_address': 'zoom', | |
| 'user_agent': 'zoom_api', | |
| 'metadata': { | |
| 'meeting_id': channel_id, | |
| 'platform': 'zoom', | |
| 'data': data | |
| } | |
| }) | |
| except Exception as e: | |
| logger.error(f"Error logging message event: {e}") | |
| # Webhook event handlers | |
| async def _handle_meeting_started(self, event_data: Dict[str, Any]): | |
| """Handle meeting started event""" | |
| try: | |
| meeting_id = event_data.get('payload', {}).get('object', {}).get('id') | |
| topic = event_data.get('payload', {}).get('object', {}).get('topic') | |
| host_id = event_data.get('payload', {}).get('object', {}).get('host_id') | |
| # Update analytics | |
| self.analytics_metrics['total_meetings'] += 1 | |
| self.analytics_metrics['meetings_today'] += 1 | |
| self.analytics_metrics['active_meetings'] += 1 | |
| # Create meeting object | |
| meeting = ZoomMeeting( | |
| meeting_id=meeting_id, | |
| topic=topic, | |
| meeting_type=ZoomMeetingType.INSTANT, | |
| host_id=host_id, | |
| start_time=datetime.utcnow(), | |
| duration=0, | |
| timezone='UTC', | |
| agenda=topic, | |
| participants=[host_id], | |
| is_recorded=False, | |
| password=None, | |
| waiting_room=False, | |
| security_level=self.zoom_config['security_level'], | |
| created_at=datetime.utcnow(), | |
| status='started', | |
| metadata={} | |
| ) | |
| self.active_meetings[meeting_id] = meeting | |
| # Trigger automation | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._trigger_automations('meeting_started', meeting, event_data) | |
| # Ingest to ATOM memory (LanceDB) | |
| try: | |
| atom_ingestion_pipeline.ingest_record("zoom", RecordType.MEETING.value, asdict(meeting)) | |
| logger.info(f"Zoom meeting {meeting_id} ingested to memory") | |
| except Exception as me: | |
| logger.error(f"Failed to ingest Zoom meeting to memory: {me}") | |
| logger.info(f"Meeting started: {meeting_id} - {topic}") | |
| except Exception as e: | |
| logger.error(f"Error handling meeting started event: {e}") | |
| async def _handle_meeting_ended(self, event_data: Dict[str, Any]): | |
| """Handle meeting ended event""" | |
| try: | |
| meeting_id = event_data.get('payload', {}).get('object', {}).get('id') | |
| # Update analytics | |
| self.analytics_metrics['active_meetings'] = max(0, self.analytics_metrics['active_meetings'] - 1) | |
| # Update meeting status | |
| if meeting_id in self.active_meetings: | |
| meeting = self.active_meetings[meeting_id] | |
| meeting.status = 'ended' | |
| meeting.duration = int((datetime.utcnow() - meeting.start_time).total_seconds()) | |
| # Trigger automation | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._trigger_automations('meeting_ended', meeting, event_data) | |
| # Update memory | |
| try: | |
| atom_ingestion_pipeline.ingest_record("zoom", RecordType.MEETING.value, asdict(meeting)) | |
| logger.info(f"Zoom meeting {meeting_id} (ended) updated in memory") | |
| except Exception as me: | |
| logger.error(f"Failed to update Zoom meeting in memory: {me}") | |
| logger.info(f"Meeting ended: {meeting_id}") | |
| except Exception as e: | |
| logger.error(f"Error handling meeting ended event: {e}") | |
| async def _handle_participant_joined(self, event_data: Dict[str, Any]): | |
| """Handle participant joined event""" | |
| try: | |
| meeting_id = event_data.get('payload', {}).get('object', {}).get('id') | |
| participant_id = event_data.get('payload', {}).get('object', {}).get('participant', {}).get('id') | |
| participant_name = event_data.get('payload', {}).get('object', {}).get('participant', {}).get('user_name') | |
| # Update analytics | |
| self.analytics_metrics['total_participants'] += 1 | |
| # Update meeting participants | |
| if meeting_id in self.active_meetings: | |
| meeting = self.active_meetings[meeting_id] | |
| if participant_id not in meeting.participants: | |
| meeting.participants.append(participant_id) | |
| # Trigger automation | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._trigger_automations('participant_joined', meeting, event_data) | |
| logger.info(f"Participant joined: {participant_name} ({participant_id}) in meeting {meeting_id}") | |
| except Exception as e: | |
| logger.error(f"Error handling participant joined event: {e}") | |
| async def _handle_participant_left(self, event_data: Dict[str, Any]): | |
| """Handle participant left event""" | |
| try: | |
| meeting_id = event_data.get('payload', {}).get('object', {}).get('id') | |
| participant_id = event_data.get('payload', {}).get('object', {}).get('participant', {}).get('id') | |
| participant_name = event_data.get('payload', {}).get('object', {}).get('participant', {}).get('user_name') | |
| # Update meeting participants | |
| if meeting_id in self.active_meetings: | |
| meeting = self.active_meetings[meeting_id] | |
| if participant_id in meeting.participants: | |
| meeting.participants.remove(participant_id) | |
| logger.info(f"Participant left: {participant_name} ({participant_id}) from meeting {meeting_id}") | |
| except Exception as e: | |
| logger.error(f"Error handling participant left event: {e}") | |
| async def _handle_recording_completed(self, event_data: Dict[str, Any]): | |
| """Handle recording completed event""" | |
| try: | |
| meeting_id = event_data.get('payload', {}).get('object', {}).get('id') | |
| recording_files = event_data.get('payload', {}).get('object', {}).get('recording_files', []) | |
| # Update analytics | |
| total_duration = sum(f.get('recording_length', 0) for f in recording_files) | |
| self.analytics_metrics['total_recording_hours'] += total_duration / 3600 | |
| # Update meeting recording status | |
| if meeting_id in self.active_meetings: | |
| meeting = self.active_meetings[meeting_id] | |
| meeting.is_recorded = True | |
| meeting.metadata['recording_files'] = recording_files | |
| # Trigger automation | |
| if self.zoom_config['enable_enterprise_features']: | |
| await self._trigger_automations('recording_completed', meeting, event_data) | |
| # Update memory with recording info | |
| try: | |
| atom_ingestion_pipeline.ingest_record("zoom", RecordType.MEETING.value, asdict(meeting)) | |
| logger.info(f"Zoom meeting {meeting_id} (recording completed) updated in memory") | |
| except Exception as me: | |
| logger.error(f"Failed to update Zoom meeting in memory: {me}") | |
| logger.info(f"Recording completed for meeting: {meeting_id}") | |
| except Exception as e: | |
| logger.error(f"Error handling recording completed event: {e}") | |
| async def _trigger_automations(self, event_type: str, meeting: ZoomMeeting, event_data: Dict[str, Any]): | |
| """Trigger automations""" | |
| try: | |
| if not self.enterprise_automation: | |
| return | |
| # Find relevant automations | |
| relevant_automations = [] | |
| for automation_id, automation in self.automation_triggers.items(): | |
| if automation['enabled'] and event_type in automation['name']: | |
| relevant_automations.append(automation) | |
| # Execute automations | |
| for automation in relevant_automations: | |
| trigger_context = { | |
| 'event_type': event_type, | |
| 'meeting': asdict(meeting), | |
| 'event_data': event_data, | |
| 'platform': 'zoom' | |
| } | |
| # Execute automation | |
| # This would integrate with automation service | |
| # await self.enterprise_automation.trigger_automation(automation_id, trigger_context) | |
| # Update metrics | |
| self.analytics_metrics['automations_triggered'] += len(relevant_automations) | |
| except Exception as e: | |
| logger.error(f"Error triggering automations: {e}") | |
| async def close(self): | |
| """Close Zoom integration""" | |
| try: | |
| if self.http_session: | |
| await self.http_session.aclose() | |
| logger.info("Zoom Integration closed") | |
| except Exception as e: | |
| logger.error(f"Error closing Zoom integration: {e}") | |
| # Global Zoom integration instance | |
| _zoom_config = { | |
| 'api_key': os.getenv('ZOOM_API_KEY'), | |
| 'api_secret': os.getenv('ZOOM_API_SECRET'), | |
| 'client_id': os.getenv('ZOOM_CLIENT_ID'), | |
| 'client_secret': os.getenv('ZOOM_CLIENT_SECRET'), | |
| 'webhook_url': os.getenv('ZOOM_WEBHOOK_URL'), | |
| 'webhook_secret': os.getenv('ZOOM_WEBHOOK_SECRET'), | |
| 'enable_enterprise_features': True, | |
| 'security_level': 'standard', | |
| 'compliance_standards': ['SOC2', 'ISO27001', 'HIPAA'], | |
| 'database': None, # Would be actual database connection | |
| 'cache': None, # Would be actual cache client | |
| } | |
| # Add optional services if they were imported successfully | |
| _security_service = globals().get('atom_enterprise_security_service') | |
| if _security_service: | |
| _zoom_config['security_service'] = _security_service | |
| _automation_service = globals().get('atom_workflow_automation_service') | |
| if _automation_service: | |
| _zoom_config['automation_service'] = _automation_service | |
| _ai_service = globals().get('ai_enhanced_service') | |
| if _ai_service: | |
| _zoom_config['ai_service'] = _ai_service | |
| atom_zoom_integration = AtomZoomIntegration(_zoom_config) |