event_handler.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562
  1. """
  2. Main event handler system for the Trixy application.
  3. This module provides the EventHandler class that manages event registration,
  4. triggering, and history tracking. It supports thread-safe operations for
  5. multi-satellite environments and provides comprehensive debugging capabilities.
  6. """
  7. import asyncio
  8. import threading
  9. import time
  10. import traceback
  11. from collections import deque, defaultdict
  12. from concurrent.futures import ThreadPoolExecutor, as_completed
  13. from dataclasses import dataclass, field
  14. from datetime import datetime, timedelta
  15. from typing import Any, Callable, Dict, List, Optional, Union, Set
  16. from enum import Enum
  17. from .event_data import TrixyEventData, EventType, EventDataFactory
  18. from .decorators import (
  19. get_event_registry, register_event_handlers, unregister_event_handlers,
  20. TrixyEventRegistry
  21. )
  22. def pprint(message: str) -> None:
  23. """
  24. Debug printing function that adapts based on mode.
  25. In production: uses proper logging, in debug: uses print.
  26. """
  27. # TODO: This should integrate with the application's logging system
  28. print(f"[EventHandler] {message}")
  29. class EventStatus(Enum):
  30. """Status of event processing."""
  31. PENDING = "pending"
  32. PROCESSING = "processing"
  33. COMPLETED = "completed"
  34. FAILED = "failed"
  35. CANCELLED = "cancelled"
  36. @dataclass
  37. class EventHistoryEntry:
  38. """Entry in the event history log."""
  39. event_id: str
  40. event_type: str
  41. event_data: TrixyEventData
  42. timestamp: datetime
  43. status: EventStatus = EventStatus.PENDING
  44. processing_start_time: Optional[datetime] = None
  45. processing_end_time: Optional[datetime] = None
  46. handlers_called: List[str] = field(default_factory=list)
  47. errors: List[str] = field(default_factory=list)
  48. execution_time_ms: float = 0.0
  49. def to_dict(self) -> Dict[str, Any]:
  50. """Convert history entry to dictionary."""
  51. return {
  52. 'event_id': self.event_id,
  53. 'event_type': self.event_type,
  54. 'event_data': self.event_data.to_dict() if self.event_data else None,
  55. 'timestamp': self.timestamp.isoformat(),
  56. 'status': self.status.value,
  57. 'processing_start_time': self.processing_start_time.isoformat() if self.processing_start_time else None,
  58. 'processing_end_time': self.processing_end_time.isoformat() if self.processing_end_time else None,
  59. 'handlers_called': self.handlers_called,
  60. 'errors': self.errors,
  61. 'execution_time_ms': self.execution_time_ms
  62. }
  63. class EventHandler:
  64. """
  65. Thread-safe event handler system for the Trixy application.
  66. This class manages event registration, triggering, and history tracking.
  67. It supports both synchronous and asynchronous event handlers and provides
  68. comprehensive debugging and monitoring capabilities.
  69. """
  70. def __init__(self, max_history_size: int = 1000, max_worker_threads: int = 10):
  71. """
  72. Initialize the event handler.
  73. Args:
  74. max_history_size: Maximum number of events to keep in history
  75. max_worker_threads: Maximum number of threads for async event processing
  76. """
  77. self._lock = threading.RLock()
  78. self._registry = get_event_registry()
  79. # Event history tracking
  80. self._history: deque[EventHistoryEntry] = deque(maxlen=max_history_size)
  81. self._history_by_type: Dict[str, deque[EventHistoryEntry]] = defaultdict(
  82. lambda: deque(maxlen=100)
  83. )
  84. # Statistics
  85. self._event_stats: Dict[str, Dict[str, int]] = defaultdict(
  86. lambda: {'triggered': 0, 'completed': 0, 'failed': 0}
  87. )
  88. # Thread pool for async operations
  89. self._executor = ThreadPoolExecutor(max_workers=max_worker_threads)
  90. # Event filtering and debugging
  91. self._enabled_events: Set[str] = set()
  92. self._disabled_events: Set[str] = set()
  93. self._debug_mode = False
  94. self._event_listeners: List[Callable] = []
  95. # Event counters for unique IDs
  96. self._event_counter = 0
  97. self._counter_lock = threading.Lock()
  98. pprint("EventHandler initialized")
  99. def _generate_event_id(self) -> str:
  100. """Generate a unique event ID."""
  101. with self._counter_lock:
  102. self._event_counter += 1
  103. return f"evt_{int(time.time())}_{self._event_counter:06d}"
  104. def enable_debug_mode(self, enabled: bool = True) -> None:
  105. """Enable or disable debug mode for verbose logging."""
  106. with self._lock:
  107. self._debug_mode = enabled
  108. pprint(f"Debug mode {'enabled' if enabled else 'disabled'}")
  109. def add_event_listener(self, listener: Callable[[EventHistoryEntry], None]) -> None:
  110. """
  111. Add an event listener that gets called for every event.
  112. Args:
  113. listener: Function that takes an EventHistoryEntry
  114. """
  115. with self._lock:
  116. if listener not in self._event_listeners:
  117. self._event_listeners.append(listener)
  118. pprint(f"Added event listener: {listener.__name__}")
  119. def remove_event_listener(self, listener: Callable) -> None:
  120. """Remove an event listener."""
  121. with self._lock:
  122. if listener in self._event_listeners:
  123. self._event_listeners.remove(listener)
  124. pprint(f"Removed event listener: {listener.__name__}")
  125. def enable_event_type(self, event_type: str) -> None:
  126. """Enable processing for a specific event type."""
  127. with self._lock:
  128. self._enabled_events.add(event_type)
  129. self._disabled_events.discard(event_type)
  130. pprint(f"Enabled event type: {event_type}")
  131. def disable_event_type(self, event_type: str) -> None:
  132. """Disable processing for a specific event type."""
  133. with self._lock:
  134. self._disabled_events.add(event_type)
  135. self._enabled_events.discard(event_type)
  136. pprint(f"Disabled event type: {event_type}")
  137. def is_event_type_enabled(self, event_type: str) -> bool:
  138. """Check if an event type is enabled for processing."""
  139. with self._lock:
  140. if self._disabled_events and event_type in self._disabled_events:
  141. return False
  142. if self._enabled_events and event_type not in self._enabled_events:
  143. return False
  144. return True
  145. def register_handler_object(self, obj: Any) -> None:
  146. """
  147. Register all @TrixyEvent decorated methods from an object.
  148. Args:
  149. obj: Object instance containing @TrixyEvent decorated methods
  150. """
  151. with self._lock:
  152. register_event_handlers(obj)
  153. pprint(f"Registered event handlers from {obj.__class__.__name__}")
  154. def unregister_handler_object(self, obj: Any) -> None:
  155. """
  156. Unregister all @TrixyEvent decorated methods from an object.
  157. Args:
  158. obj: Object instance to unregister handlers from
  159. """
  160. with self._lock:
  161. unregister_event_handlers(obj)
  162. pprint(f"Unregistered event handlers from {obj.__class__.__name__}")
  163. def trigger_event(
  164. self,
  165. event_type: Union[str, EventType],
  166. event_data: Optional[TrixyEventData] = None,
  167. **kwargs
  168. ) -> str:
  169. """
  170. Trigger an event synchronously.
  171. Args:
  172. event_type: Type of event to trigger
  173. event_data: Event data object (optional, will be created if not provided)
  174. **kwargs: Additional arguments for event data creation
  175. Returns:
  176. str: Event ID for tracking
  177. """
  178. # Normalize event type
  179. if isinstance(event_type, EventType):
  180. event_type_str = event_type.value
  181. else:
  182. event_type_str = str(event_type)
  183. # Check if event type is enabled
  184. if not self.is_event_type_enabled(event_type_str):
  185. if self._debug_mode:
  186. pprint(f"Event type {event_type_str} is disabled, skipping")
  187. return ""
  188. # Create event data if not provided
  189. if event_data is None:
  190. try:
  191. event_data = EventDataFactory.create_event_data(event_type_str, **kwargs)
  192. except ValueError as e:
  193. if self._debug_mode:
  194. pprint(f"Failed to create event data for {event_type_str}: {e}")
  195. # Create a generic event data object
  196. from .event_data import TrixyEventData
  197. event_data = TrixyEventData(**kwargs)
  198. # Generate event ID and create history entry
  199. event_id = self._generate_event_id()
  200. history_entry = EventHistoryEntry(
  201. event_id=event_id,
  202. event_type=event_type_str,
  203. event_data=event_data,
  204. timestamp=datetime.now()
  205. )
  206. try:
  207. self._process_event_sync(history_entry)
  208. except Exception as e:
  209. pprint(f"Error processing event {event_type_str}: {e}")
  210. if self._debug_mode:
  211. pprint(f"Traceback: {traceback.format_exc()}")
  212. history_entry.status = EventStatus.FAILED
  213. history_entry.errors.append(str(e))
  214. # Add to history and notify listeners
  215. with self._lock:
  216. self._history.append(history_entry)
  217. self._history_by_type[event_type_str].append(history_entry)
  218. self._event_stats[event_type_str]['triggered'] += 1
  219. if history_entry.status == EventStatus.COMPLETED:
  220. self._event_stats[event_type_str]['completed'] += 1
  221. elif history_entry.status == EventStatus.FAILED:
  222. self._event_stats[event_type_str]['failed'] += 1
  223. # Notify event listeners
  224. for listener in self._event_listeners:
  225. try:
  226. listener(history_entry)
  227. except Exception as e:
  228. pprint(f"Error in event listener {listener.__name__}: {e}")
  229. if self._debug_mode:
  230. pprint(f"Triggered event {event_type_str} with ID {event_id}")
  231. return event_id
  232. def trigger_event_async(
  233. self,
  234. event_type: Union[str, EventType],
  235. event_data: Optional[TrixyEventData] = None,
  236. **kwargs
  237. ) -> str:
  238. """
  239. Trigger an event asynchronously.
  240. Args:
  241. event_type: Type of event to trigger
  242. event_data: Event data object (optional, will be created if not provided)
  243. **kwargs: Additional arguments for event data creation
  244. Returns:
  245. str: Event ID for tracking
  246. """
  247. # Similar setup as sync version
  248. if isinstance(event_type, EventType):
  249. event_type_str = event_type.value
  250. else:
  251. event_type_str = str(event_type)
  252. if not self.is_event_type_enabled(event_type_str):
  253. if self._debug_mode:
  254. pprint(f"Event type {event_type_str} is disabled, skipping")
  255. return ""
  256. if event_data is None:
  257. try:
  258. event_data = EventDataFactory.create_event_data(event_type_str, **kwargs)
  259. except ValueError:
  260. from .event_data import TrixyEventData
  261. event_data = TrixyEventData(**kwargs)
  262. event_id = self._generate_event_id()
  263. history_entry = EventHistoryEntry(
  264. event_id=event_id,
  265. event_type=event_type_str,
  266. event_data=event_data,
  267. timestamp=datetime.now()
  268. )
  269. # Submit to thread pool
  270. future = self._executor.submit(self._process_event_sync, history_entry)
  271. # Add a callback to handle completion
  272. def on_completion(fut):
  273. try:
  274. fut.result() # This will raise any exception that occurred
  275. except Exception as e:
  276. pprint(f"Async event {event_type_str} failed: {e}")
  277. history_entry.status = EventStatus.FAILED
  278. history_entry.errors.append(str(e))
  279. # Update history and stats
  280. with self._lock:
  281. self._history.append(history_entry)
  282. self._history_by_type[event_type_str].append(history_entry)
  283. self._event_stats[event_type_str]['triggered'] += 1
  284. if history_entry.status == EventStatus.COMPLETED:
  285. self._event_stats[event_type_str]['completed'] += 1
  286. elif history_entry.status == EventStatus.FAILED:
  287. self._event_stats[event_type_str]['failed'] += 1
  288. # Notify event listeners
  289. for listener in self._event_listeners:
  290. try:
  291. listener(history_entry)
  292. except Exception as e:
  293. pprint(f"Error in event listener {listener.__name__}: {e}")
  294. future.add_done_callback(on_completion)
  295. if self._debug_mode:
  296. pprint(f"Triggered async event {event_type_str} with ID {event_id}")
  297. return event_id
  298. def _process_event_sync(self, history_entry: EventHistoryEntry) -> None:
  299. """
  300. Process an event synchronously.
  301. Args:
  302. history_entry: Event history entry to process
  303. """
  304. history_entry.status = EventStatus.PROCESSING
  305. history_entry.processing_start_time = datetime.now()
  306. handlers = self._registry.get_handlers(history_entry.event_type)
  307. if not handlers:
  308. if self._debug_mode:
  309. pprint(f"No handlers registered for event {history_entry.event_type}")
  310. history_entry.status = EventStatus.COMPLETED
  311. history_entry.processing_end_time = datetime.now()
  312. return
  313. for handler in handlers:
  314. try:
  315. handler_name = getattr(handler, '__qualname__', handler.__name__)
  316. history_entry.handlers_called.append(handler_name)
  317. if self._debug_mode:
  318. pprint(f"Calling handler {handler_name} for event {history_entry.event_type}")
  319. # Call the handler
  320. handler(history_entry.event_type, history_entry.event_data)
  321. except Exception as e:
  322. error_msg = f"Error in handler {handler_name}: {e}"
  323. history_entry.errors.append(error_msg)
  324. pprint(error_msg)
  325. if self._debug_mode:
  326. pprint(f"Traceback: {traceback.format_exc()}")
  327. history_entry.processing_end_time = datetime.now()
  328. if history_entry.processing_start_time:
  329. history_entry.execution_time_ms = (
  330. history_entry.processing_end_time - history_entry.processing_start_time
  331. ).total_seconds() * 1000
  332. history_entry.status = EventStatus.FAILED if history_entry.errors else EventStatus.COMPLETED
  333. def wait_for_event(self, event_id: str, timeout: float = 30.0) -> Optional[EventHistoryEntry]:
  334. """
  335. Wait for an event to complete processing.
  336. Args:
  337. event_id: Event ID to wait for
  338. timeout: Maximum time to wait in seconds
  339. Returns:
  340. EventHistoryEntry if found and completed, None if timeout or not found
  341. """
  342. start_time = time.time()
  343. while time.time() - start_time < timeout:
  344. with self._lock:
  345. for entry in self._history:
  346. if entry.event_id == event_id:
  347. if entry.status in [EventStatus.COMPLETED, EventStatus.FAILED, EventStatus.CANCELLED]:
  348. return entry
  349. break
  350. time.sleep(0.1) # Small delay to avoid busy waiting
  351. return None
  352. def get_event_history(
  353. self,
  354. event_type: Optional[str] = None,
  355. limit: Optional[int] = None,
  356. since: Optional[datetime] = None
  357. ) -> List[EventHistoryEntry]:
  358. """
  359. Get event history with optional filtering.
  360. Args:
  361. event_type: Filter by event type
  362. limit: Maximum number of entries to return
  363. since: Only return events after this timestamp
  364. Returns:
  365. List of event history entries
  366. """
  367. with self._lock:
  368. if event_type:
  369. entries = list(self._history_by_type.get(event_type, []))
  370. else:
  371. entries = list(self._history)
  372. # Filter by timestamp
  373. if since:
  374. entries = [e for e in entries if e.timestamp >= since]
  375. # Sort by timestamp (newest first)
  376. entries.sort(key=lambda e: e.timestamp, reverse=True)
  377. # Apply limit
  378. if limit:
  379. entries = entries[:limit]
  380. return entries
  381. def get_event_statistics(self) -> Dict[str, Dict[str, Any]]:
  382. """
  383. Get comprehensive event statistics.
  384. Returns:
  385. Dictionary containing event statistics
  386. """
  387. with self._lock:
  388. stats = {}
  389. for event_type, type_stats in self._event_stats.items():
  390. success_rate = 0.0
  391. if type_stats['triggered'] > 0:
  392. success_rate = type_stats['completed'] / type_stats['triggered'] * 100
  393. # Calculate average execution time
  394. recent_entries = list(self._history_by_type[event_type])[-10:] # Last 10 events
  395. avg_execution_time = 0.0
  396. if recent_entries:
  397. total_time = sum(e.execution_time_ms for e in recent_entries if e.execution_time_ms > 0)
  398. avg_execution_time = total_time / len(recent_entries)
  399. stats[event_type] = {
  400. 'triggered': type_stats['triggered'],
  401. 'completed': type_stats['completed'],
  402. 'failed': type_stats['failed'],
  403. 'success_rate_percent': round(success_rate, 2),
  404. 'avg_execution_time_ms': round(avg_execution_time, 2),
  405. 'handlers_registered': len(self._registry.get_handlers(event_type))
  406. }
  407. return stats
  408. def clear_history(self, event_type: Optional[str] = None) -> None:
  409. """
  410. Clear event history.
  411. Args:
  412. event_type: Clear only specific event type history, or all if None
  413. """
  414. with self._lock:
  415. if event_type:
  416. if event_type in self._history_by_type:
  417. self._history_by_type[event_type].clear()
  418. pprint(f"Cleared history for event type {event_type}")
  419. else:
  420. self._history.clear()
  421. self._history_by_type.clear()
  422. pprint("Cleared all event history")
  423. def get_registered_handlers(self) -> Dict[str, List[str]]:
  424. """
  425. Get information about all registered event handlers.
  426. Returns:
  427. Dictionary mapping event types to handler names
  428. """
  429. result = {}
  430. all_handlers = self._registry.get_all_handlers()
  431. for event_type, handlers in all_handlers.items():
  432. handler_info = []
  433. for handler in handlers:
  434. metadata = self._registry.get_handler_metadata(handler)
  435. handler_name = getattr(handler, '__qualname__', handler.__name__)
  436. priority = metadata.get('priority', 0)
  437. async_handler = metadata.get('async_handler', False)
  438. info = f"{handler_name} (priority: {priority}"
  439. if async_handler:
  440. info += ", async"
  441. info += ")"
  442. handler_info.append(info)
  443. result[event_type] = handler_info
  444. return result
  445. def shutdown(self) -> None:
  446. """Shutdown the event handler and cleanup resources."""
  447. pprint("Shutting down EventHandler...")
  448. # Shutdown thread pool
  449. self._executor.shutdown(wait=True)
  450. # Clear event listeners
  451. with self._lock:
  452. self._event_listeners.clear()
  453. pprint("EventHandler shutdown complete")
  454. def __enter__(self):
  455. """Context manager entry."""
  456. return self
  457. def __exit__(self, exc_type, exc_val, exc_tb):
  458. """Context manager exit."""
  459. self.shutdown()