main.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531
  1. # -*- coding: utf-8 -*-
  2. """
  3. LLM-basiertes NLP Plugin für Trixy.
  4. Verantwortlich für:
  5. 1. Intent-Erkennung: speech_recognized → intent_received
  6. 2. Antwort-Generierung: create_output_text → output_text_created
  7. Event-Flow:
  8. speech_recognized
  9. [NLP Plugin] Intent-Erkennung
  10. intent_received
  11. [Intent Dispatcher] Handler aufrufen
  12. intent_handled
  13. [NLP Plugin] Antwort generieren (falls nötig)
  14. output_text_created
  15. [TTS Plugin] Sprachausgabe
  16. """
  17. import asyncio
  18. import json
  19. import time
  20. import uuid
  21. from typing import Any
  22. from trixy_core.plugins.trixy_plugin import TrixyPlugin
  23. from trixy_core.plugins.extensions.points.conversation import NLPExtension
  24. from trixy_core.events.decorators import TrixyEvent
  25. from trixy_core.events.event_data.basic import (
  26. SpeechRecognized,
  27. IntentReceived,
  28. IntentHandled,
  29. CreateOutputText,
  30. OutputTextCreated,
  31. PluginLoaded,
  32. PluginUnloaded,
  33. )
  34. from trixy_core.nlp import (
  35. NLPConfig,
  36. NLPContext,
  37. NLPResult,
  38. IntentRegistry,
  39. get_system_intents_for_context,
  40. is_admin_intent,
  41. )
  42. from trixy_core.nlp.keyword_matcher import KeywordIntentMatcher
  43. from trixy_core.nlp.entities import EntityRegistry
  44. from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
  45. # Cache-Refresh-Intervall in Sekunden
  46. INTENT_CACHE_REFRESH_INTERVAL = 600 # 10 Minuten
  47. class LLMNLPExtension(NLPExtension):
  48. """NLP Extension für LLM-basierte Verarbeitung."""
  49. def __init__(self, plugin: "LLMNLPPlugin") -> None:
  50. super().__init__(
  51. extension_id="nlp_llm",
  52. name="LLM NLP",
  53. plugin_name=plugin.name,
  54. engine="llm",
  55. capabilities=["intent", "response", "entity"],
  56. offline=True,
  57. description="Lokales LLM für Intent-Erkennung und Antwortgenerierung",
  58. )
  59. self._plugin = plugin
  60. async def create_instance(self, config: dict[str, Any]):
  61. return self._plugin.provider
  62. class IntentCache:
  63. """Cache für verfügbare Intents."""
  64. def __init__(self) -> None:
  65. self._plugin_intents: list[dict[str, Any]] = []
  66. self._last_refresh: float = 0
  67. self._cached_prompt_intents: dict[str, str] = {}
  68. def refresh_plugin_intents(self) -> None:
  69. """Lädt Plugin-Intents aus der Registry neu."""
  70. registry = IntentRegistry.get_instance()
  71. self._plugin_intents = registry.get_all_as_dict()
  72. self._cached_prompt_intents.clear()
  73. self._last_refresh = time.time()
  74. pdebug(f"Intent-Cache aktualisiert: {len(self._plugin_intents)} Plugin-Intents")
  75. def needs_refresh(self) -> bool:
  76. return (time.time() - self._last_refresh) > INTENT_CACHE_REFRESH_INTERVAL
  77. def get_intents_for_context(
  78. self,
  79. wakeword_type: str = "custom",
  80. is_authenticated: bool = False,
  81. ) -> list[dict[str, Any]]:
  82. """Gibt alle verfügbaren Intents für den Kontext zurück."""
  83. system_intents = get_system_intents_for_context(
  84. wakeword_type=wakeword_type,
  85. is_authenticated=is_authenticated,
  86. )
  87. return system_intents + self._plugin_intents
  88. @property
  89. def plugin_intent_count(self) -> int:
  90. return len(self._plugin_intents)
  91. class LLMNLPPlugin(TrixyPlugin):
  92. """
  93. LLM-basiertes NLP Plugin.
  94. Aufgaben:
  95. 1. Intent-Erkennung aus speech_recognized
  96. 2. Antwort-Generierung für create_output_text
  97. """
  98. NAME = "nlp_llm"
  99. VERSION = "2.0.0"
  100. DESCRIPTION = "LLM-basierte NLP (Event-Driven)"
  101. AUTHOR = "Trixy"
  102. def __init__(self, application, plugin_path, config) -> None:
  103. super().__init__(application, plugin_path, config)
  104. self._provider = None
  105. self._extension = None
  106. self._intent_cache = IntentCache()
  107. self._keyword_matcher: KeywordIntentMatcher | None = None
  108. self._entity_registry: EntityRegistry | None = None
  109. self._refresh_task: asyncio.Task | None = None
  110. self._admin_sessions: dict[str, dict[str, Any]] = {}
  111. @property
  112. def provider(self):
  113. return self._provider
  114. async def on_load(self) -> None:
  115. """Initialisiert das Plugin."""
  116. pinfo(f"Lade {self.NAME} Plugin...")
  117. backend = self.get_config_value("backend", "llama_cpp")
  118. if backend == "ollama":
  119. from plugins.nlp_llm.backends.ollama import OllamaNLPBackend
  120. self._provider = OllamaNLPBackend()
  121. else:
  122. from plugins.nlp_llm.backends.llama_cpp import LlamaCppNLPBackend
  123. self._provider = LlamaCppNLPBackend()
  124. nlp_config = NLPConfig(
  125. backend=backend,
  126. model_name=self.get_config_value("model_name", "qwen2.5-1.5b-instruct"),
  127. model_path=self.get_config_value("model_path"),
  128. use_gpu=self.get_config_value("use_gpu", False),
  129. num_threads=self.get_config_value("num_threads", 4),
  130. temperature=self.get_config_value("temperature", 0.1),
  131. max_tokens=self.get_config_value("max_tokens", 256),
  132. context_window=self.get_config_value("context_window", 2048),
  133. extra={
  134. "ollama_host": self.get_config_value("ollama_host", "http://localhost:11434"),
  135. "response_language": self.get_config_value("response_language", "de"),
  136. }
  137. )
  138. success = await self._provider.initialize(nlp_config)
  139. if success:
  140. pinfo(f"[NLP] Provider initialisiert: {backend}, model={nlp_config.model_name}, ready={self._provider.is_ready}")
  141. else:
  142. perror(f"[NLP] Provider-Initialisierung fehlgeschlagen: {backend}")
  143. self._intent_cache.refresh_plugin_intents()
  144. # Entity-Registry initialisieren
  145. from pathlib import Path
  146. self._entity_registry = EntityRegistry()
  147. project_root = Path(__file__).resolve().parent.parent.parent
  148. entities_dir = project_root / "entities"
  149. plugins_dir = project_root / "plugins"
  150. language = self.get_config_value("response_language",
  151. nlp_config.extra.get("response_language", "de"))
  152. # Globale Entities laden (entities/{lang}/*.yml|*.json)
  153. self._entity_registry.load(entities_dir, language)
  154. # Plugin-Entities laden (plugins/*/entities/{lang}/*.yml|*.json)
  155. self._entity_registry.load_plugin_entities(plugins_dir, language)
  156. pinfo(f"[NLP] Entity-Registry geladen: {len(self._entity_registry.entity_types)} Typen")
  157. # Keyword-Matcher initialisieren (Pre-Filter vor LLM)
  158. matcher_config = self.get_config_value("keyword_matcher", {})
  159. if matcher_config.get("enabled", True):
  160. self._keyword_matcher = KeywordIntentMatcher(matcher_config)
  161. self._keyword_matcher.set_entity_registry(self._entity_registry)
  162. self._refresh_keyword_matcher()
  163. pinfo("[NLP] Keyword-Matcher aktiviert (mit Entity-Validierung)")
  164. else:
  165. self._keyword_matcher = None
  166. self._refresh_task = asyncio.create_task(self._periodic_cache_refresh())
  167. self._extension = LLMNLPExtension(self)
  168. # Extension nicht registrieren - das Plugin arbeitet direkt event-driven
  169. pinfo(f"{self.NAME} Plugin geladen (Event-Driven Mode)")
  170. async def on_unload(self) -> None:
  171. """Fährt das Plugin herunter."""
  172. if self._refresh_task:
  173. self._refresh_task.cancel()
  174. try:
  175. await self._refresh_task
  176. except asyncio.CancelledError:
  177. pass
  178. if self._provider:
  179. await self._provider.shutdown()
  180. self._provider = None
  181. self._extension = None
  182. def _refresh_keyword_matcher(self) -> None:
  183. """Aktualisiert den Keyword-Matcher mit allen verfuegbaren Intents."""
  184. if not self._keyword_matcher:
  185. return
  186. all_intents = self._intent_cache.get_intents_for_context()
  187. self._keyword_matcher.compile_intents(all_intents)
  188. async def _periodic_cache_refresh(self) -> None:
  189. """Periodischer Cache-Refresh."""
  190. while True:
  191. await asyncio.sleep(INTENT_CACHE_REFRESH_INTERVAL)
  192. self._intent_cache.refresh_plugin_intents()
  193. self._refresh_keyword_matcher()
  194. # =========================================================================
  195. # Event Handler: Plugin Changes
  196. # =========================================================================
  197. @TrixyEvent(["plugin_loaded", "plugin_unloaded"])
  198. async def on_plugin_change(self, event_name: str, event_data) -> None:
  199. """Aktualisiert Cache und Keyword-Matcher bei Plugin-Änderungen."""
  200. plugin_name = getattr(event_data, "plugin_name", "")
  201. if plugin_name != self.NAME:
  202. self._intent_cache.refresh_plugin_intents()
  203. self._refresh_keyword_matcher()
  204. # =========================================================================
  205. # Event Handler: Intent-Erkennung (speech_recognized → intent_received)
  206. # =========================================================================
  207. @TrixyEvent(["speech_recognized"])
  208. async def on_speech_recognized(self, event_name: str, event_data: SpeechRecognized) -> None:
  209. """
  210. Erkennt Intent aus Spracheingabe.
  211. Emittiert: intent_received
  212. """
  213. pinfo(f"[NLP] speech_recognized empfangen: '{getattr(event_data, 'text', '?')}'")
  214. if not self._provider:
  215. pwarn("[NLP] Kein Provider konfiguriert")
  216. return
  217. if not self._provider.is_ready:
  218. pwarn(f"[NLP] Provider nicht bereit (state={getattr(self._provider, '_state', '?')})")
  219. return
  220. text = event_data.text
  221. if not text or not text.strip():
  222. pdebug("[NLP] Leerer Text, überspringe")
  223. return
  224. # STT-Korrektur: Häufige Fehlerkennungen korrigieren
  225. text = self._correct_stt_text(text)
  226. pinfo(f"[NLP] Starte Intent-Erkennung für: '{text}'")
  227. # Session-Info holen
  228. session_info = self._get_session_info(event_data.satellite_id)
  229. wakeword_type = session_info.get("wakeword_type", "custom")
  230. is_authenticated = self._is_authenticated(event_data.satellite_id)
  231. # Cache aktualisieren falls nötig
  232. if self._intent_cache.needs_refresh():
  233. self._intent_cache.refresh_plugin_intents()
  234. available_intents = self._intent_cache.get_intents_for_context(
  235. wakeword_type=wakeword_type,
  236. is_authenticated=is_authenticated,
  237. )
  238. # Keyword-Matcher Pre-Filter (vor LLM)
  239. import time as _time
  240. metadata: dict[str, Any] = {}
  241. if self._keyword_matcher:
  242. t_kw_start = _time.monotonic()
  243. kw_match = self._keyword_matcher.match(text)
  244. t_kw_elapsed = (_time.monotonic() - t_kw_start) * 1000 # ms
  245. if kw_match:
  246. pinfo(
  247. f"[NLP] Keyword-Match: '{kw_match.intent}' "
  248. f"(confidence={kw_match.confidence:.2f}, "
  249. f"{kw_match.matched_by}, {t_kw_elapsed:.1f}ms)"
  250. )
  251. result = NLPResult(
  252. intent=kw_match.intent,
  253. confidence=kw_match.confidence,
  254. slots=kw_match.slots,
  255. success=True,
  256. processing_time=t_kw_elapsed / 1000,
  257. )
  258. if kw_match.sentiment:
  259. metadata["sentiment"] = kw_match.sentiment
  260. else:
  261. pdebug(f"[NLP] Kein Keyword-Match ({t_kw_elapsed:.1f}ms), Fallback auf LLM")
  262. kw_match = None
  263. else:
  264. kw_match = None
  265. # LLM-Fallback (nur wenn kein Keyword-Match)
  266. if kw_match is None:
  267. pinfo(f"[NLP] {len(available_intents)} Intents verfügbar, sende an LLM...")
  268. t_start = _time.monotonic()
  269. result = await self._recognize_intent(
  270. text=text,
  271. satellite_id=event_data.satellite_id,
  272. session_info=session_info,
  273. wakeword_type=wakeword_type,
  274. is_authenticated=is_authenticated,
  275. language=event_data.language or "de",
  276. )
  277. t_elapsed = _time.monotonic() - t_start
  278. if not result.success:
  279. perror(f"[NLP] Intent-Erkennung fehlgeschlagen nach {t_elapsed:.1f}s: {result.error}")
  280. return
  281. pinfo(f"[NLP] Intent erkannt in {t_elapsed:.1f}s: intent='{result.intent}', confidence={result.confidence:.2f}, slots={result.slots}")
  282. else:
  283. t_elapsed = result.processing_time
  284. # Sicherheitsprüfung
  285. if not self._check_authorization(result.intent, wakeword_type, is_authenticated):
  286. pwarn(f"[NLP] Autorisierung verweigert für Intent: {result.intent}")
  287. await self._emit_authorization_error(event_data.satellite_id, session_info)
  288. return
  289. # Konfidenz prüfen
  290. min_confidence = self.get_config_value("min_confidence", 0.3)
  291. intent = result.intent if result.confidence >= min_confidence else "unknown"
  292. if intent == "unknown" and result.intent != "unknown":
  293. pdebug(f"[NLP] Konfidenz zu niedrig ({result.confidence:.2f} < {min_confidence}), Intent -> 'unknown'")
  294. # intent_received emittieren
  295. intent_event = IntentReceived(
  296. satellite_id=event_data.satellite_id,
  297. intent=intent,
  298. confidence=result.confidence,
  299. slots=result.slots,
  300. original_text=text,
  301. session_id=session_info.get("session_id", ""),
  302. room_id=session_info.get("room_id", ""),
  303. wakeword_type=wakeword_type,
  304. is_authenticated=is_authenticated,
  305. language=event_data.language or "de",
  306. metadata=metadata,
  307. )
  308. pinfo(f"[NLP] Emittiere intent_received: '{intent}' (confidence={result.confidence:.2f})")
  309. await self._application.events.trigger("intent_received", intent_event)
  310. async def _recognize_intent(
  311. self,
  312. text: str,
  313. satellite_id: str,
  314. session_info: dict[str, Any],
  315. wakeword_type: str,
  316. is_authenticated: bool,
  317. language: str,
  318. ):
  319. """Führt Intent-Erkennung durch."""
  320. available_intents = self._intent_cache.get_intents_for_context(
  321. wakeword_type=wakeword_type,
  322. is_authenticated=is_authenticated,
  323. )
  324. context = NLPContext(
  325. text=text,
  326. satellite_id=satellite_id,
  327. room_id=session_info.get("room_id", ""),
  328. session_id=session_info.get("session_id", ""),
  329. available_intents=available_intents,
  330. language=language,
  331. )
  332. return await self._provider.process(context)
  333. # =========================================================================
  334. # Event Handler: Antwort-Generierung (create_output_text → output_text_created)
  335. # =========================================================================
  336. @TrixyEvent(["create_output_text"])
  337. async def on_create_output_text(self, event_name: str, event_data: CreateOutputText) -> None:
  338. """
  339. Generiert Antworttext aus Handler-Ergebnis.
  340. Emittiert: output_text_created
  341. """
  342. pinfo(f"[NLP] create_output_text empfangen: intent='{event_data.intent}', text='{event_data.original_text[:50]}'")
  343. if not self._provider or not self._provider.is_ready:
  344. pwarn("[NLP] Provider nicht bereit für Antwort-Generierung")
  345. return
  346. pinfo(f"[NLP] Generiere Antwort für Intent: '{event_data.intent}'...")
  347. # Antwort generieren
  348. response_text = await self._provider.generate_response(
  349. intent=event_data.intent,
  350. slots=event_data.slots,
  351. handler_result=event_data.handler_data,
  352. original_text=event_data.original_text,
  353. room_id=event_data.room_id,
  354. )
  355. if not response_text:
  356. # Fallback für unbekannte Intents
  357. response_text = await self._provider.generate_fallback_response(
  358. event_data.original_text
  359. )
  360. # output_text_created emittieren
  361. output_event = OutputTextCreated(
  362. satellite_id=event_data.satellite_id,
  363. session_id=event_data.session_id,
  364. room_id=event_data.room_id,
  365. text=response_text,
  366. intent=event_data.intent,
  367. is_followup=False,
  368. expects_response=False,
  369. )
  370. pdebug(f"Emittiere output_text_created: {response_text[:50]}...")
  371. await self._application.events.trigger("output_text_created", output_event)
  372. # =========================================================================
  373. # Hilfsmethoden
  374. # =========================================================================
  375. # STT-Korrekturen: Wort-Aliase für häufige Fehlerkennungen.
  376. # Nur angewendet wenn das Wort alleine steht (gesamter Text).
  377. _STT_WORD_ALIASES: dict[str, str] = {
  378. "shop": "stopp",
  379. "schop": "stopp",
  380. }
  381. def _correct_stt_text(self, text: str) -> str:
  382. """Korrigiert häufige STT-Fehlerkennungen."""
  383. stripped = text.strip().lower()
  384. # Einzelwort-Aliase (nur wenn das Wort alleine steht)
  385. if stripped in self._STT_WORD_ALIASES:
  386. corrected = self._STT_WORD_ALIASES[stripped]
  387. pdebug(f"[NLP] STT-Korrektur: '{text}' → '{corrected}'")
  388. return corrected
  389. return text
  390. def _get_session_info(self, satellite_id: str) -> dict[str, Any]:
  391. """Holt Session-Info."""
  392. info = {"session_id": "", "room_id": "", "wakeword_type": "custom"}
  393. if hasattr(self._application, "services"):
  394. try:
  395. conv_service = self._application.services.get_service("ConversationService")
  396. if conv_service:
  397. session = conv_service.get_active_session(satellite_id)
  398. if session:
  399. info["session_id"] = session.session_id
  400. info["room_id"] = getattr(session, "room_id", "")
  401. info["wakeword_type"] = getattr(session, "wakeword_type", "custom")
  402. except Exception:
  403. pass
  404. return info
  405. def _is_authenticated(self, satellite_id: str) -> bool:
  406. """Prüft Admin-Authentifizierung."""
  407. session = self._admin_sessions.get(satellite_id)
  408. if not session:
  409. return False
  410. if time.time() > session.get("expires", 0):
  411. del self._admin_sessions[satellite_id]
  412. return False
  413. return session.get("authenticated", False)
  414. def _check_authorization(
  415. self, intent: str, wakeword_type: str, is_authenticated: bool
  416. ) -> bool:
  417. """Prüft Intent-Autorisierung."""
  418. if not is_admin_intent(intent):
  419. return True
  420. if wakeword_type != "system_command":
  421. return False
  422. if intent == "system_login":
  423. return True
  424. return is_authenticated
  425. async def _emit_authorization_error(
  426. self, satellite_id: str, session_info: dict[str, Any]
  427. ) -> None:
  428. """Emittiert Autorisierungs-Fehler."""
  429. error_event = OutputTextCreated(
  430. satellite_id=satellite_id,
  431. session_id=session_info.get("session_id", ""),
  432. room_id=session_info.get("room_id", ""),
  433. text="Dieser Befehl erfordert Administrator-Rechte. "
  434. "Bitte verwende das System-Wakeword und melde dich an.",
  435. intent="authorization_error",
  436. )
  437. await self._application.events.trigger("output_text_created", error_event)