conversation_manager.py 61 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594
  1. """
  2. Main Conversation Manager Implementation
  3. This module implements the ConversationManager class that serves as the central
  4. management system for all conversations in the Trixy voice assistant system.
  5. It coordinates conversation sessions, state management, context preservation,
  6. and multi-turn conversation workflows as specified in CLAUDE.md.
  7. Key Features:
  8. - Session management for multi-turn conversations
  9. - Conversation ID tracking across events
  10. - Context preservation and state management
  11. - Plugin integration for multi-turn workflows
  12. - Thread-safe operations for multi-satellite support
  13. - Conversation timeout handling and cleanup
  14. - Comprehensive logging and analytics
  15. - Integration with event system and application container
  16. The ConversationManager implements the conversation workflow specified in CLAUDE.md:
  17. 1. Wakeword detection starts conversation session
  18. 2. Raw audio input with conversation_id
  19. 3. Text processing with same conversation_id
  20. 4. Intent handling with conversation tracking
  21. 5. Multi-turn conversations where plugins ask questions and wait for responses
  22. Example workflow:
  23. - User: "Trixy, Please order a pizza"
  24. - Trixy: "Ok, what kind of pizza do you want?"
  25. - User: "A spicy one"
  26. - Trixy: "Ok, i will send an order to your local pizza service"
  27. """
  28. import asyncio
  29. import json
  30. import threading
  31. import time
  32. import uuid
  33. import weakref
  34. from datetime import datetime, timezone, timedelta
  35. from dataclasses import dataclass, field
  36. from pathlib import Path
  37. from typing import Any, Dict, List, Optional, Set, Callable, Tuple, Union
  38. from collections import defaultdict, OrderedDict
  39. from concurrent.futures import ThreadPoolExecutor
  40. from .conversation_session import (
  41. ConversationSession, ConversationSessionState, ConversationTurn,
  42. ConversationRole, create_conversation_session
  43. )
  44. from .conversation_context import (
  45. ConversationContext, ContextScope, ContextRetentionPolicy,
  46. create_conversation_context
  47. )
  48. from .conversation_state import (
  49. ConversationState, ConversationStateManager, ConversationPhase,
  50. create_state_manager
  51. )
  52. from ..events.event_data import SatelliteInfo, SpeakerInfo
  53. def pprint(message: str) -> None:
  54. """Conversation manager logging function."""
  55. print(f"[CONVERSATION_MANAGER] {message}")
  56. @dataclass
  57. class ConversationStats:
  58. """Statistics for conversation management."""
  59. total_conversations: int = 0
  60. active_conversations: int = 0
  61. completed_conversations: int = 0
  62. timed_out_conversations: int = 0
  63. error_conversations: int = 0
  64. total_turns: int = 0
  65. average_turns_per_conversation: float = 0.0
  66. average_conversation_duration: float = 0.0
  67. plugin_interactions: int = 0
  68. context_entries: int = 0
  69. uptime_seconds: float = 0.0
  70. last_updated: Optional[datetime] = None
  71. def to_dict(self) -> Dict[str, Any]:
  72. """Convert to dictionary representation."""
  73. return {
  74. "total_conversations": self.total_conversations,
  75. "active_conversations": self.active_conversations,
  76. "completed_conversations": self.completed_conversations,
  77. "timed_out_conversations": self.timed_out_conversations,
  78. "error_conversations": self.error_conversations,
  79. "total_turns": self.total_turns,
  80. "average_turns_per_conversation": self.average_turns_per_conversation,
  81. "average_conversation_duration": self.average_conversation_duration,
  82. "plugin_interactions": self.plugin_interactions,
  83. "context_entries": self.context_entries,
  84. "uptime_seconds": self.uptime_seconds,
  85. "last_updated": self.last_updated.isoformat() if self.last_updated else None,
  86. }
  87. class ConversationManagerError(Exception):
  88. """Base exception for conversation manager errors."""
  89. pass
  90. class ConversationNotFoundError(ConversationManagerError):
  91. """Raised when a requested conversation is not found."""
  92. pass
  93. class ConversationTimeoutError(ConversationManagerError):
  94. """Raised when a conversation operation times out."""
  95. pass
  96. class ConversationStateError(ConversationManagerError):
  97. """Raised when a conversation is in an invalid state for the requested operation."""
  98. pass
  99. class ConversationManager:
  100. """
  101. Central management system for all conversations in the Trixy voice assistant.
  102. This class implements the conversation management system as specified in CLAUDE.md,
  103. providing comprehensive session management, context preservation, and multi-turn
  104. conversation support for the Trixy voice assistant system.
  105. Key Features:
  106. - Session management for multi-turn conversations
  107. - Conversation ID tracking across events
  108. - Context preservation across turns and sessions
  109. - Plugin integration for multi-turn workflows
  110. - Thread-safe operations for multi-satellite support
  111. - Automatic timeout handling and cleanup
  112. - Comprehensive analytics and monitoring
  113. - Integration with event system and application container
  114. The manager supports the conversation workflow specified in CLAUDE.md:
  115. 1. Wakeword detection starts conversation session
  116. 2. Raw audio input with conversation_id
  117. 3. Text processing with same conversation_id
  118. 4. Intent handling with conversation tracking
  119. 5. Multi-turn conversations where plugins ask questions and wait for responses
  120. Usage:
  121. # Create conversation manager
  122. manager = ConversationManager(application)
  123. # Start a conversation
  124. session = manager.start_conversation(
  125. satellite_info=satellite_info,
  126. speaker_info=speaker_info,
  127. trigger_event="wakeword_received"
  128. )
  129. # Process conversation workflow
  130. manager.process_audio_input(session.conversation_id, audio_data)
  131. manager.process_text_input(session.conversation_id, "order pizza")
  132. manager.process_intent(session.conversation_id, intent_data)
  133. # Plugin interaction
  134. manager.ask_question(session.conversation_id, "What kind of pizza?")
  135. manager.process_text_input(session.conversation_id, "spicy one")
  136. # Complete conversation
  137. manager.complete_conversation(session.conversation_id, "Order placed!")
  138. """
  139. def __init__(
  140. self,
  141. application,
  142. max_concurrent_conversations: int = 100,
  143. default_conversation_timeout: float = 300.0, # 5 minutes
  144. plugin_question_timeout: float = 30.0, # 30 seconds
  145. cleanup_interval_seconds: float = 60.0, # 1 minute
  146. max_conversation_history: int = 1000,
  147. enable_context_persistence: bool = False,
  148. context_storage_path: Optional[str] = None,
  149. enable_analytics: bool = True
  150. ):
  151. """
  152. Initialize the conversation manager.
  153. Args:
  154. application: Application container reference
  155. max_concurrent_conversations: Maximum concurrent conversations
  156. default_conversation_timeout: Default timeout for conversations
  157. plugin_question_timeout: Timeout for plugin questions
  158. cleanup_interval_seconds: Interval for cleanup tasks
  159. max_conversation_history: Maximum conversations to keep in history
  160. enable_context_persistence: Enable context persistence to disk
  161. context_storage_path: Path for storing context data
  162. enable_analytics: Enable analytics collection
  163. """
  164. # Core dependencies
  165. self._application = application
  166. self._application_ref = weakref.ref(application) if application else None
  167. # Configuration
  168. self._max_concurrent_conversations = max_concurrent_conversations
  169. self._default_timeout = default_conversation_timeout
  170. self._plugin_timeout = plugin_question_timeout
  171. self._cleanup_interval = cleanup_interval_seconds
  172. self._max_history = max_conversation_history
  173. self._enable_context_persistence = enable_context_persistence
  174. self._context_storage_path = Path(context_storage_path or "./conversation_context")
  175. self._enable_analytics = enable_analytics
  176. # Active conversation storage
  177. self._active_conversations: Dict[str, ConversationSession] = OrderedDict()
  178. self._conversation_contexts: Dict[str, ConversationContext] = {}
  179. self._conversation_states: Dict[str, ConversationStateManager] = {}
  180. # Conversation history
  181. self._completed_conversations: Dict[str, ConversationSession] = OrderedDict()
  182. # Context management across conversations
  183. self._global_context: ConversationContext = create_conversation_context(
  184. conversation_id="global",
  185. enable_persistence=enable_context_persistence,
  186. persistence_path=self._context_storage_path / "global"
  187. )
  188. # Speaker and satellite context tracking
  189. self._speaker_contexts: Dict[str, ConversationContext] = {}
  190. self._satellite_contexts: Dict[str, ConversationContext] = {}
  191. self._room_contexts: Dict[str, ConversationContext] = {}
  192. # Plugin interaction tracking
  193. self._plugin_questions: Dict[str, Dict[str, Any]] = {} # conversation_id -> plugin data
  194. self._pending_callbacks: Dict[str, Callable] = {}
  195. # Thread safety
  196. self._lock = threading.RLock()
  197. # Background tasks
  198. self._cleanup_timer: Optional[threading.Timer] = None
  199. self._executor = ThreadPoolExecutor(max_workers=4, thread_name_prefix="ConversationManager")
  200. # Statistics and monitoring
  201. self._stats = ConversationStats()
  202. self._start_time = datetime.now(timezone.utc)
  203. # Event integration
  204. self._event_handler = None
  205. self._satellite_manager = None
  206. # Create storage directories
  207. if self._enable_context_persistence:
  208. self._context_storage_path.mkdir(parents=True, exist_ok=True)
  209. pprint(f"ConversationManager initialized (max_concurrent: {max_concurrent_conversations})")
  210. # Start background cleanup
  211. self._start_cleanup_timer()
  212. # Load persisted contexts
  213. if self._enable_context_persistence:
  214. self._load_persisted_contexts()
  215. # Core Conversation Management
  216. def start_conversation(
  217. self,
  218. satellite_info: SatelliteInfo,
  219. speaker_info: Optional[SpeakerInfo] = None,
  220. trigger_event: str = "wakeword_received",
  221. conversation_id: Optional[str] = None,
  222. timeout_seconds: Optional[float] = None,
  223. metadata: Optional[Dict[str, Any]] = None
  224. ) -> ConversationSession:
  225. """
  226. Start a new conversation session.
  227. Args:
  228. satellite_info: Information about the satellite device
  229. speaker_info: Information about the detected speaker
  230. trigger_event: Event that triggered the conversation
  231. conversation_id: Optional specific conversation ID
  232. timeout_seconds: Conversation timeout (uses default if None)
  233. metadata: Additional metadata for the conversation
  234. Returns:
  235. ConversationSession: The created conversation session
  236. Raises:
  237. ConversationManagerError: If conversation cannot be started
  238. """
  239. with self._lock:
  240. # Check capacity
  241. if len(self._active_conversations) >= self._max_concurrent_conversations:
  242. # Try to clean up completed conversations
  243. self._cleanup_completed_conversations()
  244. if len(self._active_conversations) >= self._max_concurrent_conversations:
  245. raise ConversationManagerError(
  246. f"Maximum concurrent conversations ({self._max_concurrent_conversations}) reached"
  247. )
  248. # Generate conversation ID if not provided
  249. if conversation_id is None:
  250. conversation_id = self._generate_conversation_id()
  251. # Check for duplicate conversation ID
  252. if conversation_id in self._active_conversations:
  253. raise ConversationManagerError(f"Conversation {conversation_id} already exists")
  254. # Create conversation session
  255. session = create_conversation_session(
  256. conversation_id=conversation_id,
  257. satellite_info=satellite_info,
  258. speaker_info=speaker_info,
  259. trigger_event=trigger_event,
  260. timeout_seconds=timeout_seconds or self._default_timeout,
  261. application_ref=self._application_ref
  262. )
  263. # Create conversation context
  264. context = create_conversation_context(
  265. conversation_id=conversation_id,
  266. speaker_id=speaker_info.speaker_id if speaker_info else None,
  267. satellite_id=satellite_info.satellite_id,
  268. room_id=satellite_info.room_id,
  269. enable_persistence=self._enable_context_persistence,
  270. persistence_path=self._context_storage_path / conversation_id
  271. )
  272. # Create state manager
  273. state_manager = create_state_manager(
  274. conversation_id=conversation_id,
  275. initial_state=ConversationState.INITIALIZING,
  276. application_ref=self._application_ref
  277. )
  278. # Add metadata to session
  279. if metadata:
  280. session._metadata.update(metadata)
  281. # Store in active conversations
  282. self._active_conversations[conversation_id] = session
  283. self._conversation_contexts[conversation_id] = context
  284. self._conversation_states[conversation_id] = state_manager
  285. # Set up speaker context if needed
  286. if speaker_info and speaker_info.speaker_id not in self._speaker_contexts:
  287. self._speaker_contexts[speaker_info.speaker_id] = create_conversation_context(
  288. conversation_id=f"speaker_{speaker_info.speaker_id}",
  289. speaker_id=speaker_info.speaker_id,
  290. enable_persistence=self._enable_context_persistence,
  291. persistence_path=self._context_storage_path / "speakers" / speaker_info.speaker_id
  292. )
  293. # Set up satellite context if needed
  294. if satellite_info.satellite_id not in self._satellite_contexts:
  295. self._satellite_contexts[satellite_info.satellite_id] = create_conversation_context(
  296. conversation_id=f"satellite_{satellite_info.satellite_id}",
  297. satellite_id=satellite_info.satellite_id,
  298. enable_persistence=self._enable_context_persistence,
  299. persistence_path=self._context_storage_path / "satellites" / satellite_info.satellite_id
  300. )
  301. # Set up room context if needed
  302. if satellite_info.room_id not in self._room_contexts:
  303. self._room_contexts[satellite_info.room_id] = create_conversation_context(
  304. conversation_id=f"room_{satellite_info.room_id}",
  305. room_id=satellite_info.room_id,
  306. enable_persistence=self._enable_context_persistence,
  307. persistence_path=self._context_storage_path / "rooms" / satellite_info.room_id
  308. )
  309. # Transition to ready state
  310. state_manager.transition_to(
  311. ConversationState.READY,
  312. trigger="conversation_started",
  313. source="conversation_manager"
  314. )
  315. # Update statistics
  316. self._update_stats()
  317. # Trigger conversation started event
  318. self._trigger_conversation_event("conversation_started", session)
  319. pprint(f"Started conversation {conversation_id} for {satellite_info.alias}")
  320. return session
  321. def get_conversation(self, conversation_id: str) -> Optional[ConversationSession]:
  322. """
  323. Get an active conversation session.
  324. Args:
  325. conversation_id: Conversation identifier
  326. Returns:
  327. ConversationSession or None if not found
  328. """
  329. with self._lock:
  330. return self._active_conversations.get(conversation_id)
  331. def get_conversation_context(self, conversation_id: str) -> Optional[ConversationContext]:
  332. """
  333. Get conversation context.
  334. Args:
  335. conversation_id: Conversation identifier
  336. Returns:
  337. ConversationContext or None if not found
  338. """
  339. with self._lock:
  340. return self._conversation_contexts.get(conversation_id)
  341. def get_conversation_state(self, conversation_id: str) -> Optional[ConversationStateManager]:
  342. """
  343. Get conversation state manager.
  344. Args:
  345. conversation_id: Conversation identifier
  346. Returns:
  347. ConversationStateManager or None if not found
  348. """
  349. with self._lock:
  350. return self._conversation_states.get(conversation_id)
  351. def end_conversation(
  352. self,
  353. conversation_id: str,
  354. reason: str = "completed",
  355. final_response: Optional[str] = None
  356. ) -> bool:
  357. """
  358. End a conversation session.
  359. Args:
  360. conversation_id: Conversation identifier
  361. reason: Reason for ending the conversation
  362. final_response: Optional final response to the user
  363. Returns:
  364. True if conversation was ended successfully
  365. """
  366. with self._lock:
  367. session = self._active_conversations.get(conversation_id)
  368. if not session:
  369. return False
  370. # Complete the session
  371. if reason == "completed":
  372. session.complete(final_response)
  373. elif reason == "timeout":
  374. session.timeout()
  375. elif reason == "cancelled":
  376. session.cancel()
  377. else:
  378. session.error(f"Ended: {reason}")
  379. # Get context and state manager
  380. context = self._conversation_contexts.get(conversation_id)
  381. state_manager = self._conversation_states.get(conversation_id)
  382. # End session context
  383. if context:
  384. context.end_session()
  385. # Transition to appropriate terminal state
  386. if state_manager and not state_manager.is_terminal:
  387. if reason == "completed":
  388. state_manager.transition_to(ConversationState.COMPLETED)
  389. elif reason == "timeout":
  390. state_manager.transition_to(ConversationState.TIMED_OUT)
  391. elif reason == "cancelled":
  392. state_manager.transition_to(ConversationState.CANCELLED)
  393. else:
  394. state_manager.transition_to(ConversationState.ERROR)
  395. # Move to completed conversations
  396. self._completed_conversations[conversation_id] = session
  397. # Clean up active conversation
  398. del self._active_conversations[conversation_id]
  399. if conversation_id in self._conversation_contexts:
  400. del self._conversation_contexts[conversation_id]
  401. if conversation_id in self._conversation_states:
  402. del self._conversation_states[conversation_id]
  403. if conversation_id in self._plugin_questions:
  404. del self._plugin_questions[conversation_id]
  405. if conversation_id in self._pending_callbacks:
  406. del self._pending_callbacks[conversation_id]
  407. # Trim completed conversations if needed
  408. if len(self._completed_conversations) > self._max_history:
  409. # Remove oldest completed conversations
  410. oldest_key = next(iter(self._completed_conversations))
  411. del self._completed_conversations[oldest_key]
  412. # Update statistics
  413. self._update_stats()
  414. # Trigger conversation ended event
  415. self._trigger_conversation_event("conversation_ended", session, {"reason": reason})
  416. pprint(f"Ended conversation {conversation_id}: {reason}")
  417. return True
  418. # Audio and Input Processing
  419. def process_audio_input(
  420. self,
  421. conversation_id: str,
  422. audio_data: bytes,
  423. speaker_info: Optional[SpeakerInfo] = None,
  424. metadata: Optional[Dict[str, Any]] = None
  425. ) -> Optional[ConversationTurn]:
  426. """
  427. Process raw audio input for a conversation.
  428. Args:
  429. conversation_id: Conversation identifier
  430. audio_data: Raw audio data
  431. speaker_info: Speaker information
  432. metadata: Additional metadata
  433. Returns:
  434. ConversationTurn or None if conversation not found
  435. """
  436. session = self.get_conversation(conversation_id)
  437. context = self.get_conversation_context(conversation_id)
  438. state_manager = self.get_conversation_state(conversation_id)
  439. if not session or not context or not state_manager:
  440. pprint(f"Conversation {conversation_id} not found for audio input")
  441. return None
  442. # Transition to listening state
  443. state_manager.transition_to(
  444. ConversationState.LISTENING,
  445. trigger="audio_input_received",
  446. source="conversation_manager"
  447. )
  448. # Add audio turn
  449. turn = session.add_turn(
  450. role=ConversationRole.USER,
  451. input_type="audio",
  452. input_audio_data=audio_data,
  453. speaker_info=speaker_info,
  454. metadata=metadata or {}
  455. )
  456. # Update context with audio metadata
  457. context.set("last_audio_input", {
  458. "turn_id": turn.turn_id,
  459. "audio_length": len(audio_data),
  460. "timestamp": turn.timestamp.isoformat()
  461. }, scope=ContextScope.SESSION)
  462. # Trigger audio input event
  463. self._trigger_audio_input_event(conversation_id, audio_data, speaker_info, session.satellite_info)
  464. pprint(f"Processed audio input for conversation {conversation_id}")
  465. return turn
  466. def process_text_input(
  467. self,
  468. conversation_id: str,
  469. text: str,
  470. confidence: float = 1.0,
  471. speaker_info: Optional[SpeakerInfo] = None,
  472. metadata: Optional[Dict[str, Any]] = None
  473. ) -> Optional[ConversationTurn]:
  474. """
  475. Process text input for a conversation.
  476. Args:
  477. conversation_id: Conversation identifier
  478. text: Input text
  479. confidence: Confidence score
  480. speaker_info: Speaker information
  481. metadata: Additional metadata
  482. Returns:
  483. ConversationTurn or None if conversation not found
  484. """
  485. session = self.get_conversation(conversation_id)
  486. context = self.get_conversation_context(conversation_id)
  487. state_manager = self.get_conversation_state(conversation_id)
  488. if not session or not context or not state_manager:
  489. pprint(f"Conversation {conversation_id} not found for text input")
  490. return None
  491. # Check if this is an answer to a plugin question
  492. if conversation_id in self._plugin_questions:
  493. return self._handle_plugin_answer(conversation_id, text, confidence, speaker_info)
  494. # Transition to processing state
  495. state_manager.transition_to(
  496. ConversationState.TRANSCRIBING,
  497. trigger="text_input_received",
  498. source="conversation_manager"
  499. )
  500. # Add text turn
  501. turn = session.add_turn(
  502. role=ConversationRole.USER,
  503. input_type="text",
  504. input_text=text,
  505. input_confidence=confidence,
  506. speaker_info=speaker_info,
  507. metadata=metadata or {}
  508. )
  509. # Update context with text
  510. context.set("last_text_input", {
  511. "turn_id": turn.turn_id,
  512. "text": text,
  513. "confidence": confidence,
  514. "timestamp": turn.timestamp.isoformat()
  515. }, scope=ContextScope.SESSION)
  516. # Store in speaker context if available
  517. if speaker_info:
  518. speaker_context = self._speaker_contexts.get(speaker_info.speaker_id)
  519. if speaker_context:
  520. speaker_context.set("last_utterance", text, scope=ContextScope.SPEAKER)
  521. # Trigger text received event
  522. self._trigger_text_received_event(conversation_id, text, confidence, speaker_info, session.satellite_info)
  523. pprint(f"Processed text input for conversation {conversation_id}: {text}")
  524. return turn
  525. def process_intent(
  526. self,
  527. conversation_id: str,
  528. intent: str,
  529. entities: Optional[Dict[str, Any]] = None,
  530. confidence: float = 1.0,
  531. metadata: Optional[Dict[str, Any]] = None
  532. ) -> bool:
  533. """
  534. Process intent for a conversation.
  535. Args:
  536. conversation_id: Conversation identifier
  537. intent: Detected intent
  538. entities: Extracted entities
  539. confidence: Intent confidence
  540. metadata: Additional metadata
  541. Returns:
  542. True if intent was processed successfully
  543. """
  544. session = self.get_conversation(conversation_id)
  545. context = self.get_conversation_context(conversation_id)
  546. state_manager = self.get_conversation_state(conversation_id)
  547. if not session or not context or not state_manager:
  548. pprint(f"Conversation {conversation_id} not found for intent processing")
  549. return False
  550. # Transition to intent processing state
  551. state_manager.transition_to(
  552. ConversationState.PROCESSING_INTENT,
  553. trigger="intent_received",
  554. source="conversation_manager"
  555. )
  556. # Update current turn with intent information
  557. current_turn = session.get_current_turn()
  558. if current_turn:
  559. session.update_turn_intent(
  560. current_turn.turn_id,
  561. intent=intent,
  562. entities=entities or {},
  563. confidence=confidence
  564. )
  565. # Update context with intent
  566. context.set("current_intent", {
  567. "intent": intent,
  568. "entities": entities or {},
  569. "confidence": confidence,
  570. "timestamp": datetime.now(timezone.utc).isoformat()
  571. }, scope=ContextScope.SESSION)
  572. # Add to conversation history in context
  573. conversation_history = context.get("conversation_history", [], scope=ContextScope.SESSION)
  574. conversation_history.append({
  575. "type": "intent",
  576. "intent": intent,
  577. "entities": entities or {},
  578. "timestamp": datetime.now(timezone.utc).isoformat()
  579. })
  580. context.set("conversation_history", conversation_history, scope=ContextScope.SESSION)
  581. # Trigger intent received event
  582. self._trigger_intent_received_event(
  583. conversation_id, intent, entities or {}, confidence,
  584. current_turn.input_text if current_turn else "",
  585. session.speaker_info, session.satellite_info
  586. )
  587. pprint(f"Processed intent for conversation {conversation_id}: {intent}")
  588. return True
  589. # Response Generation
  590. def generate_response(
  591. self,
  592. conversation_id: str,
  593. response_text: str,
  594. voice_settings: Optional[Dict[str, Any]] = None,
  595. plugin_name: Optional[str] = None,
  596. metadata: Optional[Dict[str, Any]] = None
  597. ) -> bool:
  598. """
  599. Generate a response for a conversation.
  600. Args:
  601. conversation_id: Conversation identifier
  602. response_text: Response text
  603. voice_settings: Voice synthesis settings
  604. plugin_name: Plugin that generated the response
  605. metadata: Additional metadata
  606. Returns:
  607. True if response was generated successfully
  608. """
  609. session = self.get_conversation(conversation_id)
  610. context = self.get_conversation_context(conversation_id)
  611. state_manager = self.get_conversation_state(conversation_id)
  612. if not session or not context or not state_manager:
  613. pprint(f"Conversation {conversation_id} not found for response generation")
  614. return False
  615. # Transition to response generation state
  616. state_manager.transition_to(
  617. ConversationState.GENERATING_RESPONSE,
  618. trigger="generate_response",
  619. source=plugin_name or "conversation_manager"
  620. )
  621. # Update current turn or create assistant turn
  622. current_turn = session.get_current_turn()
  623. if current_turn and current_turn.role == ConversationRole.USER:
  624. # Update existing user turn with response
  625. session.update_turn_response(
  626. current_turn.turn_id,
  627. response_text=response_text,
  628. response_voice_settings=voice_settings or {},
  629. plugin_name=plugin_name
  630. )
  631. else:
  632. # Create new assistant turn
  633. turn = session.add_turn(
  634. role=ConversationRole.ASSISTANT,
  635. input_type="response",
  636. input_text=response_text,
  637. metadata=metadata or {}
  638. )
  639. session.update_turn_response(
  640. turn.turn_id,
  641. response_text=response_text,
  642. response_voice_settings=voice_settings or {},
  643. plugin_name=plugin_name
  644. )
  645. # Update context with response
  646. context.set("last_response", {
  647. "text": response_text,
  648. "plugin": plugin_name,
  649. "timestamp": datetime.now(timezone.utc).isoformat()
  650. }, scope=ContextScope.SESSION)
  651. # Add to conversation history
  652. conversation_history = context.get("conversation_history", [], scope=ContextScope.SESSION)
  653. conversation_history.append({
  654. "type": "response",
  655. "text": response_text,
  656. "plugin": plugin_name,
  657. "timestamp": datetime.now(timezone.utc).isoformat()
  658. })
  659. context.set("conversation_history", conversation_history, scope=ContextScope.SESSION)
  660. pprint(f"Generated response for conversation {conversation_id}: {response_text[:50]}...")
  661. return True
  662. def synthesize_speech(
  663. self,
  664. conversation_id: str,
  665. audio_data: bytes,
  666. text: str,
  667. voice_settings: Optional[Dict[str, Any]] = None
  668. ) -> bool:
  669. """
  670. Process synthesized speech for a conversation.
  671. Args:
  672. conversation_id: Conversation identifier
  673. audio_data: Synthesized audio data
  674. text: Original text
  675. voice_settings: Voice settings used
  676. Returns:
  677. True if speech was processed successfully
  678. """
  679. session = self.get_conversation(conversation_id)
  680. state_manager = self.get_conversation_state(conversation_id)
  681. if not session or not state_manager:
  682. return False
  683. # Transition to speech synthesis state
  684. state_manager.transition_to(
  685. ConversationState.SYNTHESIZING_SPEECH,
  686. trigger="speech_synthesized",
  687. source="conversation_manager"
  688. )
  689. # Update current turn with audio data
  690. current_turn = session.get_current_turn()
  691. if current_turn:
  692. session.update_turn_response(
  693. current_turn.turn_id,
  694. response_audio_data=audio_data,
  695. response_voice_settings=voice_settings or {}
  696. )
  697. # Trigger TTS received event
  698. self._trigger_tts_received_event(conversation_id, audio_data, text, voice_settings or {})
  699. # Transition to playing response
  700. state_manager.transition_to(
  701. ConversationState.PLAYING_RESPONSE,
  702. trigger="playing_response",
  703. source="conversation_manager"
  704. )
  705. pprint(f"Synthesized speech for conversation {conversation_id}")
  706. return True
  707. # Plugin Interaction Support
  708. def ask_question(
  709. self,
  710. conversation_id: str,
  711. question: str,
  712. plugin_name: str,
  713. callback: Optional[Callable] = None,
  714. timeout_seconds: Optional[float] = None,
  715. voice_settings: Optional[Dict[str, Any]] = None,
  716. metadata: Optional[Dict[str, Any]] = None
  717. ) -> bool:
  718. """
  719. Ask a question and wait for user response (plugin interaction).
  720. Args:
  721. conversation_id: Conversation identifier
  722. question: Question to ask the user
  723. plugin_name: Name of the plugin asking the question
  724. callback: Optional callback function for the response
  725. timeout_seconds: Timeout for the response
  726. voice_settings: Voice settings for TTS
  727. metadata: Additional metadata
  728. Returns:
  729. True if question was asked successfully
  730. """
  731. session = self.get_conversation(conversation_id)
  732. context = self.get_conversation_context(conversation_id)
  733. state_manager = self.get_conversation_state(conversation_id)
  734. if not session or not context or not state_manager:
  735. pprint(f"Conversation {conversation_id} not found for plugin question")
  736. return False
  737. # Transition to plugin interaction state
  738. state_manager.transition_to(
  739. ConversationState.PLUGIN_INTERACTION,
  740. trigger="plugin_question",
  741. source=plugin_name
  742. )
  743. # Ask question through session
  744. turn = session.ask_question(
  745. question=question,
  746. plugin_name=plugin_name,
  747. callback=callback,
  748. timeout_seconds=timeout_seconds or self._plugin_timeout,
  749. voice_settings=voice_settings
  750. )
  751. # Store plugin question info
  752. self._plugin_questions[conversation_id] = {
  753. "plugin_name": plugin_name,
  754. "question": question,
  755. "turn_id": turn.turn_id,
  756. "callback": callback,
  757. "timeout": datetime.now(timezone.utc) + timedelta(
  758. seconds=timeout_seconds or self._plugin_timeout
  759. ),
  760. "metadata": metadata or {}
  761. }
  762. if callback:
  763. self._pending_callbacks[conversation_id] = callback
  764. # Update context
  765. context.set("plugin_question", {
  766. "plugin": plugin_name,
  767. "question": question,
  768. "turn_id": turn.turn_id,
  769. "timestamp": datetime.now(timezone.utc).isoformat()
  770. }, scope=ContextScope.SESSION)
  771. # Transition to waiting for user
  772. state_manager.transition_to(
  773. ConversationState.WAITING_FOR_USER,
  774. trigger="waiting_for_answer",
  775. source=plugin_name
  776. )
  777. # Generate TTS response for the question
  778. self.generate_response(
  779. conversation_id=conversation_id,
  780. response_text=question,
  781. voice_settings=voice_settings,
  782. plugin_name=plugin_name,
  783. metadata={"is_question": True}
  784. )
  785. pprint(f"Plugin {plugin_name} asked question in conversation {conversation_id}: {question}")
  786. return True
  787. def _handle_plugin_answer(
  788. self,
  789. conversation_id: str,
  790. answer: str,
  791. confidence: float,
  792. speaker_info: Optional[SpeakerInfo]
  793. ) -> Optional[ConversationTurn]:
  794. """Handle an answer to a plugin question."""
  795. session = self.get_conversation(conversation_id)
  796. context = self.get_conversation_context(conversation_id)
  797. state_manager = self.get_conversation_state(conversation_id)
  798. if not session or not context or not state_manager:
  799. return None
  800. plugin_data = self._plugin_questions.get(conversation_id)
  801. if not plugin_data:
  802. return None
  803. # Answer the question through session
  804. turn = session.answer_question(answer, confidence, speaker_info)
  805. # Update statistics
  806. if self._enable_analytics:
  807. self._stats.plugin_interactions += 1
  808. # Clean up plugin question
  809. if conversation_id in self._plugin_questions:
  810. del self._plugin_questions[conversation_id]
  811. if conversation_id in self._pending_callbacks:
  812. del self._pending_callbacks[conversation_id]
  813. # Clear plugin question context
  814. context.delete("plugin_question", scope=ContextScope.SESSION)
  815. # Update context with answer
  816. context.set("plugin_answer", {
  817. "plugin": plugin_data["plugin_name"],
  818. "question": plugin_data["question"],
  819. "answer": answer,
  820. "confidence": confidence,
  821. "timestamp": datetime.now(timezone.utc).isoformat()
  822. }, scope=ContextScope.SESSION)
  823. pprint(f"Received plugin answer in conversation {conversation_id}: {answer}")
  824. return turn
  825. def complete_conversation(
  826. self,
  827. conversation_id: str,
  828. final_response: Optional[str] = None,
  829. metadata: Optional[Dict[str, Any]] = None
  830. ) -> bool:
  831. """
  832. Complete a conversation with optional final response.
  833. Args:
  834. conversation_id: Conversation identifier
  835. final_response: Optional final response
  836. metadata: Additional metadata
  837. Returns:
  838. True if conversation was completed successfully
  839. """
  840. if final_response:
  841. self.generate_response(
  842. conversation_id=conversation_id,
  843. response_text=final_response,
  844. metadata=metadata
  845. )
  846. return self.end_conversation(conversation_id, "completed", final_response)
  847. # Context Management
  848. def get_global_context(self) -> ConversationContext:
  849. """Get the global context manager."""
  850. return self._global_context
  851. def get_speaker_context(self, speaker_id: str) -> Optional[ConversationContext]:
  852. """Get context for a specific speaker."""
  853. return self._speaker_contexts.get(speaker_id)
  854. def get_satellite_context(self, satellite_id: str) -> Optional[ConversationContext]:
  855. """Get context for a specific satellite."""
  856. return self._satellite_contexts.get(satellite_id)
  857. def get_room_context(self, room_id: str) -> Optional[ConversationContext]:
  858. """Get context for a specific room."""
  859. return self._room_contexts.get(room_id)
  860. def set_global_context(self, key: str, value: Any, **kwargs) -> None:
  861. """Set a global context value."""
  862. self._global_context.set(key, value, scope=ContextScope.GLOBAL, **kwargs)
  863. def get_context_value(
  864. self,
  865. conversation_id: str,
  866. key: str,
  867. scopes: Optional[List[ContextScope]] = None,
  868. default: Any = None
  869. ) -> Any:
  870. """
  871. Get a context value searching multiple scopes.
  872. Args:
  873. conversation_id: Conversation identifier
  874. key: Context key
  875. scopes: Scopes to search (defaults to all relevant scopes)
  876. default: Default value if not found
  877. Returns:
  878. Context value or default
  879. """
  880. context = self.get_conversation_context(conversation_id)
  881. if not context:
  882. return default
  883. if scopes is None:
  884. scopes = [ContextScope.SESSION, ContextScope.SPEAKER, ContextScope.SATELLITE,
  885. ContextScope.ROOM, ContextScope.GLOBAL]
  886. return context.get_multi_scope(key, scopes, default)
  887. # Monitoring and Analytics
  888. def get_active_conversations(self) -> List[ConversationSession]:
  889. """Get all active conversation sessions."""
  890. with self._lock:
  891. return list(self._active_conversations.values())
  892. def get_completed_conversations(self) -> List[ConversationSession]:
  893. """Get completed conversation sessions."""
  894. with self._lock:
  895. return list(self._completed_conversations.values())
  896. def get_conversation_by_speaker(self, speaker_id: str) -> List[ConversationSession]:
  897. """Get conversations for a specific speaker."""
  898. with self._lock:
  899. conversations = []
  900. for session in self._active_conversations.values():
  901. if session.speaker_info and session.speaker_info.speaker_id == speaker_id:
  902. conversations.append(session)
  903. return conversations
  904. def get_conversation_by_satellite(self, satellite_id: str) -> List[ConversationSession]:
  905. """Get conversations for a specific satellite."""
  906. with self._lock:
  907. conversations = []
  908. for session in self._active_conversations.values():
  909. if session.satellite_info.satellite_id == satellite_id:
  910. conversations.append(session)
  911. return conversations
  912. def _update_stats(self) -> None:
  913. """Update conversation statistics."""
  914. if not self._enable_analytics:
  915. return
  916. with self._lock:
  917. now = datetime.now(timezone.utc)
  918. # Count conversations by state
  919. active_count = len(self._active_conversations)
  920. completed_count = len(self._completed_conversations)
  921. # Count conversations by final state
  922. timed_out_count = 0
  923. error_count = 0
  924. for session in self._completed_conversations.values():
  925. if session.state == ConversationSessionState.TIMED_OUT:
  926. timed_out_count += 1
  927. elif session.state == ConversationSessionState.ERROR:
  928. error_count += 1
  929. # Calculate turn statistics
  930. total_turns = 0
  931. total_duration = 0.0
  932. for session in list(self._active_conversations.values()) + list(self._completed_conversations.values()):
  933. total_turns += session.turn_count
  934. total_duration += session.duration_seconds
  935. total_conversations = active_count + completed_count
  936. avg_turns = total_turns / max(total_conversations, 1)
  937. avg_duration = total_duration / max(total_conversations, 1)
  938. # Count context entries
  939. context_entries = 0
  940. for context in self._conversation_contexts.values():
  941. for scope_entries in context._contexts.values():
  942. context_entries += len(scope_entries)
  943. # Update statistics
  944. self._stats = ConversationStats(
  945. total_conversations=total_conversations,
  946. active_conversations=active_count,
  947. completed_conversations=completed_count,
  948. timed_out_conversations=timed_out_count,
  949. error_conversations=error_count,
  950. total_turns=total_turns,
  951. average_turns_per_conversation=avg_turns,
  952. average_conversation_duration=avg_duration,
  953. plugin_interactions=self._stats.plugin_interactions,
  954. context_entries=context_entries,
  955. uptime_seconds=(now - self._start_time).total_seconds(),
  956. last_updated=now
  957. )
  958. def get_stats(self) -> ConversationStats:
  959. """Get conversation statistics."""
  960. self._update_stats()
  961. return self._stats
  962. def get_status(self) -> Dict[str, Any]:
  963. """Get comprehensive status information."""
  964. with self._lock:
  965. stats = self.get_stats()
  966. return {
  967. "statistics": stats.to_dict(),
  968. "configuration": {
  969. "max_concurrent_conversations": self._max_concurrent_conversations,
  970. "default_timeout_seconds": self._default_timeout,
  971. "plugin_timeout_seconds": self._plugin_timeout,
  972. "cleanup_interval_seconds": self._cleanup_interval,
  973. "max_history": self._max_history,
  974. "enable_context_persistence": self._enable_context_persistence,
  975. "enable_analytics": self._enable_analytics,
  976. },
  977. "active_conversations": [
  978. {
  979. "conversation_id": session.conversation_id,
  980. "state": session.state.value,
  981. "duration_seconds": session.duration_seconds,
  982. "turn_count": session.turn_count,
  983. "speaker_id": session.speaker_info.speaker_id if session.speaker_info else None,
  984. "satellite_id": session.satellite_info.satellite_id,
  985. "room_id": session.satellite_info.room_id
  986. }
  987. for session in self._active_conversations.values()
  988. ],
  989. "plugin_questions": len(self._plugin_questions),
  990. "pending_callbacks": len(self._pending_callbacks),
  991. "context_managers": {
  992. "conversation_contexts": len(self._conversation_contexts),
  993. "speaker_contexts": len(self._speaker_contexts),
  994. "satellite_contexts": len(self._satellite_contexts),
  995. "room_contexts": len(self._room_contexts),
  996. },
  997. "timestamp": datetime.now(timezone.utc).isoformat()
  998. }
  999. # Cleanup and Maintenance
  1000. def _start_cleanup_timer(self) -> None:
  1001. """Start the background cleanup timer."""
  1002. if self._cleanup_timer:
  1003. self._cleanup_timer.cancel()
  1004. self._cleanup_timer = threading.Timer(self._cleanup_interval, self._periodic_cleanup)
  1005. self._cleanup_timer.daemon = True
  1006. self._cleanup_timer.start()
  1007. def _periodic_cleanup(self) -> None:
  1008. """Perform periodic cleanup tasks."""
  1009. try:
  1010. # Check for timed out conversations
  1011. self._check_conversation_timeouts()
  1012. # Check for timed out plugin questions
  1013. self._check_plugin_timeouts()
  1014. # Clean up expired contexts
  1015. self._cleanup_expired_contexts()
  1016. # Clean up completed conversations
  1017. self._cleanup_completed_conversations()
  1018. # Update statistics
  1019. self._update_stats()
  1020. except Exception as e:
  1021. pprint(f"Error during periodic cleanup: {e}")
  1022. finally:
  1023. # Restart timer
  1024. self._start_cleanup_timer()
  1025. def _check_conversation_timeouts(self) -> None:
  1026. """Check for and handle conversation timeouts."""
  1027. with self._lock:
  1028. timed_out = []
  1029. for conversation_id, session in self._active_conversations.items():
  1030. if session.is_timed_out:
  1031. timed_out.append(conversation_id)
  1032. for conversation_id in timed_out:
  1033. pprint(f"Conversation {conversation_id} timed out")
  1034. self.end_conversation(conversation_id, "timeout")
  1035. def _check_plugin_timeouts(self) -> None:
  1036. """Check for and handle plugin question timeouts."""
  1037. with self._lock:
  1038. now = datetime.now(timezone.utc)
  1039. timed_out = []
  1040. for conversation_id, plugin_data in self._plugin_questions.items():
  1041. if now > plugin_data["timeout"]:
  1042. timed_out.append(conversation_id)
  1043. for conversation_id in timed_out:
  1044. session = self._active_conversations.get(conversation_id)
  1045. if session:
  1046. session.cancel_plugin_question("timeout")
  1047. pprint(f"Plugin question timed out for conversation {conversation_id}")
  1048. # Clean up
  1049. if conversation_id in self._plugin_questions:
  1050. del self._plugin_questions[conversation_id]
  1051. if conversation_id in self._pending_callbacks:
  1052. del self._pending_callbacks[conversation_id]
  1053. def _cleanup_expired_contexts(self) -> None:
  1054. """Clean up expired context entries."""
  1055. # Clean up conversation contexts
  1056. for context in self._conversation_contexts.values():
  1057. context.cleanup_expired()
  1058. # Clean up speaker contexts
  1059. for context in self._speaker_contexts.values():
  1060. context.cleanup_expired()
  1061. # Clean up satellite contexts
  1062. for context in self._satellite_contexts.values():
  1063. context.cleanup_expired()
  1064. # Clean up room contexts
  1065. for context in self._room_contexts.values():
  1066. context.cleanup_expired()
  1067. # Clean up global context
  1068. self._global_context.cleanup_expired()
  1069. def _cleanup_completed_conversations(self) -> None:
  1070. """Clean up old completed conversations."""
  1071. with self._lock:
  1072. if len(self._completed_conversations) <= self._max_history:
  1073. return
  1074. # Remove oldest conversations
  1075. excess_count = len(self._completed_conversations) - self._max_history
  1076. for _ in range(excess_count):
  1077. if self._completed_conversations:
  1078. oldest_key = next(iter(self._completed_conversations))
  1079. del self._completed_conversations[oldest_key]
  1080. # Event Integration
  1081. def set_event_handler(self, event_handler) -> None:
  1082. """Set the event handler for triggering conversation events."""
  1083. self._event_handler = event_handler
  1084. def set_satellite_manager(self, satellite_manager) -> None:
  1085. """Set the satellite manager for satellite integration."""
  1086. self._satellite_manager = satellite_manager
  1087. def _trigger_conversation_event(
  1088. self,
  1089. event_type: str,
  1090. session: ConversationSession,
  1091. additional_data: Optional[Dict[str, Any]] = None
  1092. ) -> None:
  1093. """Trigger a conversation-related event."""
  1094. try:
  1095. if self._application_ref and self._application_ref():
  1096. app = self._application_ref()
  1097. event_handler = app.get_event_handler()
  1098. from ..events import EventDataFactory
  1099. event_data = EventDataFactory.create_event_data(
  1100. event_type,
  1101. conversation_id=session.conversation_id,
  1102. speaker_info=session.speaker_info,
  1103. satellite_info=session.satellite_info,
  1104. state=session.state.value,
  1105. turn_count=session.turn_count,
  1106. duration_seconds=session.duration_seconds,
  1107. **(additional_data or {})
  1108. )
  1109. event_handler.trigger_event(event_type, event_data)
  1110. except Exception as e:
  1111. pprint(f"Error triggering conversation event: {e}")
  1112. def _trigger_audio_input_event(
  1113. self,
  1114. conversation_id: str,
  1115. audio_data: bytes,
  1116. speaker_info: Optional[SpeakerInfo],
  1117. satellite_info: SatelliteInfo
  1118. ) -> None:
  1119. """Trigger raw audio input received event."""
  1120. try:
  1121. if self._application_ref and self._application_ref():
  1122. app = self._application_ref()
  1123. event_handler = app.get_event_handler()
  1124. from ..events import EventDataFactory
  1125. audio_data = EventDataFactory.create_event_data(
  1126. "raw_audio_input_received",
  1127. conversation_id=conversation_id,
  1128. audio_data=audio_data,
  1129. speaker_info=speaker_info,
  1130. satellite_info=satellite_info,
  1131. sample_rate=16000,
  1132. channels=1,
  1133. bit_depth=16,
  1134. duration_seconds=len(audio_data) / (16000 * 2) # Assuming 16kHz 16-bit
  1135. )
  1136. event_handler.trigger_event("raw_audio_input_received", audio_data)
  1137. except Exception as e:
  1138. pprint(f"Error triggering audio input event: {e}")
  1139. def _trigger_text_received_event(
  1140. self,
  1141. conversation_id: str,
  1142. text: str,
  1143. confidence: float,
  1144. speaker_info: Optional[SpeakerInfo],
  1145. satellite_info: SatelliteInfo
  1146. ) -> None:
  1147. """Trigger text received event."""
  1148. try:
  1149. if self._application_ref and self._application_ref():
  1150. app = self._application_ref()
  1151. event_handler = app.get_event_handler()
  1152. from ..events import EventDataFactory
  1153. text_data = EventDataFactory.create_event_data(
  1154. "text_received",
  1155. conversation_id=conversation_id,
  1156. text=text,
  1157. confidence=confidence,
  1158. speaker_info=speaker_info,
  1159. satellite_info=satellite_info
  1160. )
  1161. event_handler.trigger_event("text_received", text_data)
  1162. except Exception as e:
  1163. pprint(f"Error triggering text received event: {e}")
  1164. def _trigger_intent_received_event(
  1165. self,
  1166. conversation_id: str,
  1167. intent: str,
  1168. entities: Dict[str, Any],
  1169. confidence: float,
  1170. original_text: str,
  1171. speaker_info: Optional[SpeakerInfo],
  1172. satellite_info: SatelliteInfo
  1173. ) -> None:
  1174. """Trigger intent received event."""
  1175. try:
  1176. if self._application_ref and self._application_ref():
  1177. app = self._application_ref()
  1178. event_handler = app.get_event_handler()
  1179. from ..events import EventDataFactory
  1180. intent_data = EventDataFactory.create_event_data(
  1181. "intent_received",
  1182. conversation_id=conversation_id,
  1183. intent=intent,
  1184. entities=entities,
  1185. confidence=confidence,
  1186. original_text=original_text,
  1187. speaker_info=speaker_info,
  1188. satellite_info=satellite_info
  1189. )
  1190. event_handler.trigger_event("intent_received", intent_data)
  1191. except Exception as e:
  1192. pprint(f"Error triggering intent received event: {e}")
  1193. def _trigger_tts_received_event(
  1194. self,
  1195. conversation_id: str,
  1196. audio_data: bytes,
  1197. text: str,
  1198. voice_settings: Dict[str, Any]
  1199. ) -> None:
  1200. """Trigger TTS received event."""
  1201. try:
  1202. if self._application_ref and self._application_ref():
  1203. app = self._application_ref()
  1204. event_handler = app.get_event_handler()
  1205. from ..events import EventDataFactory
  1206. tts_data = EventDataFactory.create_event_data(
  1207. "tts_received",
  1208. conversation_id=conversation_id,
  1209. audio_data=audio_data,
  1210. text=text,
  1211. voice_settings=voice_settings,
  1212. sample_rate=16000,
  1213. channels=1,
  1214. bit_depth=16,
  1215. duration_seconds=len(audio_data) / (16000 * 2)
  1216. )
  1217. event_handler.trigger_event("tts_received", tts_data)
  1218. except Exception as e:
  1219. pprint(f"Error triggering TTS received event: {e}")
  1220. # Persistence Operations
  1221. def _load_persisted_contexts(self) -> None:
  1222. """Load persisted context data."""
  1223. if not self._enable_context_persistence or not self._context_storage_path.exists():
  1224. return
  1225. try:
  1226. # Load global context
  1227. global_path = self._context_storage_path / "global"
  1228. if global_path.exists():
  1229. self._global_context._load_persistent_context()
  1230. # Load speaker contexts
  1231. speakers_path = self._context_storage_path / "speakers"
  1232. if speakers_path.exists():
  1233. for speaker_dir in speakers_path.iterdir():
  1234. if speaker_dir.is_dir():
  1235. speaker_id = speaker_dir.name
  1236. self._speaker_contexts[speaker_id] = create_conversation_context(
  1237. conversation_id=f"speaker_{speaker_id}",
  1238. speaker_id=speaker_id,
  1239. enable_persistence=True,
  1240. persistence_path=speaker_dir
  1241. )
  1242. # Load satellite contexts
  1243. satellites_path = self._context_storage_path / "satellites"
  1244. if satellites_path.exists():
  1245. for satellite_dir in satellites_path.iterdir():
  1246. if satellite_dir.is_dir():
  1247. satellite_id = satellite_dir.name
  1248. self._satellite_contexts[satellite_id] = create_conversation_context(
  1249. conversation_id=f"satellite_{satellite_id}",
  1250. satellite_id=satellite_id,
  1251. enable_persistence=True,
  1252. persistence_path=satellite_dir
  1253. )
  1254. # Load room contexts
  1255. rooms_path = self._context_storage_path / "rooms"
  1256. if rooms_path.exists():
  1257. for room_dir in rooms_path.iterdir():
  1258. if room_dir.is_dir():
  1259. room_id = room_dir.name
  1260. self._room_contexts[room_id] = create_conversation_context(
  1261. conversation_id=f"room_{room_id}",
  1262. room_id=room_id,
  1263. enable_persistence=True,
  1264. persistence_path=room_dir
  1265. )
  1266. pprint("Loaded persisted conversation contexts")
  1267. except Exception as e:
  1268. pprint(f"Error loading persisted contexts: {e}")
  1269. # Utility Methods
  1270. def _generate_conversation_id(self) -> str:
  1271. """Generate a unique conversation ID."""
  1272. timestamp = int(time.time() * 1000)
  1273. return f"conv_{timestamp}_{uuid.uuid4().hex[:8]}"
  1274. def shutdown(self) -> None:
  1275. """Shutdown the conversation manager gracefully."""
  1276. with self._lock:
  1277. pprint("Shutting down conversation manager...")
  1278. # Cancel cleanup timer
  1279. if self._cleanup_timer:
  1280. self._cleanup_timer.cancel()
  1281. self._cleanup_timer = None
  1282. # End all active conversations
  1283. active_ids = list(self._active_conversations.keys())
  1284. for conversation_id in active_ids:
  1285. self.end_conversation(conversation_id, "shutdown")
  1286. # Shutdown executor
  1287. self._executor.shutdown(wait=True)
  1288. # Clear all data
  1289. self._active_conversations.clear()
  1290. self._conversation_contexts.clear()
  1291. self._conversation_states.clear()
  1292. self._plugin_questions.clear()
  1293. self._pending_callbacks.clear()
  1294. pprint("Conversation manager shutdown complete")
  1295. def __str__(self) -> str:
  1296. """String representation of the conversation manager."""
  1297. stats = self.get_stats()
  1298. return (f"ConversationManager(active={stats.active_conversations}, "
  1299. f"completed={stats.completed_conversations})")
  1300. def __repr__(self) -> str:
  1301. """Detailed representation of the conversation manager."""
  1302. return (f"ConversationManager(max_concurrent={self._max_concurrent_conversations}, "
  1303. f"active={len(self._active_conversations)}, "
  1304. f"completed={len(self._completed_conversations)})")
  1305. # Factory function
  1306. def create_conversation_manager(
  1307. application,
  1308. **kwargs
  1309. ) -> ConversationManager:
  1310. """
  1311. Create a ConversationManager with default configuration.
  1312. Args:
  1313. application: Application container reference
  1314. **kwargs: Additional configuration parameters
  1315. Returns:
  1316. ConversationManager: Configured conversation manager
  1317. """
  1318. return ConversationManager(application, **kwargs)
  1319. # Module exports
  1320. __all__ = [
  1321. "ConversationManager",
  1322. "ConversationManagerError",
  1323. "ConversationNotFoundError",
  1324. "ConversationTimeoutError",
  1325. "ConversationStateError",
  1326. "ConversationStats",
  1327. "create_conversation_manager"
  1328. ]