profiler.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387
  1. # -*- coding: utf-8 -*-
  2. """
  3. Conversation Profiler — misst Phasen-Dauern im Debug-Modus.
  4. Registriert sich mit SYSTEM-Prioritaet auf alle relevanten Events
  5. der Conversation-Pipeline und zeichnet Timestamps auf.
  6. WICHTIG: SYSTEM (0) statt MONITOR (100) ist notwendig, weil die Events
  7. verschachtelt sind — ein Handler fuer raw_audio_received loest intern
  8. stt_completed → speech_recognized → intent_received → ... → tts_completed aus.
  9. Mit MONITOR wuerden die Timestamps erst NACH der gesamten inneren Kette
  10. aufgezeichnet, was zu fehlenden Phasen fuehrt.
  11. Tracking ueber satellite_id als primaeren Key, da dies das einzige
  12. Feld ist das in allen Events konsistent vorhanden ist.
  13. Ausgabe erfolgt bei tts_completed (letztes Event der Pipeline),
  14. da conversation_ended im Server-Modus nicht als Event emittiert wird.
  15. Null Overhead wenn nicht im Debug-Modus (Guard in jedem Handler).
  16. """
  17. from __future__ import annotations
  18. import time
  19. from dataclasses import dataclass, field
  20. from typing import Any, TYPE_CHECKING
  21. from trixy_core.service.iservice import IService
  22. from trixy_core.service.enums import ServicePriority, ServiceGroup
  23. from trixy_core.events.decorators import TrixyEvent
  24. from trixy_core.events.enums import EventPriority
  25. from trixy_core.utils.debug import pinfo, pdebug
  26. if TYPE_CHECKING:
  27. from trixy_core.application import IApplication
  28. @dataclass
  29. class SessionProfile:
  30. """Timing-Daten fuer eine Conversation-Session."""
  31. satellite_id: str = ""
  32. session_id: str = ""
  33. timestamps: dict[str, float] = field(default_factory=dict)
  34. providers: dict[str, str] = field(default_factory=dict)
  35. text: str = ""
  36. intent: str = ""
  37. confidence: float = 0.0
  38. is_keyboard: bool = False
  39. def _get_field(data: Any, name: str, default: Any = "") -> Any:
  40. """Extrahiert ein Feld aus EventData (trigger) oder dict-metadata (emit)."""
  41. # Direktes Attribut (typisierte EventData via trigger())
  42. val = getattr(data, name, None)
  43. if val is not None and val != "":
  44. return val
  45. # metadata-Dict (emit()-Events: Felder in EventData.metadata)
  46. if hasattr(data, "get"):
  47. val = data.get(name)
  48. if val is not None and val != "":
  49. return val
  50. return default
  51. # Phasen-Definition: (Name, Start-Event, End-Event, Default-Komponente)
  52. _PHASES: list[tuple[str, str, str, str]] = [
  53. ("Wakeword", "wakeword_detected", "conversation_started", "WakewordService"),
  54. ("Recording", "conversation_started", "raw_audio_received", "AudioAccumulator"),
  55. ("STT", "raw_audio_received", "stt_completed", "STT"),
  56. ("STT-Korrektur", "stt_completed", "speech_recognized", "STTCorrector"),
  57. ("NLP", "speech_recognized", "intent_received", "NLP"),
  58. ("Handler", "intent_received", "intent_handled", "Handler"),
  59. ("LLM", "create_output_text", "output_text_created", "LLM"),
  60. ("TTS", "output_text_created", "tts_completed", "TTS"),
  61. ]
  62. class ConversationProfiler(IService):
  63. """
  64. Conversation Profiler — misst Phasen-Dauern im Debug-Modus.
  65. Tracking ueber satellite_id als primaeren Key.
  66. Ausgabe bei tts_completed (letztes Event der Pipeline).
  67. Null Overhead wenn nicht im Debug-Modus.
  68. """
  69. PRIORITY = ServicePriority.OPTIONAL
  70. GROUP = ServiceGroup.CONVERSATION
  71. DEPENDENCIES: list[str] = []
  72. NAME = "ConversationProfiler"
  73. def __init__(self, application: "IApplication") -> None:
  74. super().__init__(application)
  75. # satellite_id → SessionProfile
  76. self._profiles: dict[str, SessionProfile] = {}
  77. async def start(self) -> None:
  78. """Startet den Profiler."""
  79. if self._application.debug:
  80. pinfo("ConversationProfiler gestartet (Debug-Modus)")
  81. async def stop(self) -> None:
  82. """Stoppt den Profiler."""
  83. self._profiles.clear()
  84. # === Hilfsmethoden ===
  85. def _get_sat_id(self, event_data: Any) -> str:
  86. """Extrahiert satellite_id aus Event-Daten."""
  87. return str(_get_field(event_data, "satellite_id", ""))
  88. def _get_profile(self, sat_id: str) -> SessionProfile | None:
  89. """Holt ein bestehendes Profil fuer den Satellite."""
  90. if not sat_id:
  91. return None
  92. return self._profiles.get(sat_id)
  93. def _ensure_profile(self, sat_id: str) -> SessionProfile | None:
  94. """Holt oder erstellt ein Profil fuer den Satellite."""
  95. if not sat_id:
  96. return None
  97. if sat_id not in self._profiles:
  98. self._profiles[sat_id] = SessionProfile(satellite_id=sat_id)
  99. return self._profiles[sat_id]
  100. # === Event-Handler (alle mit SYSTEM-Prioritaet) ===
  101. @TrixyEvent(["wakeword_detected"], priority=EventPriority.SYSTEM)
  102. async def on_wakeword_detected(self, event_name: str, event_data: Any) -> None:
  103. """Wakeword erkannt — Startpunkt der Pipeline."""
  104. if not self._application.debug:
  105. return
  106. sat_id = self._get_sat_id(event_data)
  107. if not sat_id:
  108. return
  109. # Neues Profil starten (altes verwerfen falls vorhanden)
  110. self._profiles[sat_id] = SessionProfile(satellite_id=sat_id)
  111. profile = self._profiles[sat_id]
  112. profile.timestamps["wakeword_detected"] = time.monotonic()
  113. profile.session_id = str(_get_field(event_data, "session_id", ""))
  114. @TrixyEvent(["conversation_started"], priority=EventPriority.SYSTEM)
  115. async def on_conversation_started(self, event_name: str, event_data: Any) -> None:
  116. """Conversation gestartet (nach Arbitrierung)."""
  117. if not self._application.debug:
  118. return
  119. sat_id = self._get_sat_id(event_data)
  120. if not sat_id:
  121. return
  122. profile = self._ensure_profile(sat_id)
  123. if profile:
  124. profile.timestamps["conversation_started"] = time.monotonic()
  125. # conversation_id als session_id uebernehmen (Server-seitige ID)
  126. conv_id = str(_get_field(event_data, "conversation_id", ""))
  127. if conv_id:
  128. profile.session_id = conv_id
  129. @TrixyEvent(["raw_audio_received"], priority=EventPriority.SYSTEM)
  130. async def on_raw_audio_received(self, event_name: str, event_data: Any) -> None:
  131. """Rohe Audio-Daten empfangen."""
  132. if not self._application.debug:
  133. return
  134. sat_id = self._get_sat_id(event_data)
  135. profile = self._get_profile(sat_id)
  136. if profile and "raw_audio_received" not in profile.timestamps:
  137. profile.timestamps["raw_audio_received"] = time.monotonic()
  138. @TrixyEvent(["stt_completed"], priority=EventPriority.SYSTEM)
  139. async def on_stt_completed(self, event_name: str, event_data: Any) -> None:
  140. """STT abgeschlossen."""
  141. if not self._application.debug:
  142. return
  143. sat_id = self._get_sat_id(event_data)
  144. profile = self._get_profile(sat_id)
  145. if profile:
  146. profile.timestamps["stt_completed"] = time.monotonic()
  147. provider = str(_get_field(event_data, "provider", ""))
  148. if provider:
  149. profile.providers["STT"] = provider
  150. @TrixyEvent(["speech_recognized"], priority=EventPriority.SYSTEM)
  151. async def on_speech_recognized(self, event_name: str, event_data: Any) -> None:
  152. """Sprache erkannt (nach STT-Korrektur oder Keyboard-Input)."""
  153. if not self._application.debug:
  154. return
  155. source = str(_get_field(event_data, "source", "stt"))
  156. sat_id = self._get_sat_id(event_data)
  157. text = str(_get_field(event_data, "text", ""))
  158. # Keyboard-Input: Neues Profil ab diesem Punkt
  159. if source == "keyboard":
  160. # Fuer Standalone: satellite_id kann leer sein
  161. key = sat_id or "standalone"
  162. self._profiles[key] = SessionProfile(
  163. satellite_id=key,
  164. is_keyboard=True,
  165. text=text,
  166. )
  167. self._profiles[key].timestamps["speech_recognized"] = time.monotonic()
  168. return
  169. # Normaler STT-Pfad
  170. profile = self._get_profile(sat_id)
  171. if profile:
  172. profile.timestamps["speech_recognized"] = time.monotonic()
  173. if text:
  174. profile.text = text
  175. @TrixyEvent(["intent_received"], priority=EventPriority.SYSTEM)
  176. async def on_intent_received(self, event_name: str, event_data: Any) -> None:
  177. """Intent erkannt."""
  178. if not self._application.debug:
  179. return
  180. sat_id = self._get_sat_id(event_data)
  181. profile = self._get_profile(sat_id)
  182. if not profile:
  183. return
  184. profile.timestamps["intent_received"] = time.monotonic()
  185. intent = str(_get_field(event_data, "intent", ""))
  186. confidence = float(_get_field(event_data, "confidence", 0.0))
  187. if intent:
  188. profile.intent = intent
  189. profile.providers["Handler"] = intent
  190. if confidence:
  191. profile.confidence = confidence
  192. if not profile.text:
  193. profile.text = str(_get_field(event_data, "original_text", ""))
  194. @TrixyEvent(["intent_handled"], priority=EventPriority.SYSTEM)
  195. async def on_intent_handled(self, event_name: str, event_data: Any) -> None:
  196. """Intent verarbeitet."""
  197. if not self._application.debug:
  198. return
  199. sat_id = self._get_sat_id(event_data)
  200. profile = self._get_profile(sat_id)
  201. if profile:
  202. profile.timestamps["intent_handled"] = time.monotonic()
  203. @TrixyEvent(["create_output_text"], priority=EventPriority.SYSTEM)
  204. async def on_create_output_text(self, event_name: str, event_data: Any) -> None:
  205. """LLM-Textgenerierung angefordert."""
  206. if not self._application.debug:
  207. return
  208. sat_id = self._get_sat_id(event_data)
  209. profile = self._get_profile(sat_id)
  210. if profile:
  211. profile.timestamps["create_output_text"] = time.monotonic()
  212. @TrixyEvent(["output_text_created"], priority=EventPriority.SYSTEM)
  213. async def on_output_text_created(self, event_name: str, event_data: Any) -> None:
  214. """Antworttext generiert."""
  215. if not self._application.debug:
  216. return
  217. sat_id = self._get_sat_id(event_data)
  218. profile = self._get_profile(sat_id)
  219. if profile:
  220. profile.timestamps["output_text_created"] = time.monotonic()
  221. @TrixyEvent(["tts_completed"], priority=EventPriority.SYSTEM)
  222. async def on_tts_completed(self, event_name: str, event_data: Any) -> None:
  223. """TTS abgeschlossen — Profil ausgeben (letztes Event der Pipeline)."""
  224. if not self._application.debug:
  225. return
  226. sat_id = self._get_sat_id(event_data)
  227. profile = self._get_profile(sat_id)
  228. if not profile:
  229. return
  230. profile.timestamps["tts_completed"] = time.monotonic()
  231. provider = str(_get_field(event_data, "provider", ""))
  232. if provider:
  233. profile.providers["TTS"] = provider
  234. # Profil ausgeben und aufraeumen
  235. self._print_profile(profile)
  236. self._profiles.pop(sat_id, None)
  237. # === Ausgabe ===
  238. def _print_profile(self, profile: SessionProfile) -> None:
  239. """Gibt das Profil als formatierte Tabelle aus."""
  240. ts = profile.timestamps
  241. # Phasen berechnen
  242. phases: list[tuple[str, str, float]] = []
  243. if profile.is_keyboard:
  244. # Keyboard-Input: Nur Phasen ab NLP (kein Wakeword/Recording/STT)
  245. phase_list = _PHASES[4:]
  246. else:
  247. phase_list = _PHASES
  248. for phase_name, start_evt, end_evt, default_component in phase_list:
  249. if start_evt in ts and end_evt in ts:
  250. duration = ts[end_evt] - ts[start_evt]
  251. component = profile.providers.get(phase_name, default_component)
  252. phases.append((phase_name, component, duration))
  253. if not phases:
  254. return
  255. # Gesamt-Dauer
  256. all_ts = sorted(ts.values())
  257. total_duration = all_ts[-1] - all_ts[0] if len(all_ts) >= 2 else 0.0
  258. # Tabelle formatieren
  259. w_phase = 16
  260. w_comp = 16
  261. w_dur = 19
  262. w_inner = w_phase + 1 + w_comp + 1 + w_dur
  263. # Session-ID kuerzen fuer Titel
  264. display_id = profile.session_id or profile.satellite_id
  265. short_id = display_id[:16] if len(display_id) > 16 else display_id
  266. lines: list[str] = []
  267. lines.append(f"╔{'═' * w_inner}╗")
  268. title = f"Conversation Profil — {short_id}"
  269. lines.append(f"║{title:^{w_inner}}║")
  270. lines.append(f"╠{'═' * w_inner}╣")
  271. # Header
  272. lines.append(
  273. f"║ {'Phase':<{w_phase - 1}}│ {'Komponente':<{w_comp - 1}}│ {'Dauer':<{w_dur - 1}}║"
  274. )
  275. lines.append(
  276. f"╠{'═' * w_phase}╪{'═' * w_comp}╪{'═' * w_dur}╣"
  277. )
  278. # Phasen-Zeilen
  279. for phase_name, component, duration in phases:
  280. dur_str = f"{duration:>7.3f} sec"
  281. lines.append(
  282. f"║ {phase_name:<{w_phase - 1}}│ {component:<{w_comp - 1}}│ {dur_str:<{w_dur - 1}}║"
  283. )
  284. # Gesamt
  285. lines.append(
  286. f"╠{'═' * w_phase}╪{'═' * w_comp}╪{'═' * w_dur}╣"
  287. )
  288. total_str = f"{total_duration:>7.3f} sec"
  289. lines.append(
  290. f"║ {'GESAMT':<{w_phase - 1}}│ {'':<{w_comp - 1}}│ {total_str:<{w_dur - 1}}║"
  291. )
  292. # Meta-Informationen
  293. lines.append(
  294. f"╠{'═' * w_phase}╧{'═' * w_comp}╧{'═' * w_dur}╣"
  295. )
  296. if profile.text:
  297. text_display = profile.text[:w_inner - 9]
  298. lines.append(f"║ Text: \"{text_display}\"{' ' * max(0, w_inner - 9 - len(text_display))}║")
  299. if profile.intent:
  300. intent_info = f"Intent: {profile.intent}"
  301. if profile.confidence:
  302. intent_info += f" (confidence: {profile.confidence:.2f})"
  303. intent_display = intent_info[:w_inner - 2]
  304. lines.append(f"║ {intent_display:<{w_inner - 2}} ║")
  305. if profile.satellite_id:
  306. sat_display = f"Satellite: {profile.satellite_id}"[:w_inner - 2]
  307. lines.append(f"║ {sat_display:<{w_inner - 2}} ║")
  308. if profile.is_keyboard:
  309. lines.append(f"║ {'Quelle: Keyboard-Input':<{w_inner - 2}} ║")
  310. lines.append(f"╚{'═' * w_inner}╝")
  311. # Ausgabe ueber pinfo()
  312. for line in lines:
  313. pinfo(line)