| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531 |
- # -*- 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)
|