# -*- coding: utf-8 -*- """ LLM-basiertes NLP Plugin für Trixy. Verantwortlich für: 1. Intent-Erkennung: speech_recognized → intent_received 2. Antwort-Generierung: create_output_text → output_text_created Event-Flow: speech_recognized ↓ [NLP Plugin] Intent-Erkennung ↓ intent_received ↓ [Intent Dispatcher] Handler aufrufen ↓ intent_handled ↓ [NLP Plugin] Antwort generieren (falls nötig) ↓ output_text_created ↓ [TTS Plugin] Sprachausgabe """ import asyncio import json import time import uuid from typing import Any from trixy_core.plugins.trixy_plugin import TrixyPlugin from trixy_core.plugins.extensions.points.conversation import NLPExtension from trixy_core.events.decorators import TrixyEvent from trixy_core.events.event_data.basic import ( SpeechRecognized, IntentReceived, IntentHandled, CreateOutputText, OutputTextCreated, PluginLoaded, PluginUnloaded, ) from trixy_core.nlp import ( NLPConfig, NLPContext, NLPResult, IntentRegistry, get_system_intents_for_context, is_admin_intent, ) from trixy_core.nlp.keyword_matcher import KeywordIntentMatcher from trixy_core.nlp.entities import EntityRegistry from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn # Cache-Refresh-Intervall in Sekunden INTENT_CACHE_REFRESH_INTERVAL = 600 # 10 Minuten class LLMNLPExtension(NLPExtension): """NLP Extension für LLM-basierte Verarbeitung.""" def __init__(self, plugin: "LLMNLPPlugin") -> None: super().__init__( extension_id="nlp_llm", name="LLM NLP", plugin_name=plugin.name, engine="llm", capabilities=["intent", "response", "entity"], offline=True, description="Lokales LLM für Intent-Erkennung und Antwortgenerierung", ) self._plugin = plugin async def create_instance(self, config: dict[str, Any]): return self._plugin.provider class IntentCache: """Cache für verfügbare Intents.""" def __init__(self) -> None: self._plugin_intents: list[dict[str, Any]] = [] self._last_refresh: float = 0 self._cached_prompt_intents: dict[str, str] = {} def refresh_plugin_intents(self) -> None: """Lädt Plugin-Intents aus der Registry neu.""" registry = IntentRegistry.get_instance() self._plugin_intents = registry.get_all_as_dict() self._cached_prompt_intents.clear() self._last_refresh = time.time() pdebug(f"Intent-Cache aktualisiert: {len(self._plugin_intents)} Plugin-Intents") def needs_refresh(self) -> bool: return (time.time() - self._last_refresh) > INTENT_CACHE_REFRESH_INTERVAL def get_intents_for_context( self, wakeword_type: str = "custom", is_authenticated: bool = False, ) -> list[dict[str, Any]]: """Gibt alle verfügbaren Intents für den Kontext zurück.""" system_intents = get_system_intents_for_context( wakeword_type=wakeword_type, is_authenticated=is_authenticated, ) return system_intents + self._plugin_intents @property def plugin_intent_count(self) -> int: return len(self._plugin_intents) class LLMNLPPlugin(TrixyPlugin): """ LLM-basiertes NLP Plugin. Aufgaben: 1. Intent-Erkennung aus speech_recognized 2. Antwort-Generierung für create_output_text """ NAME = "nlp_llm" VERSION = "2.0.0" DESCRIPTION = "LLM-basierte NLP (Event-Driven)" AUTHOR = "Trixy" def __init__(self, application, plugin_path, config) -> None: super().__init__(application, plugin_path, config) self._provider = None self._extension = None self._intent_cache = IntentCache() self._keyword_matcher: KeywordIntentMatcher | None = None self._entity_registry: EntityRegistry | None = None self._refresh_task: asyncio.Task | None = None self._admin_sessions: dict[str, dict[str, Any]] = {} @property def provider(self): return self._provider async def on_load(self) -> None: """Initialisiert das Plugin.""" pinfo(f"Lade {self.NAME} Plugin...") backend = self.get_config_value("backend", "llama_cpp") if backend == "ollama": from plugins.nlp_llm.backends.ollama import OllamaNLPBackend self._provider = OllamaNLPBackend() else: from plugins.nlp_llm.backends.llama_cpp import LlamaCppNLPBackend self._provider = LlamaCppNLPBackend() nlp_config = NLPConfig( backend=backend, model_name=self.get_config_value("model_name", "qwen2.5-1.5b-instruct"), model_path=self.get_config_value("model_path"), use_gpu=self.get_config_value("use_gpu", False), num_threads=self.get_config_value("num_threads", 4), temperature=self.get_config_value("temperature", 0.1), max_tokens=self.get_config_value("max_tokens", 256), context_window=self.get_config_value("context_window", 2048), extra={ "ollama_host": self.get_config_value("ollama_host", "http://localhost:11434"), "response_language": self.get_config_value("response_language", "de"), } ) success = await self._provider.initialize(nlp_config) if success: pinfo(f"[NLP] Provider initialisiert: {backend}, model={nlp_config.model_name}, ready={self._provider.is_ready}") else: perror(f"[NLP] Provider-Initialisierung fehlgeschlagen: {backend}") self._intent_cache.refresh_plugin_intents() # Entity-Registry initialisieren from pathlib import Path self._entity_registry = EntityRegistry() project_root = Path(__file__).resolve().parent.parent.parent entities_dir = project_root / "entities" plugins_dir = project_root / "plugins" language = self.get_config_value("response_language", nlp_config.extra.get("response_language", "de")) # Globale Entities laden (entities/{lang}/*.yml|*.json) self._entity_registry.load(entities_dir, language) # Plugin-Entities laden (plugins/*/entities/{lang}/*.yml|*.json) self._entity_registry.load_plugin_entities(plugins_dir, language) pinfo(f"[NLP] Entity-Registry geladen: {len(self._entity_registry.entity_types)} Typen") # Keyword-Matcher initialisieren (Pre-Filter vor LLM) matcher_config = self.get_config_value("keyword_matcher", {}) if matcher_config.get("enabled", True): self._keyword_matcher = KeywordIntentMatcher(matcher_config) self._keyword_matcher.set_entity_registry(self._entity_registry) self._refresh_keyword_matcher() pinfo("[NLP] Keyword-Matcher aktiviert (mit Entity-Validierung)") else: self._keyword_matcher = None self._refresh_task = asyncio.create_task(self._periodic_cache_refresh()) self._extension = LLMNLPExtension(self) # Extension nicht registrieren - das Plugin arbeitet direkt event-driven pinfo(f"{self.NAME} Plugin geladen (Event-Driven Mode)") async def on_unload(self) -> None: """Fährt das Plugin herunter.""" if self._refresh_task: self._refresh_task.cancel() try: await self._refresh_task except asyncio.CancelledError: pass if self._provider: await self._provider.shutdown() self._provider = None self._extension = None def _refresh_keyword_matcher(self) -> None: """Aktualisiert den Keyword-Matcher mit allen verfuegbaren Intents.""" if not self._keyword_matcher: return all_intents = self._intent_cache.get_intents_for_context() self._keyword_matcher.compile_intents(all_intents) async def _periodic_cache_refresh(self) -> None: """Periodischer Cache-Refresh.""" while True: await asyncio.sleep(INTENT_CACHE_REFRESH_INTERVAL) self._intent_cache.refresh_plugin_intents() self._refresh_keyword_matcher() # ========================================================================= # Event Handler: Plugin Changes # ========================================================================= @TrixyEvent(["plugin_loaded", "plugin_unloaded"]) async def on_plugin_change(self, event_name: str, event_data) -> None: """Aktualisiert Cache und Keyword-Matcher bei Plugin-Änderungen.""" plugin_name = getattr(event_data, "plugin_name", "") if plugin_name != self.NAME: self._intent_cache.refresh_plugin_intents() self._refresh_keyword_matcher() # ========================================================================= # Event Handler: Intent-Erkennung (speech_recognized → intent_received) # ========================================================================= @TrixyEvent(["speech_recognized"]) async def on_speech_recognized(self, event_name: str, event_data: SpeechRecognized) -> None: """ Erkennt Intent aus Spracheingabe. Emittiert: intent_received """ pinfo(f"[NLP] speech_recognized empfangen: '{getattr(event_data, 'text', '?')}'") if not self._provider: pwarn("[NLP] Kein Provider konfiguriert") return if not self._provider.is_ready: pwarn(f"[NLP] Provider nicht bereit (state={getattr(self._provider, '_state', '?')})") return text = event_data.text if not text or not text.strip(): pdebug("[NLP] Leerer Text, überspringe") return # STT-Korrektur: Häufige Fehlerkennungen korrigieren text = self._correct_stt_text(text) pinfo(f"[NLP] Starte Intent-Erkennung für: '{text}'") # Session-Info holen session_info = self._get_session_info(event_data.satellite_id) wakeword_type = session_info.get("wakeword_type", "custom") is_authenticated = self._is_authenticated(event_data.satellite_id) # Cache aktualisieren falls nötig if self._intent_cache.needs_refresh(): self._intent_cache.refresh_plugin_intents() available_intents = self._intent_cache.get_intents_for_context( wakeword_type=wakeword_type, is_authenticated=is_authenticated, ) # Keyword-Matcher Pre-Filter (vor LLM) import time as _time metadata: dict[str, Any] = {} if self._keyword_matcher: t_kw_start = _time.monotonic() kw_match = self._keyword_matcher.match(text) t_kw_elapsed = (_time.monotonic() - t_kw_start) * 1000 # ms if kw_match: pinfo( f"[NLP] Keyword-Match: '{kw_match.intent}' " f"(confidence={kw_match.confidence:.2f}, " f"{kw_match.matched_by}, {t_kw_elapsed:.1f}ms)" ) result = NLPResult( intent=kw_match.intent, confidence=kw_match.confidence, slots=kw_match.slots, success=True, processing_time=t_kw_elapsed / 1000, ) if kw_match.sentiment: metadata["sentiment"] = kw_match.sentiment else: pdebug(f"[NLP] Kein Keyword-Match ({t_kw_elapsed:.1f}ms), Fallback auf LLM") kw_match = None else: kw_match = None # LLM-Fallback (nur wenn kein Keyword-Match) if kw_match is None: pinfo(f"[NLP] {len(available_intents)} Intents verfügbar, sende an LLM...") t_start = _time.monotonic() result = await self._recognize_intent( text=text, satellite_id=event_data.satellite_id, session_info=session_info, wakeword_type=wakeword_type, is_authenticated=is_authenticated, language=event_data.language or "de", ) t_elapsed = _time.monotonic() - t_start if not result.success: perror(f"[NLP] Intent-Erkennung fehlgeschlagen nach {t_elapsed:.1f}s: {result.error}") return pinfo(f"[NLP] Intent erkannt in {t_elapsed:.1f}s: intent='{result.intent}', confidence={result.confidence:.2f}, slots={result.slots}") else: t_elapsed = result.processing_time # Sicherheitsprüfung if not self._check_authorization(result.intent, wakeword_type, is_authenticated): pwarn(f"[NLP] Autorisierung verweigert für Intent: {result.intent}") await self._emit_authorization_error(event_data.satellite_id, session_info) return # Konfidenz prüfen min_confidence = self.get_config_value("min_confidence", 0.3) intent = result.intent if result.confidence >= min_confidence else "unknown" if intent == "unknown" and result.intent != "unknown": pdebug(f"[NLP] Konfidenz zu niedrig ({result.confidence:.2f} < {min_confidence}), Intent -> 'unknown'") # intent_received emittieren intent_event = IntentReceived( satellite_id=event_data.satellite_id, intent=intent, confidence=result.confidence, slots=result.slots, original_text=text, session_id=session_info.get("session_id", ""), room_id=session_info.get("room_id", ""), wakeword_type=wakeword_type, is_authenticated=is_authenticated, language=event_data.language or "de", metadata=metadata, ) pinfo(f"[NLP] Emittiere intent_received: '{intent}' (confidence={result.confidence:.2f})") await self._application.events.trigger("intent_received", intent_event) async def _recognize_intent( self, text: str, satellite_id: str, session_info: dict[str, Any], wakeword_type: str, is_authenticated: bool, language: str, ): """Führt Intent-Erkennung durch.""" available_intents = self._intent_cache.get_intents_for_context( wakeword_type=wakeword_type, is_authenticated=is_authenticated, ) context = NLPContext( text=text, satellite_id=satellite_id, room_id=session_info.get("room_id", ""), session_id=session_info.get("session_id", ""), available_intents=available_intents, language=language, ) return await self._provider.process(context) # ========================================================================= # Event Handler: Antwort-Generierung (create_output_text → output_text_created) # ========================================================================= @TrixyEvent(["create_output_text"]) async def on_create_output_text(self, event_name: str, event_data: CreateOutputText) -> None: """ Generiert Antworttext aus Handler-Ergebnis. Emittiert: output_text_created """ pinfo(f"[NLP] create_output_text empfangen: intent='{event_data.intent}', text='{event_data.original_text[:50]}'") if not self._provider or not self._provider.is_ready: pwarn("[NLP] Provider nicht bereit für Antwort-Generierung") return pinfo(f"[NLP] Generiere Antwort für Intent: '{event_data.intent}'...") # Antwort generieren response_text = await self._provider.generate_response( intent=event_data.intent, slots=event_data.slots, handler_result=event_data.handler_data, original_text=event_data.original_text, room_id=event_data.room_id, ) if not response_text: # Fallback für unbekannte Intents response_text = await self._provider.generate_fallback_response( event_data.original_text ) # output_text_created emittieren output_event = OutputTextCreated( satellite_id=event_data.satellite_id, session_id=event_data.session_id, room_id=event_data.room_id, text=response_text, intent=event_data.intent, is_followup=False, expects_response=False, ) pdebug(f"Emittiere output_text_created: {response_text[:50]}...") await self._application.events.trigger("output_text_created", output_event) # ========================================================================= # Hilfsmethoden # ========================================================================= # STT-Korrekturen: Wort-Aliase für häufige Fehlerkennungen. # Nur angewendet wenn das Wort alleine steht (gesamter Text). _STT_WORD_ALIASES: dict[str, str] = { "shop": "stopp", "schop": "stopp", } def _correct_stt_text(self, text: str) -> str: """Korrigiert häufige STT-Fehlerkennungen.""" stripped = text.strip().lower() # Einzelwort-Aliase (nur wenn das Wort alleine steht) if stripped in self._STT_WORD_ALIASES: corrected = self._STT_WORD_ALIASES[stripped] pdebug(f"[NLP] STT-Korrektur: '{text}' → '{corrected}'") return corrected return text def _get_session_info(self, satellite_id: str) -> dict[str, Any]: """Holt Session-Info.""" info = {"session_id": "", "room_id": "", "wakeword_type": "custom"} if hasattr(self._application, "services"): try: conv_service = self._application.services.get_service("ConversationService") if conv_service: session = conv_service.get_active_session(satellite_id) if session: info["session_id"] = session.session_id info["room_id"] = getattr(session, "room_id", "") info["wakeword_type"] = getattr(session, "wakeword_type", "custom") except Exception: pass return info def _is_authenticated(self, satellite_id: str) -> bool: """Prüft Admin-Authentifizierung.""" session = self._admin_sessions.get(satellite_id) if not session: return False if time.time() > session.get("expires", 0): del self._admin_sessions[satellite_id] return False return session.get("authenticated", False) def _check_authorization( self, intent: str, wakeword_type: str, is_authenticated: bool ) -> bool: """Prüft Intent-Autorisierung.""" if not is_admin_intent(intent): return True if wakeword_type != "system_command": return False if intent == "system_login": return True return is_authenticated async def _emit_authorization_error( self, satellite_id: str, session_info: dict[str, Any] ) -> None: """Emittiert Autorisierungs-Fehler.""" error_event = OutputTextCreated( satellite_id=satellite_id, session_id=session_info.get("session_id", ""), room_id=session_info.get("room_id", ""), text="Dieser Befehl erfordert Administrator-Rechte. " "Bitte verwende das System-Wakeword und melde dich an.", intent="authorization_error", ) await self._application.events.trigger("output_text_created", error_event)