""" 🏢 ENTERPRISE FEATURES - DataVision Production-Ready Capabilities ================================================================== Enterprise-grade features: - Action Engine (exports, alerts, webhooks) - Audit Logging - Rate Limiting - Multi-user Support - Scheduled Reports Ready for production deployment. """ import json import logging import os from typing import Dict, List, Optional, Any, Callable from dataclasses import dataclass, field from datetime import datetime, timedelta from enum import Enum import hashlib import threading from collections import defaultdict logger = logging.getLogger(__name__) # ============================================================================= # ACTION ENGINE - Exports, Alerts, Webhooks # ============================================================================= class ActionType(Enum): """Types of actions""" EXPORT_CSV = "export_csv" EXPORT_EXCEL = "export_excel" EXPORT_PDF = "export_pdf" SEND_EMAIL = "send_email" WEBHOOK = "webhook" SLACK = "slack" SCHEDULE = "schedule" @dataclass class ActionResult: """Result of an action""" success: bool action_type: ActionType message: str output_path: Optional[str] = None metadata: Dict[str, Any] = field(default_factory=dict) class ActionEngine: """ ⚡ Action Engine Execute actions on data and insights: - Export to CSV, Excel, PDF - Send email alerts - Trigger webhooks - Schedule recurring actions """ def __init__(self, storage_path: str = "storage/exports"): self.storage_path = storage_path os.makedirs(storage_path, exist_ok=True) self.scheduled_actions: Dict[str, Dict] = {} async def execute( self, action_type: ActionType, data: Any, config: Dict[str, Any], user_id: str = "default" ) -> ActionResult: """ Execute an action Args: action_type: Type of action data: Data to process config: Action configuration user_id: User identifier Returns: Action result """ try: if action_type == ActionType.EXPORT_CSV: return await self._export_csv(data, config, user_id) elif action_type == ActionType.EXPORT_EXCEL: return await self._export_excel(data, config, user_id) elif action_type == ActionType.EXPORT_PDF: return await self._export_pdf(data, config, user_id) elif action_type == ActionType.SEND_EMAIL: return await self._send_email(data, config) elif action_type == ActionType.WEBHOOK: return await self._trigger_webhook(data, config) elif action_type == ActionType.SCHEDULE: return await self._schedule_action(config, user_id) else: return ActionResult( success=False, action_type=action_type, message=f"Unknown action type: {action_type}" ) except Exception as e: logger.error(f"Action execution error: {e}") return ActionResult( success=False, action_type=action_type, message=str(e) ) async def _export_csv( self, data: Any, config: Dict, user_id: str ) -> ActionResult: """Export data to CSV""" import pandas as pd filename = config.get("filename", f"export_{datetime.now().strftime('%Y%m%d_%H%M%S')}.csv") filepath = os.path.join(self.storage_path, user_id, filename) os.makedirs(os.path.dirname(filepath), exist_ok=True) if isinstance(data, pd.DataFrame): data.to_csv(filepath, index=False) elif isinstance(data, dict): pd.DataFrame([data]).to_csv(filepath, index=False) elif isinstance(data, list): pd.DataFrame(data).to_csv(filepath, index=False) else: return ActionResult( success=False, action_type=ActionType.EXPORT_CSV, message="Invalid data type for CSV export" ) return ActionResult( success=True, action_type=ActionType.EXPORT_CSV, message=f"Exported to {filename}", output_path=filepath ) async def _export_excel( self, data: Any, config: Dict, user_id: str ) -> ActionResult: """Export data to Excel""" import pandas as pd try: filename = config.get("filename", f"export_{datetime.now().strftime('%Y%m%d_%H%M%S')}.xlsx") filepath = os.path.join(self.storage_path, user_id, filename) os.makedirs(os.path.dirname(filepath), exist_ok=True) if isinstance(data, pd.DataFrame): data.to_excel(filepath, index=False, engine='openpyxl') else: return ActionResult( success=False, action_type=ActionType.EXPORT_EXCEL, message="Invalid data type for Excel export" ) return ActionResult( success=True, action_type=ActionType.EXPORT_EXCEL, message=f"Exported to {filename}", output_path=filepath ) except ImportError: return ActionResult( success=False, action_type=ActionType.EXPORT_EXCEL, message="openpyxl not installed. Use: pip install openpyxl" ) async def _export_pdf( self, data: Any, config: Dict, user_id: str ) -> ActionResult: """Export report to PDF""" try: from reportlab.lib import colors from reportlab.lib.pagesizes import letter from reportlab.platypus import SimpleDocTemplate, Table, TableStyle, Paragraph from reportlab.lib.styles import getSampleStyleSheet filename = config.get("filename", f"report_{datetime.now().strftime('%Y%m%d_%H%M%S')}.pdf") filepath = os.path.join(self.storage_path, user_id, filename) os.makedirs(os.path.dirname(filepath), exist_ok=True) doc = SimpleDocTemplate(filepath, pagesize=letter) elements = [] styles = getSampleStyleSheet() # Add title title = config.get("title", "DataVision Report") elements.append(Paragraph(title, styles['Title'])) # Add content if isinstance(data, str): elements.append(Paragraph(data, styles['Normal'])) elif isinstance(data, dict): for key, value in data.items(): elements.append(Paragraph(f"{key}: {value}", styles['Normal'])) doc.build(elements) return ActionResult( success=True, action_type=ActionType.EXPORT_PDF, message=f"Report exported to {filename}", output_path=filepath ) except ImportError: return ActionResult( success=False, action_type=ActionType.EXPORT_PDF, message="reportlab not installed. Use: pip install reportlab" ) async def _send_email(self, data: Any, config: Dict) -> ActionResult: """Send email (placeholder - integrate with email service)""" to_email = config.get("to") subject = config.get("subject", "DataVision Notification") if not to_email: return ActionResult( success=False, action_type=ActionType.SEND_EMAIL, message="No recipient email specified" ) # Placeholder - would integrate with email service logger.info(f"Email would be sent to {to_email}: {subject}") return ActionResult( success=True, action_type=ActionType.SEND_EMAIL, message=f"Email queued for {to_email}", metadata={"to": to_email, "subject": subject} ) async def _trigger_webhook(self, data: Any, config: Dict) -> ActionResult: """Trigger a webhook""" try: import httpx url = config.get("url") if not url: return ActionResult( success=False, action_type=ActionType.WEBHOOK, message="No webhook URL specified" ) headers = config.get("headers", {"Content-Type": "application/json"}) async with httpx.AsyncClient() as client: response = await client.post( url, json=data if isinstance(data, dict) else {"data": str(data)}, headers=headers, timeout=30 ) return ActionResult( success=response.is_success, action_type=ActionType.WEBHOOK, message=f"Webhook triggered: {response.status_code}", metadata={"status_code": response.status_code} ) except ImportError: return ActionResult( success=False, action_type=ActionType.WEBHOOK, message="httpx not installed. Use: pip install httpx" ) except Exception as e: return ActionResult( success=False, action_type=ActionType.WEBHOOK, message=str(e) ) async def _schedule_action(self, config: Dict, user_id: str) -> ActionResult: """Schedule a recurring action""" schedule_id = config.get("id", hashlib.md5(str(config).encode()).hexdigest()[:8]) self.scheduled_actions[schedule_id] = { "user_id": user_id, "config": config, "created_at": datetime.now().isoformat(), "next_run": config.get("next_run"), "frequency": config.get("frequency", "daily") } return ActionResult( success=True, action_type=ActionType.SCHEDULE, message=f"Action scheduled: {schedule_id}", metadata={"schedule_id": schedule_id} ) # ============================================================================= # AUDIT LOGGING # ============================================================================= @dataclass class AuditEntry: """An audit log entry""" timestamp: str user_id: str action: str resource: str details: Dict[str, Any] ip_address: Optional[str] = None success: bool = True class AuditLogger: """ 📋 Audit Logger Track all user actions for compliance: - Query logs - Data access - Exports - Configuration changes """ def __init__(self, storage_path: str = "storage/audit"): self.storage_path = storage_path os.makedirs(storage_path, exist_ok=True) self.logs: List[AuditEntry] = [] self._lock = threading.Lock() def log( self, user_id: str, action: str, resource: str, details: Dict[str, Any] = None, ip_address: str = None, success: bool = True ): """Log an action""" entry = AuditEntry( timestamp=datetime.now().isoformat(), user_id=user_id, action=action, resource=resource, details=details or {}, ip_address=ip_address, success=success ) with self._lock: self.logs.append(entry) # Persist every 100 entries if len(self.logs) >= 100: self._persist_logs() def _persist_logs(self): """Persist logs to disk""" if not self.logs: return filename = f"audit_{datetime.now().strftime('%Y%m%d')}.jsonl" filepath = os.path.join(self.storage_path, filename) with open(filepath, 'a') as f: for entry in self.logs: f.write(json.dumps({ "timestamp": entry.timestamp, "user_id": entry.user_id, "action": entry.action, "resource": entry.resource, "details": entry.details, "ip_address": entry.ip_address, "success": entry.success }) + "\n") self.logs = [] def get_user_logs( self, user_id: str, start_date: datetime = None, end_date: datetime = None, action_filter: str = None ) -> List[Dict]: """Get logs for a user""" results = [] for entry in self.logs: if entry.user_id != user_id: continue if action_filter and action_filter not in entry.action: continue entry_time = datetime.fromisoformat(entry.timestamp) if start_date and entry_time < start_date: continue if end_date and entry_time > end_date: continue results.append({ "timestamp": entry.timestamp, "action": entry.action, "resource": entry.resource, "success": entry.success }) return results def flush(self): """Force persist all logs""" with self._lock: self._persist_logs() # ============================================================================= # RATE LIMITER — Delegated to core.rate_limiter # ============================================================================= # The unified rate limiter lives in core/rate_limiter.py (Redis + in-memory). # This module re-exports a convenience function for backward compatibility. # ============================================================================= # MULTI-USER SESSION MANAGER # ============================================================================= @dataclass class UserSession: """User session data""" user_id: str session_id: str created_at: datetime last_activity: datetime metadata: Dict[str, Any] = field(default_factory=dict) class SessionManager: """ 👥 Multi-User Session Manager Manage user sessions: - Session creation/validation - Activity tracking - Session cleanup """ def __init__(self, session_timeout_minutes: int = 60): self.session_timeout = timedelta(minutes=session_timeout_minutes) self.sessions: Dict[str, UserSession] = {} self._lock = threading.Lock() def create_session(self, user_id: str, metadata: Dict = None) -> str: """Create a new session""" session_id = hashlib.sha256( f"{user_id}{datetime.now().isoformat()}{os.urandom(16).hex()}".encode() ).hexdigest()[:32] session = UserSession( user_id=user_id, session_id=session_id, created_at=datetime.now(), last_activity=datetime.now(), metadata=metadata or {} ) with self._lock: self.sessions[session_id] = session return session_id def validate_session(self, session_id: str) -> Optional[UserSession]: """Validate and refresh a session""" with self._lock: session = self.sessions.get(session_id) if not session: return None # Check timeout if datetime.now() - session.last_activity > self.session_timeout: del self.sessions[session_id] return None # Refresh activity session.last_activity = datetime.now() return session def end_session(self, session_id: str) -> bool: """End a session""" with self._lock: if session_id in self.sessions: del self.sessions[session_id] return True return False def get_active_sessions(self, user_id: str) -> List[Dict]: """Get all active sessions for a user""" results = [] now = datetime.now() with self._lock: for sid, session in list(self.sessions.items()): if session.user_id != user_id: continue if now - session.last_activity > self.session_timeout: del self.sessions[sid] continue results.append({ "session_id": sid, "created_at": session.created_at.isoformat(), "last_activity": session.last_activity.isoformat() }) return results def cleanup_expired(self): """Remove all expired sessions""" now = datetime.now() with self._lock: expired = [ sid for sid, session in self.sessions.items() if now - session.last_activity > self.session_timeout ] for sid in expired: del self.sessions[sid] return len(expired) # ============================================================================= # EXPORTS # ============================================================================= action_engine = ActionEngine() audit_logger = AuditLogger() # rate_limiter — now managed by core.rate_limiter module session_manager = SessionManager() async def export_data( data: Any, format: str, # csv, excel, pdf user_id: str, config: Dict = None ) -> Dict[str, Any]: """Quick function to export data""" action_map = { "csv": ActionType.EXPORT_CSV, "excel": ActionType.EXPORT_EXCEL, "pdf": ActionType.EXPORT_PDF } action_type = action_map.get(format.lower(), ActionType.EXPORT_CSV) result = await action_engine.execute(action_type, data, config or {}, user_id) return { "success": result.success, "message": result.message, "path": result.output_path } def log_action( user_id: str, action: str, resource: str, details: Dict = None ): """Quick function to log an action""" audit_logger.log(user_id, action, resource, details) def check_rate_limit(user_id: str) -> Dict[str, Any]: """Quick synchronous rate limit check (delegates to core.rate_limiter).""" from core.rate_limiter import get_rate_limiter, RATE_LIMITS import asyncio import time limiter = get_rate_limiter() limits = RATE_LIMITS["default"] key = f"user:{user_id}:default" try: loop = asyncio.get_event_loop() if loop.is_running(): # Can't await in a sync context inside a running loop return {"allowed": True, "retry_after": None} is_limited, remaining, retry_after = loop.run_until_complete( limiter.is_rate_limited(key, limits["max_requests"], limits["window_seconds"]) ) except RuntimeError: return {"allowed": True, "retry_after": None} return { "allowed": not is_limited, "retry_after": retry_after if is_limited else None, "requests_last_minute": limits["max_requests"] - remaining, "limit_per_minute": limits["max_requests"] }