| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562 |
- """
- Main event handler system for the Trixy application.
- This module provides the EventHandler class that manages event registration,
- triggering, and history tracking. It supports thread-safe operations for
- multi-satellite environments and provides comprehensive debugging capabilities.
- """
- import asyncio
- import threading
- import time
- import traceback
- from collections import deque, defaultdict
- from concurrent.futures import ThreadPoolExecutor, as_completed
- from dataclasses import dataclass, field
- from datetime import datetime, timedelta
- from typing import Any, Callable, Dict, List, Optional, Union, Set
- from enum import Enum
- from .event_data import TrixyEventData, EventType, EventDataFactory
- from .decorators import (
- get_event_registry, register_event_handlers, unregister_event_handlers,
- TrixyEventRegistry
- )
- def pprint(message: str) -> None:
- """
- Debug printing function that adapts based on mode.
- In production: uses proper logging, in debug: uses print.
- """
- # TODO: This should integrate with the application's logging system
- print(f"[EventHandler] {message}")
- class EventStatus(Enum):
- """Status of event processing."""
- PENDING = "pending"
- PROCESSING = "processing"
- COMPLETED = "completed"
- FAILED = "failed"
- CANCELLED = "cancelled"
- @dataclass
- class EventHistoryEntry:
- """Entry in the event history log."""
- event_id: str
- event_type: str
- event_data: TrixyEventData
- timestamp: datetime
- status: EventStatus = EventStatus.PENDING
- processing_start_time: Optional[datetime] = None
- processing_end_time: Optional[datetime] = None
- handlers_called: List[str] = field(default_factory=list)
- errors: List[str] = field(default_factory=list)
- execution_time_ms: float = 0.0
-
- def to_dict(self) -> Dict[str, Any]:
- """Convert history entry to dictionary."""
- return {
- 'event_id': self.event_id,
- 'event_type': self.event_type,
- 'event_data': self.event_data.to_dict() if self.event_data else None,
- 'timestamp': self.timestamp.isoformat(),
- 'status': self.status.value,
- 'processing_start_time': self.processing_start_time.isoformat() if self.processing_start_time else None,
- 'processing_end_time': self.processing_end_time.isoformat() if self.processing_end_time else None,
- 'handlers_called': self.handlers_called,
- 'errors': self.errors,
- 'execution_time_ms': self.execution_time_ms
- }
- class EventHandler:
- """
- Thread-safe event handler system for the Trixy application.
-
- This class manages event registration, triggering, and history tracking.
- It supports both synchronous and asynchronous event handlers and provides
- comprehensive debugging and monitoring capabilities.
- """
-
- def __init__(self, max_history_size: int = 1000, max_worker_threads: int = 10):
- """
- Initialize the event handler.
-
- Args:
- max_history_size: Maximum number of events to keep in history
- max_worker_threads: Maximum number of threads for async event processing
- """
- self._lock = threading.RLock()
- self._registry = get_event_registry()
-
- # Event history tracking
- self._history: deque[EventHistoryEntry] = deque(maxlen=max_history_size)
- self._history_by_type: Dict[str, deque[EventHistoryEntry]] = defaultdict(
- lambda: deque(maxlen=100)
- )
-
- # Statistics
- self._event_stats: Dict[str, Dict[str, int]] = defaultdict(
- lambda: {'triggered': 0, 'completed': 0, 'failed': 0}
- )
-
- # Thread pool for async operations
- self._executor = ThreadPoolExecutor(max_workers=max_worker_threads)
-
- # Event filtering and debugging
- self._enabled_events: Set[str] = set()
- self._disabled_events: Set[str] = set()
- self._debug_mode = False
- self._event_listeners: List[Callable] = []
-
- # Event counters for unique IDs
- self._event_counter = 0
- self._counter_lock = threading.Lock()
-
- pprint("EventHandler initialized")
-
- def _generate_event_id(self) -> str:
- """Generate a unique event ID."""
- with self._counter_lock:
- self._event_counter += 1
- return f"evt_{int(time.time())}_{self._event_counter:06d}"
-
- def enable_debug_mode(self, enabled: bool = True) -> None:
- """Enable or disable debug mode for verbose logging."""
- with self._lock:
- self._debug_mode = enabled
- pprint(f"Debug mode {'enabled' if enabled else 'disabled'}")
-
- def add_event_listener(self, listener: Callable[[EventHistoryEntry], None]) -> None:
- """
- Add an event listener that gets called for every event.
-
- Args:
- listener: Function that takes an EventHistoryEntry
- """
- with self._lock:
- if listener not in self._event_listeners:
- self._event_listeners.append(listener)
- pprint(f"Added event listener: {listener.__name__}")
-
- def remove_event_listener(self, listener: Callable) -> None:
- """Remove an event listener."""
- with self._lock:
- if listener in self._event_listeners:
- self._event_listeners.remove(listener)
- pprint(f"Removed event listener: {listener.__name__}")
-
- def enable_event_type(self, event_type: str) -> None:
- """Enable processing for a specific event type."""
- with self._lock:
- self._enabled_events.add(event_type)
- self._disabled_events.discard(event_type)
- pprint(f"Enabled event type: {event_type}")
-
- def disable_event_type(self, event_type: str) -> None:
- """Disable processing for a specific event type."""
- with self._lock:
- self._disabled_events.add(event_type)
- self._enabled_events.discard(event_type)
- pprint(f"Disabled event type: {event_type}")
-
- def is_event_type_enabled(self, event_type: str) -> bool:
- """Check if an event type is enabled for processing."""
- with self._lock:
- if self._disabled_events and event_type in self._disabled_events:
- return False
- if self._enabled_events and event_type not in self._enabled_events:
- return False
- return True
-
- def register_handler_object(self, obj: Any) -> None:
- """
- Register all @TrixyEvent decorated methods from an object.
-
- Args:
- obj: Object instance containing @TrixyEvent decorated methods
- """
- with self._lock:
- register_event_handlers(obj)
- pprint(f"Registered event handlers from {obj.__class__.__name__}")
-
- def unregister_handler_object(self, obj: Any) -> None:
- """
- Unregister all @TrixyEvent decorated methods from an object.
-
- Args:
- obj: Object instance to unregister handlers from
- """
- with self._lock:
- unregister_event_handlers(obj)
- pprint(f"Unregistered event handlers from {obj.__class__.__name__}")
-
- def trigger_event(
- self,
- event_type: Union[str, EventType],
- event_data: Optional[TrixyEventData] = None,
- **kwargs
- ) -> str:
- """
- Trigger an event synchronously.
-
- Args:
- event_type: Type of event to trigger
- event_data: Event data object (optional, will be created if not provided)
- **kwargs: Additional arguments for event data creation
-
- Returns:
- str: Event ID for tracking
- """
- # Normalize event type
- if isinstance(event_type, EventType):
- event_type_str = event_type.value
- else:
- event_type_str = str(event_type)
-
- # Check if event type is enabled
- if not self.is_event_type_enabled(event_type_str):
- if self._debug_mode:
- pprint(f"Event type {event_type_str} is disabled, skipping")
- return ""
-
- # Create event data if not provided
- if event_data is None:
- try:
- event_data = EventDataFactory.create_event_data(event_type_str, **kwargs)
- except ValueError as e:
- if self._debug_mode:
- pprint(f"Failed to create event data for {event_type_str}: {e}")
- # Create a generic event data object
- from .event_data import TrixyEventData
- event_data = TrixyEventData(**kwargs)
-
- # Generate event ID and create history entry
- event_id = self._generate_event_id()
- history_entry = EventHistoryEntry(
- event_id=event_id,
- event_type=event_type_str,
- event_data=event_data,
- timestamp=datetime.now()
- )
-
- try:
- self._process_event_sync(history_entry)
- except Exception as e:
- pprint(f"Error processing event {event_type_str}: {e}")
- if self._debug_mode:
- pprint(f"Traceback: {traceback.format_exc()}")
- history_entry.status = EventStatus.FAILED
- history_entry.errors.append(str(e))
-
- # Add to history and notify listeners
- with self._lock:
- self._history.append(history_entry)
- self._history_by_type[event_type_str].append(history_entry)
- self._event_stats[event_type_str]['triggered'] += 1
-
- if history_entry.status == EventStatus.COMPLETED:
- self._event_stats[event_type_str]['completed'] += 1
- elif history_entry.status == EventStatus.FAILED:
- self._event_stats[event_type_str]['failed'] += 1
-
- # Notify event listeners
- for listener in self._event_listeners:
- try:
- listener(history_entry)
- except Exception as e:
- pprint(f"Error in event listener {listener.__name__}: {e}")
-
- if self._debug_mode:
- pprint(f"Triggered event {event_type_str} with ID {event_id}")
-
- return event_id
-
- def trigger_event_async(
- self,
- event_type: Union[str, EventType],
- event_data: Optional[TrixyEventData] = None,
- **kwargs
- ) -> str:
- """
- Trigger an event asynchronously.
-
- Args:
- event_type: Type of event to trigger
- event_data: Event data object (optional, will be created if not provided)
- **kwargs: Additional arguments for event data creation
-
- Returns:
- str: Event ID for tracking
- """
- # Similar setup as sync version
- if isinstance(event_type, EventType):
- event_type_str = event_type.value
- else:
- event_type_str = str(event_type)
-
- if not self.is_event_type_enabled(event_type_str):
- if self._debug_mode:
- pprint(f"Event type {event_type_str} is disabled, skipping")
- return ""
-
- if event_data is None:
- try:
- event_data = EventDataFactory.create_event_data(event_type_str, **kwargs)
- except ValueError:
- from .event_data import TrixyEventData
- event_data = TrixyEventData(**kwargs)
-
- event_id = self._generate_event_id()
- history_entry = EventHistoryEntry(
- event_id=event_id,
- event_type=event_type_str,
- event_data=event_data,
- timestamp=datetime.now()
- )
-
- # Submit to thread pool
- future = self._executor.submit(self._process_event_sync, history_entry)
-
- # Add a callback to handle completion
- def on_completion(fut):
- try:
- fut.result() # This will raise any exception that occurred
- except Exception as e:
- pprint(f"Async event {event_type_str} failed: {e}")
- history_entry.status = EventStatus.FAILED
- history_entry.errors.append(str(e))
-
- # Update history and stats
- with self._lock:
- self._history.append(history_entry)
- self._history_by_type[event_type_str].append(history_entry)
- self._event_stats[event_type_str]['triggered'] += 1
-
- if history_entry.status == EventStatus.COMPLETED:
- self._event_stats[event_type_str]['completed'] += 1
- elif history_entry.status == EventStatus.FAILED:
- self._event_stats[event_type_str]['failed'] += 1
-
- # Notify event listeners
- for listener in self._event_listeners:
- try:
- listener(history_entry)
- except Exception as e:
- pprint(f"Error in event listener {listener.__name__}: {e}")
-
- future.add_done_callback(on_completion)
-
- if self._debug_mode:
- pprint(f"Triggered async event {event_type_str} with ID {event_id}")
-
- return event_id
-
- def _process_event_sync(self, history_entry: EventHistoryEntry) -> None:
- """
- Process an event synchronously.
-
- Args:
- history_entry: Event history entry to process
- """
- history_entry.status = EventStatus.PROCESSING
- history_entry.processing_start_time = datetime.now()
-
- handlers = self._registry.get_handlers(history_entry.event_type)
-
- if not handlers:
- if self._debug_mode:
- pprint(f"No handlers registered for event {history_entry.event_type}")
- history_entry.status = EventStatus.COMPLETED
- history_entry.processing_end_time = datetime.now()
- return
-
- for handler in handlers:
- try:
- handler_name = getattr(handler, '__qualname__', handler.__name__)
- history_entry.handlers_called.append(handler_name)
-
- if self._debug_mode:
- pprint(f"Calling handler {handler_name} for event {history_entry.event_type}")
-
- # Call the handler
- handler(history_entry.event_type, history_entry.event_data)
-
- except Exception as e:
- error_msg = f"Error in handler {handler_name}: {e}"
- history_entry.errors.append(error_msg)
- pprint(error_msg)
- if self._debug_mode:
- pprint(f"Traceback: {traceback.format_exc()}")
-
- history_entry.processing_end_time = datetime.now()
- if history_entry.processing_start_time:
- history_entry.execution_time_ms = (
- history_entry.processing_end_time - history_entry.processing_start_time
- ).total_seconds() * 1000
-
- history_entry.status = EventStatus.FAILED if history_entry.errors else EventStatus.COMPLETED
-
- def wait_for_event(self, event_id: str, timeout: float = 30.0) -> Optional[EventHistoryEntry]:
- """
- Wait for an event to complete processing.
-
- Args:
- event_id: Event ID to wait for
- timeout: Maximum time to wait in seconds
-
- Returns:
- EventHistoryEntry if found and completed, None if timeout or not found
- """
- start_time = time.time()
-
- while time.time() - start_time < timeout:
- with self._lock:
- for entry in self._history:
- if entry.event_id == event_id:
- if entry.status in [EventStatus.COMPLETED, EventStatus.FAILED, EventStatus.CANCELLED]:
- return entry
- break
-
- time.sleep(0.1) # Small delay to avoid busy waiting
-
- return None
-
- def get_event_history(
- self,
- event_type: Optional[str] = None,
- limit: Optional[int] = None,
- since: Optional[datetime] = None
- ) -> List[EventHistoryEntry]:
- """
- Get event history with optional filtering.
-
- Args:
- event_type: Filter by event type
- limit: Maximum number of entries to return
- since: Only return events after this timestamp
-
- Returns:
- List of event history entries
- """
- with self._lock:
- if event_type:
- entries = list(self._history_by_type.get(event_type, []))
- else:
- entries = list(self._history)
-
- # Filter by timestamp
- if since:
- entries = [e for e in entries if e.timestamp >= since]
-
- # Sort by timestamp (newest first)
- entries.sort(key=lambda e: e.timestamp, reverse=True)
-
- # Apply limit
- if limit:
- entries = entries[:limit]
-
- return entries
-
- def get_event_statistics(self) -> Dict[str, Dict[str, Any]]:
- """
- Get comprehensive event statistics.
-
- Returns:
- Dictionary containing event statistics
- """
- with self._lock:
- stats = {}
-
- for event_type, type_stats in self._event_stats.items():
- success_rate = 0.0
- if type_stats['triggered'] > 0:
- success_rate = type_stats['completed'] / type_stats['triggered'] * 100
-
- # Calculate average execution time
- recent_entries = list(self._history_by_type[event_type])[-10:] # Last 10 events
- avg_execution_time = 0.0
- if recent_entries:
- total_time = sum(e.execution_time_ms for e in recent_entries if e.execution_time_ms > 0)
- avg_execution_time = total_time / len(recent_entries)
-
- stats[event_type] = {
- 'triggered': type_stats['triggered'],
- 'completed': type_stats['completed'],
- 'failed': type_stats['failed'],
- 'success_rate_percent': round(success_rate, 2),
- 'avg_execution_time_ms': round(avg_execution_time, 2),
- 'handlers_registered': len(self._registry.get_handlers(event_type))
- }
-
- return stats
-
- def clear_history(self, event_type: Optional[str] = None) -> None:
- """
- Clear event history.
-
- Args:
- event_type: Clear only specific event type history, or all if None
- """
- with self._lock:
- if event_type:
- if event_type in self._history_by_type:
- self._history_by_type[event_type].clear()
- pprint(f"Cleared history for event type {event_type}")
- else:
- self._history.clear()
- self._history_by_type.clear()
- pprint("Cleared all event history")
-
- def get_registered_handlers(self) -> Dict[str, List[str]]:
- """
- Get information about all registered event handlers.
-
- Returns:
- Dictionary mapping event types to handler names
- """
- result = {}
- all_handlers = self._registry.get_all_handlers()
-
- for event_type, handlers in all_handlers.items():
- handler_info = []
- for handler in handlers:
- metadata = self._registry.get_handler_metadata(handler)
- handler_name = getattr(handler, '__qualname__', handler.__name__)
- priority = metadata.get('priority', 0)
- async_handler = metadata.get('async_handler', False)
-
- info = f"{handler_name} (priority: {priority}"
- if async_handler:
- info += ", async"
- info += ")"
-
- handler_info.append(info)
-
- result[event_type] = handler_info
-
- return result
-
- def shutdown(self) -> None:
- """Shutdown the event handler and cleanup resources."""
- pprint("Shutting down EventHandler...")
-
- # Shutdown thread pool
- self._executor.shutdown(wait=True)
-
- # Clear event listeners
- with self._lock:
- self._event_listeners.clear()
-
- pprint("EventHandler shutdown complete")
-
- def __enter__(self):
- """Context manager entry."""
- return self
-
- def __exit__(self, exc_type, exc_val, exc_tb):
- """Context manager exit."""
- self.shutdown()
|