| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387 |
- # -*- coding: utf-8 -*-
- """
- Conversation Profiler — misst Phasen-Dauern im Debug-Modus.
- Registriert sich mit SYSTEM-Prioritaet auf alle relevanten Events
- der Conversation-Pipeline und zeichnet Timestamps auf.
- WICHTIG: SYSTEM (0) statt MONITOR (100) ist notwendig, weil die Events
- verschachtelt sind — ein Handler fuer raw_audio_received loest intern
- stt_completed → speech_recognized → intent_received → ... → tts_completed aus.
- Mit MONITOR wuerden die Timestamps erst NACH der gesamten inneren Kette
- aufgezeichnet, was zu fehlenden Phasen fuehrt.
- Tracking ueber satellite_id als primaeren Key, da dies das einzige
- Feld ist das in allen Events konsistent vorhanden ist.
- Ausgabe erfolgt bei tts_completed (letztes Event der Pipeline),
- da conversation_ended im Server-Modus nicht als Event emittiert wird.
- Null Overhead wenn nicht im Debug-Modus (Guard in jedem Handler).
- """
- from __future__ import annotations
- import time
- from dataclasses import dataclass, field
- from typing import Any, TYPE_CHECKING
- from trixy_core.service.iservice import IService
- from trixy_core.service.enums import ServicePriority, ServiceGroup
- from trixy_core.events.decorators import TrixyEvent
- from trixy_core.events.enums import EventPriority
- from trixy_core.utils.debug import pinfo, pdebug
- if TYPE_CHECKING:
- from trixy_core.application import IApplication
- @dataclass
- class SessionProfile:
- """Timing-Daten fuer eine Conversation-Session."""
- satellite_id: str = ""
- session_id: str = ""
- timestamps: dict[str, float] = field(default_factory=dict)
- providers: dict[str, str] = field(default_factory=dict)
- text: str = ""
- intent: str = ""
- confidence: float = 0.0
- is_keyboard: bool = False
- def _get_field(data: Any, name: str, default: Any = "") -> Any:
- """Extrahiert ein Feld aus EventData (trigger) oder dict-metadata (emit)."""
- # Direktes Attribut (typisierte EventData via trigger())
- val = getattr(data, name, None)
- if val is not None and val != "":
- return val
- # metadata-Dict (emit()-Events: Felder in EventData.metadata)
- if hasattr(data, "get"):
- val = data.get(name)
- if val is not None and val != "":
- return val
- return default
- # Phasen-Definition: (Name, Start-Event, End-Event, Default-Komponente)
- _PHASES: list[tuple[str, str, str, str]] = [
- ("Wakeword", "wakeword_detected", "conversation_started", "WakewordService"),
- ("Recording", "conversation_started", "raw_audio_received", "AudioAccumulator"),
- ("STT", "raw_audio_received", "stt_completed", "STT"),
- ("STT-Korrektur", "stt_completed", "speech_recognized", "STTCorrector"),
- ("NLP", "speech_recognized", "intent_received", "NLP"),
- ("Handler", "intent_received", "intent_handled", "Handler"),
- ("LLM", "create_output_text", "output_text_created", "LLM"),
- ("TTS", "output_text_created", "tts_completed", "TTS"),
- ]
- class ConversationProfiler(IService):
- """
- Conversation Profiler — misst Phasen-Dauern im Debug-Modus.
- Tracking ueber satellite_id als primaeren Key.
- Ausgabe bei tts_completed (letztes Event der Pipeline).
- Null Overhead wenn nicht im Debug-Modus.
- """
- PRIORITY = ServicePriority.OPTIONAL
- GROUP = ServiceGroup.CONVERSATION
- DEPENDENCIES: list[str] = []
- NAME = "ConversationProfiler"
- def __init__(self, application: "IApplication") -> None:
- super().__init__(application)
- # satellite_id → SessionProfile
- self._profiles: dict[str, SessionProfile] = {}
- async def start(self) -> None:
- """Startet den Profiler."""
- if self._application.debug:
- pinfo("ConversationProfiler gestartet (Debug-Modus)")
- async def stop(self) -> None:
- """Stoppt den Profiler."""
- self._profiles.clear()
- # === Hilfsmethoden ===
- def _get_sat_id(self, event_data: Any) -> str:
- """Extrahiert satellite_id aus Event-Daten."""
- return str(_get_field(event_data, "satellite_id", ""))
- def _get_profile(self, sat_id: str) -> SessionProfile | None:
- """Holt ein bestehendes Profil fuer den Satellite."""
- if not sat_id:
- return None
- return self._profiles.get(sat_id)
- def _ensure_profile(self, sat_id: str) -> SessionProfile | None:
- """Holt oder erstellt ein Profil fuer den Satellite."""
- if not sat_id:
- return None
- if sat_id not in self._profiles:
- self._profiles[sat_id] = SessionProfile(satellite_id=sat_id)
- return self._profiles[sat_id]
- # === Event-Handler (alle mit SYSTEM-Prioritaet) ===
- @TrixyEvent(["wakeword_detected"], priority=EventPriority.SYSTEM)
- async def on_wakeword_detected(self, event_name: str, event_data: Any) -> None:
- """Wakeword erkannt — Startpunkt der Pipeline."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- if not sat_id:
- return
- # Neues Profil starten (altes verwerfen falls vorhanden)
- self._profiles[sat_id] = SessionProfile(satellite_id=sat_id)
- profile = self._profiles[sat_id]
- profile.timestamps["wakeword_detected"] = time.monotonic()
- profile.session_id = str(_get_field(event_data, "session_id", ""))
- @TrixyEvent(["conversation_started"], priority=EventPriority.SYSTEM)
- async def on_conversation_started(self, event_name: str, event_data: Any) -> None:
- """Conversation gestartet (nach Arbitrierung)."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- if not sat_id:
- return
- profile = self._ensure_profile(sat_id)
- if profile:
- profile.timestamps["conversation_started"] = time.monotonic()
- # conversation_id als session_id uebernehmen (Server-seitige ID)
- conv_id = str(_get_field(event_data, "conversation_id", ""))
- if conv_id:
- profile.session_id = conv_id
- @TrixyEvent(["raw_audio_received"], priority=EventPriority.SYSTEM)
- async def on_raw_audio_received(self, event_name: str, event_data: Any) -> None:
- """Rohe Audio-Daten empfangen."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if profile and "raw_audio_received" not in profile.timestamps:
- profile.timestamps["raw_audio_received"] = time.monotonic()
- @TrixyEvent(["stt_completed"], priority=EventPriority.SYSTEM)
- async def on_stt_completed(self, event_name: str, event_data: Any) -> None:
- """STT abgeschlossen."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if profile:
- profile.timestamps["stt_completed"] = time.monotonic()
- provider = str(_get_field(event_data, "provider", ""))
- if provider:
- profile.providers["STT"] = provider
- @TrixyEvent(["speech_recognized"], priority=EventPriority.SYSTEM)
- async def on_speech_recognized(self, event_name: str, event_data: Any) -> None:
- """Sprache erkannt (nach STT-Korrektur oder Keyboard-Input)."""
- if not self._application.debug:
- return
- source = str(_get_field(event_data, "source", "stt"))
- sat_id = self._get_sat_id(event_data)
- text = str(_get_field(event_data, "text", ""))
- # Keyboard-Input: Neues Profil ab diesem Punkt
- if source == "keyboard":
- # Fuer Standalone: satellite_id kann leer sein
- key = sat_id or "standalone"
- self._profiles[key] = SessionProfile(
- satellite_id=key,
- is_keyboard=True,
- text=text,
- )
- self._profiles[key].timestamps["speech_recognized"] = time.monotonic()
- return
- # Normaler STT-Pfad
- profile = self._get_profile(sat_id)
- if profile:
- profile.timestamps["speech_recognized"] = time.monotonic()
- if text:
- profile.text = text
- @TrixyEvent(["intent_received"], priority=EventPriority.SYSTEM)
- async def on_intent_received(self, event_name: str, event_data: Any) -> None:
- """Intent erkannt."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if not profile:
- return
- profile.timestamps["intent_received"] = time.monotonic()
- intent = str(_get_field(event_data, "intent", ""))
- confidence = float(_get_field(event_data, "confidence", 0.0))
- if intent:
- profile.intent = intent
- profile.providers["Handler"] = intent
- if confidence:
- profile.confidence = confidence
- if not profile.text:
- profile.text = str(_get_field(event_data, "original_text", ""))
- @TrixyEvent(["intent_handled"], priority=EventPriority.SYSTEM)
- async def on_intent_handled(self, event_name: str, event_data: Any) -> None:
- """Intent verarbeitet."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if profile:
- profile.timestamps["intent_handled"] = time.monotonic()
- @TrixyEvent(["create_output_text"], priority=EventPriority.SYSTEM)
- async def on_create_output_text(self, event_name: str, event_data: Any) -> None:
- """LLM-Textgenerierung angefordert."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if profile:
- profile.timestamps["create_output_text"] = time.monotonic()
- @TrixyEvent(["output_text_created"], priority=EventPriority.SYSTEM)
- async def on_output_text_created(self, event_name: str, event_data: Any) -> None:
- """Antworttext generiert."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if profile:
- profile.timestamps["output_text_created"] = time.monotonic()
- @TrixyEvent(["tts_completed"], priority=EventPriority.SYSTEM)
- async def on_tts_completed(self, event_name: str, event_data: Any) -> None:
- """TTS abgeschlossen — Profil ausgeben (letztes Event der Pipeline)."""
- if not self._application.debug:
- return
- sat_id = self._get_sat_id(event_data)
- profile = self._get_profile(sat_id)
- if not profile:
- return
- profile.timestamps["tts_completed"] = time.monotonic()
- provider = str(_get_field(event_data, "provider", ""))
- if provider:
- profile.providers["TTS"] = provider
- # Profil ausgeben und aufraeumen
- self._print_profile(profile)
- self._profiles.pop(sat_id, None)
- # === Ausgabe ===
- def _print_profile(self, profile: SessionProfile) -> None:
- """Gibt das Profil als formatierte Tabelle aus."""
- ts = profile.timestamps
- # Phasen berechnen
- phases: list[tuple[str, str, float]] = []
- if profile.is_keyboard:
- # Keyboard-Input: Nur Phasen ab NLP (kein Wakeword/Recording/STT)
- phase_list = _PHASES[4:]
- else:
- phase_list = _PHASES
- for phase_name, start_evt, end_evt, default_component in phase_list:
- if start_evt in ts and end_evt in ts:
- duration = ts[end_evt] - ts[start_evt]
- component = profile.providers.get(phase_name, default_component)
- phases.append((phase_name, component, duration))
- if not phases:
- return
- # Gesamt-Dauer
- all_ts = sorted(ts.values())
- total_duration = all_ts[-1] - all_ts[0] if len(all_ts) >= 2 else 0.0
- # Tabelle formatieren
- w_phase = 16
- w_comp = 16
- w_dur = 19
- w_inner = w_phase + 1 + w_comp + 1 + w_dur
- # Session-ID kuerzen fuer Titel
- display_id = profile.session_id or profile.satellite_id
- short_id = display_id[:16] if len(display_id) > 16 else display_id
- lines: list[str] = []
- lines.append(f"╔{'═' * w_inner}╗")
- title = f"Conversation Profil — {short_id}"
- lines.append(f"║{title:^{w_inner}}║")
- lines.append(f"╠{'═' * w_inner}╣")
- # Header
- lines.append(
- f"║ {'Phase':<{w_phase - 1}}│ {'Komponente':<{w_comp - 1}}│ {'Dauer':<{w_dur - 1}}║"
- )
- lines.append(
- f"╠{'═' * w_phase}╪{'═' * w_comp}╪{'═' * w_dur}╣"
- )
- # Phasen-Zeilen
- for phase_name, component, duration in phases:
- dur_str = f"{duration:>7.3f} sec"
- lines.append(
- f"║ {phase_name:<{w_phase - 1}}│ {component:<{w_comp - 1}}│ {dur_str:<{w_dur - 1}}║"
- )
- # Gesamt
- lines.append(
- f"╠{'═' * w_phase}╪{'═' * w_comp}╪{'═' * w_dur}╣"
- )
- total_str = f"{total_duration:>7.3f} sec"
- lines.append(
- f"║ {'GESAMT':<{w_phase - 1}}│ {'':<{w_comp - 1}}│ {total_str:<{w_dur - 1}}║"
- )
- # Meta-Informationen
- lines.append(
- f"╠{'═' * w_phase}╧{'═' * w_comp}╧{'═' * w_dur}╣"
- )
- if profile.text:
- text_display = profile.text[:w_inner - 9]
- lines.append(f"║ Text: \"{text_display}\"{' ' * max(0, w_inner - 9 - len(text_display))}║")
- if profile.intent:
- intent_info = f"Intent: {profile.intent}"
- if profile.confidence:
- intent_info += f" (confidence: {profile.confidence:.2f})"
- intent_display = intent_info[:w_inner - 2]
- lines.append(f"║ {intent_display:<{w_inner - 2}} ║")
- if profile.satellite_id:
- sat_display = f"Satellite: {profile.satellite_id}"[:w_inner - 2]
- lines.append(f"║ {sat_display:<{w_inner - 2}} ║")
- if profile.is_keyboard:
- lines.append(f"║ {'Quelle: Keyboard-Input':<{w_inner - 2}} ║")
- lines.append(f"╚{'═' * w_inner}╝")
- # Ausgabe ueber pinfo()
- for line in lines:
- pinfo(line)
|