tracker.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493
  1. # -*- coding: utf-8 -*-
  2. """
  3. ActivityTracker — Zentrales Activity-Tracking als IService.
  4. Protokolliert Verbindungen, Intents, Fehler, Trainings und
  5. Conversations. Daten werden in JSON-Dateien unter data/activity/
  6. persistiert mit konfigurierbarer Retention.
  7. """
  8. from __future__ import annotations
  9. import asyncio
  10. import json
  11. import time
  12. from dataclasses import asdict
  13. from datetime import datetime, timedelta, timezone
  14. from pathlib import Path
  15. from typing import TYPE_CHECKING, Any
  16. from trixy_core.service.iservice import IService
  17. from trixy_core.service.enums import ServicePriority, ServiceGroup
  18. from trixy_core.events.decorators import TrixyEvent
  19. from trixy_core.activity.models import (
  20. SatelliteConnectionEntry,
  21. ConfigConnectionEntry,
  22. TrainingRunEntry,
  23. IntentLogEntry,
  24. ErrorLogEntry,
  25. ConversationEntry,
  26. )
  27. from trixy_core.utils.debug import pinfo, pdebug, perror
  28. if TYPE_CHECKING:
  29. from trixy_core.application import IApplication
  30. # Maximale Eintraege pro Kategorie
  31. MAX_ENTRIES = 10000
  32. # Flush-Intervall in Sekunden
  33. FLUSH_INTERVAL = 60
  34. # Cleanup-Intervall in Sekunden (1 Stunde)
  35. CLEANUP_INTERVAL = 3600
  36. def _now_iso() -> str:
  37. """Gibt den aktuellen Zeitstempel als ISO-8601 String zurueck."""
  38. return datetime.now(timezone.utc).isoformat()
  39. def _parse_iso(iso_str: str) -> datetime | None:
  40. """Parst einen ISO-8601 String, gibt None bei Fehler zurueck."""
  41. if not iso_str:
  42. return None
  43. try:
  44. return datetime.fromisoformat(iso_str)
  45. except (ValueError, TypeError):
  46. return None
  47. class ActivityTracker(IService):
  48. """
  49. Zentraler Activity-Tracker fuer Server-Aktivitaeten.
  50. Lauscht auf Events und bietet Logging-Methoden fuer externe
  51. Aufrufer. Persistiert Daten in JSON-Dateien und bereinigt
  52. alte Eintraege basierend auf retention_days.
  53. """
  54. NAME = "ActivityTracker"
  55. PRIORITY = ServicePriority.CORE
  56. GROUP = ServiceGroup.UTILITY
  57. DEPENDENCIES: list[str] = []
  58. # Dateinamen fuer JSON-Persistenz
  59. _FILES = {
  60. "satellite_connections": "satellite_connections.json",
  61. "config_connections": "config_connections.json",
  62. "training_runs": "training_runs.json",
  63. "intent_log": "intent_log.json",
  64. "error_log": "error_log.json",
  65. "conversations": "conversations.json",
  66. }
  67. def __init__(
  68. self,
  69. application: "IApplication",
  70. retention_days: int = 7,
  71. data_directory: str = "data/activity",
  72. ) -> None:
  73. super().__init__(application)
  74. self._retention_days = retention_days
  75. self._data_dir = Path(data_directory)
  76. # In-Memory-Listen
  77. self._satellite_connections: list[dict[str, Any]] = []
  78. self._config_connections: list[dict[str, Any]] = []
  79. self._training_runs: list[dict[str, Any]] = []
  80. self._intent_log: list[dict[str, Any]] = []
  81. self._error_log: list[dict[str, Any]] = []
  82. self._conversations: list[dict[str, Any]] = []
  83. # Hintergrund-Tasks
  84. self._flush_task: asyncio.Task | None = None
  85. self._cleanup_task: asyncio.Task | None = None
  86. self._dirty = False
  87. async def start(self) -> None:
  88. """Startet den Tracker: Daten laden, Events registrieren, Hintergrund-Tasks starten."""
  89. self._data_dir.mkdir(parents=True, exist_ok=True)
  90. self._load_all()
  91. self._cleanup_old_entries()
  92. # Event-Handler registrieren
  93. self._application.events.register("satellite_connected", self._on_satellite_connected)
  94. self._application.events.register("satellite_disconnected", self._on_satellite_disconnected)
  95. self._application.events.register("intent_handled", self._on_intent_handled)
  96. self._application.events.register("conversation_started", self._on_conversation_started)
  97. self._application.events.register("conversation_ended", self._on_conversation_ended)
  98. # Periodische Tasks starten
  99. self._flush_task = asyncio.create_task(self._periodic_flush())
  100. self._cleanup_task = asyncio.create_task(self._periodic_cleanup())
  101. pinfo(f"ActivityTracker gestartet (Retention: {self._retention_days} Tage, "
  102. f"Pfad: {self._data_dir})")
  103. async def stop(self) -> None:
  104. """Stoppt den Tracker und flusht ausstehende Daten."""
  105. if self._flush_task:
  106. self._flush_task.cancel()
  107. self._flush_task = None
  108. if self._cleanup_task:
  109. self._cleanup_task.cancel()
  110. self._cleanup_task = None
  111. self._flush_all()
  112. pinfo("ActivityTracker gestoppt")
  113. # =========================================================================
  114. # Event-Handler
  115. # =========================================================================
  116. async def _on_satellite_connected(self, event_name: str, event_data: Any) -> None:
  117. """Neuen Satellite-Verbindungseintrag erstellen."""
  118. entry = SatelliteConnectionEntry(
  119. satellite_id=getattr(event_data, "satellite_id", ""),
  120. alias=getattr(event_data, "alias", ""),
  121. room=getattr(event_data, "room", ""),
  122. connected_at=_now_iso(),
  123. ip_address=getattr(event_data, "ip_address", ""),
  124. )
  125. self._satellite_connections.append(asdict(entry))
  126. self._trim_list(self._satellite_connections)
  127. self._dirty = True
  128. pdebug(f"[ActivityTracker] Satellite verbunden: {entry.alias}")
  129. async def _on_satellite_disconnected(self, event_name: str, event_data: Any) -> None:
  130. """Disconnected-Zeitpunkt im letzten passenden Eintrag setzen."""
  131. sid = getattr(event_data, "satellite_id", "")
  132. now = _now_iso()
  133. # Letzten offenen Eintrag fuer diesen Satellite finden
  134. for entry in reversed(self._satellite_connections):
  135. if entry.get("satellite_id") == sid and not entry.get("disconnected_at"):
  136. entry["disconnected_at"] = now
  137. self._dirty = True
  138. pdebug(f"[ActivityTracker] Satellite getrennt: {sid}")
  139. return
  140. async def _on_intent_handled(self, event_name: str, event_data: Any) -> None:
  141. """Neuen Intent-Log-Eintrag erstellen."""
  142. # Plugin-Name aus IntentRegistry ermitteln
  143. plugin_name = ""
  144. try:
  145. from trixy_core.nlp.intent_registry import IntentRegistry
  146. registry = IntentRegistry.get_instance()
  147. intent_name = getattr(event_data, "intent", "")
  148. definition = registry.get(intent_name)
  149. if definition:
  150. plugin_name = definition.plugin_name
  151. except Exception:
  152. pass
  153. entry = IntentLogEntry(
  154. timestamp=_now_iso(),
  155. intent=getattr(event_data, "intent", ""),
  156. plugin=plugin_name,
  157. satellite_id=getattr(event_data, "satellite_id", ""),
  158. room_id=getattr(event_data, "room_id", ""),
  159. original_text=getattr(event_data, "original_text", ""),
  160. success=getattr(event_data, "success", True),
  161. )
  162. self._intent_log.append(asdict(entry))
  163. self._trim_list(self._intent_log)
  164. self._dirty = True
  165. async def _on_conversation_started(self, event_name: str, event_data: Any) -> None:
  166. """Neuen Conversation-Eintrag erstellen."""
  167. entry = ConversationEntry(
  168. session_id=getattr(event_data, "conversation_id", ""),
  169. satellite_id=getattr(event_data, "satellite_id", ""),
  170. started_at=_now_iso(),
  171. )
  172. self._conversations.append(asdict(entry))
  173. self._trim_list(self._conversations)
  174. self._dirty = True
  175. async def _on_conversation_ended(self, event_name: str, event_data: Any) -> None:
  176. """Completed-Zeitpunkt im passenden Conversation-Eintrag setzen."""
  177. session_id = getattr(event_data, "conversation_id", "")
  178. now = _now_iso()
  179. for entry in reversed(self._conversations):
  180. if entry.get("session_id") == session_id and not entry.get("completed_at"):
  181. entry["completed_at"] = now
  182. self._dirty = True
  183. return
  184. # =========================================================================
  185. # Oeffentliche Logging-Methoden (fuer externe Aufrufer)
  186. # =========================================================================
  187. def log_config_connection(self, peer: str) -> None:
  188. """Protokolliert eine neue Config-Tool-Verbindung."""
  189. entry = ConfigConnectionEntry(
  190. peer_address=peer,
  191. connected_at=_now_iso(),
  192. )
  193. self._config_connections.append(asdict(entry))
  194. self._trim_list(self._config_connections)
  195. self._dirty = True
  196. def log_config_disconnection(self, peer: str) -> None:
  197. """Setzt den Disconnected-Zeitpunkt fuer eine Config-Tool-Verbindung."""
  198. now = _now_iso()
  199. for entry in reversed(self._config_connections):
  200. if entry.get("peer_address") == peer and not entry.get("disconnected_at"):
  201. entry["disconnected_at"] = now
  202. self._dirty = True
  203. return
  204. def log_training_start(self, trainer_id: str, model_name: str) -> None:
  205. """Protokolliert den Start eines Trainings."""
  206. entry = TrainingRunEntry(
  207. trainer_id=trainer_id,
  208. model_name=model_name,
  209. started_at=_now_iso(),
  210. )
  211. self._training_runs.append(asdict(entry))
  212. self._trim_list(self._training_runs)
  213. self._dirty = True
  214. def log_training_complete(self, trainer_id: str, success: bool, error: str = "") -> None:
  215. """Protokolliert den Abschluss eines Trainings."""
  216. now = _now_iso()
  217. for entry in reversed(self._training_runs):
  218. if entry.get("trainer_id") == trainer_id and not entry.get("completed_at"):
  219. entry["completed_at"] = now
  220. entry["success"] = success
  221. entry["error"] = error
  222. # Dauer berechnen
  223. started = _parse_iso(entry.get("started_at", ""))
  224. if started:
  225. completed = _parse_iso(now)
  226. if completed:
  227. entry["duration_seconds"] = (completed - started).total_seconds()
  228. self._dirty = True
  229. return
  230. def log_error(self, source: str, message: str, satellite_id: str = "") -> None:
  231. """Protokolliert einen Fehler."""
  232. entry = ErrorLogEntry(
  233. timestamp=_now_iso(),
  234. source=source,
  235. message=message,
  236. satellite_id=satellite_id,
  237. )
  238. self._error_log.append(asdict(entry))
  239. self._trim_list(self._error_log)
  240. self._dirty = True
  241. # =========================================================================
  242. # Query-Methoden (fuer AdminCommands Plugin)
  243. # =========================================================================
  244. def get_last_intents(self, count: int = 5) -> list[dict[str, Any]]:
  245. """Liefert die letzten N Intents."""
  246. return self._intent_log[-count:]
  247. def get_last_intents_by_room(self, room: str, count: int = 5) -> list[dict[str, Any]]:
  248. """Liefert die letzten N Intents fuer einen bestimmten Raum."""
  249. filtered = [e for e in self._intent_log if e.get("room_id") == room]
  250. return filtered[-count:]
  251. def get_last_conversations(self, count: int = 5) -> list[dict[str, Any]]:
  252. """Liefert die letzten N Conversations."""
  253. return self._conversations[-count:]
  254. def get_last_errors(self, count: int = 5) -> list[dict[str, Any]]:
  255. """Liefert die letzten N Fehler."""
  256. return self._error_log[-count:]
  257. def get_error_count_since(self, hours: int = 24) -> int:
  258. """Zaehlt Fehler der letzten N Stunden."""
  259. cutoff = datetime.now(timezone.utc) - timedelta(hours=hours)
  260. count = 0
  261. for entry in self._error_log:
  262. ts = _parse_iso(entry.get("timestamp", ""))
  263. if ts and ts >= cutoff:
  264. count += 1
  265. return count
  266. def get_last_training(self, trainer_id: str = "") -> dict[str, Any] | None:
  267. """Liefert den letzten Trainings-Lauf (optional gefiltert nach trainer_id)."""
  268. for entry in reversed(self._training_runs):
  269. if not trainer_id or entry.get("trainer_id") == trainer_id:
  270. return entry
  271. return None
  272. def get_last_config_connection(self) -> dict[str, Any] | None:
  273. """Liefert die letzte Config-Tool-Verbindung."""
  274. if self._config_connections:
  275. return self._config_connections[-1]
  276. return None
  277. def get_satellite_connection_history(self, alias: str) -> list[dict[str, Any]]:
  278. """Liefert die Verbindungshistorie fuer einen Satellite (nach Alias)."""
  279. return [
  280. e for e in self._satellite_connections
  281. if e.get("alias", "").lower() == alias.lower()
  282. ]
  283. def get_intent_stats(self, days: int = 7) -> dict[str, int]:
  284. """Liefert Intent-Statistiken der letzten N Tage."""
  285. cutoff = datetime.now(timezone.utc) - timedelta(days=days)
  286. stats: dict[str, int] = {}
  287. for entry in self._intent_log:
  288. ts = _parse_iso(entry.get("timestamp", ""))
  289. if ts and ts >= cutoff:
  290. intent_name = entry.get("intent", "unknown")
  291. stats[intent_name] = stats.get(intent_name, 0) + 1
  292. return stats
  293. def clear_all(self) -> None:
  294. """Loescht alle Activity-Daten."""
  295. self._satellite_connections.clear()
  296. self._config_connections.clear()
  297. self._training_runs.clear()
  298. self._intent_log.clear()
  299. self._error_log.clear()
  300. self._conversations.clear()
  301. self._dirty = True
  302. self._flush_all()
  303. pinfo("[ActivityTracker] Alle Daten geloescht")
  304. # =========================================================================
  305. # Persistenz
  306. # =========================================================================
  307. def _load_all(self) -> None:
  308. """Laedt alle JSON-Dateien in den Speicher."""
  309. mapping = {
  310. "satellite_connections": "_satellite_connections",
  311. "config_connections": "_config_connections",
  312. "training_runs": "_training_runs",
  313. "intent_log": "_intent_log",
  314. "error_log": "_error_log",
  315. "conversations": "_conversations",
  316. }
  317. for key, attr in mapping.items():
  318. file_path = self._data_dir / self._FILES[key]
  319. data = self._load_json(file_path)
  320. setattr(self, attr, data)
  321. total = sum(
  322. len(getattr(self, attr)) for attr in mapping.values()
  323. )
  324. pdebug(f"[ActivityTracker] {total} Eintraege geladen")
  325. def _flush_all(self) -> None:
  326. """Schreibt alle In-Memory-Daten in JSON-Dateien."""
  327. if not self._dirty:
  328. return
  329. mapping = {
  330. "satellite_connections": self._satellite_connections,
  331. "config_connections": self._config_connections,
  332. "training_runs": self._training_runs,
  333. "intent_log": self._intent_log,
  334. "error_log": self._error_log,
  335. "conversations": self._conversations,
  336. }
  337. for key, data in mapping.items():
  338. file_path = self._data_dir / self._FILES[key]
  339. self._save_json(file_path, data)
  340. self._dirty = False
  341. def _load_json(self, path: Path) -> list[dict[str, Any]]:
  342. """Laedt eine JSON-Datei, gibt leere Liste bei Fehler zurueck."""
  343. if not path.exists():
  344. return []
  345. try:
  346. with open(path, "r", encoding="utf-8") as f:
  347. data = json.load(f)
  348. if isinstance(data, list):
  349. return data
  350. except (json.JSONDecodeError, OSError) as e:
  351. perror(f"[ActivityTracker] Fehler beim Laden von {path}: {e}")
  352. return []
  353. def _save_json(self, path: Path, data: list[dict[str, Any]]) -> None:
  354. """Speichert Daten als JSON-Datei."""
  355. try:
  356. with open(path, "w", encoding="utf-8") as f:
  357. json.dump(data, f, ensure_ascii=False, indent=2)
  358. except OSError as e:
  359. perror(f"[ActivityTracker] Fehler beim Speichern von {path}: {e}")
  360. def _trim_list(self, lst: list) -> None:
  361. """Begrenzt eine Liste auf MAX_ENTRIES Eintraege."""
  362. if len(lst) > MAX_ENTRIES:
  363. del lst[: len(lst) - MAX_ENTRIES]
  364. # =========================================================================
  365. # Retention / Cleanup
  366. # =========================================================================
  367. def _cleanup_old_entries(self) -> None:
  368. """Entfernt Eintraege die aelter als retention_days sind."""
  369. cutoff = datetime.now(timezone.utc) - timedelta(days=self._retention_days)
  370. before = sum(len(lst) for lst in [
  371. self._satellite_connections, self._config_connections,
  372. self._training_runs, self._intent_log,
  373. self._error_log, self._conversations,
  374. ])
  375. # Timestamp-Feld pro Kategorie
  376. ts_fields = {
  377. "_satellite_connections": "connected_at",
  378. "_config_connections": "connected_at",
  379. "_training_runs": "started_at",
  380. "_intent_log": "timestamp",
  381. "_error_log": "timestamp",
  382. "_conversations": "started_at",
  383. }
  384. for attr, ts_field in ts_fields.items():
  385. lst: list[dict] = getattr(self, attr)
  386. filtered = []
  387. for entry in lst:
  388. ts = _parse_iso(entry.get(ts_field, ""))
  389. if ts is None or ts >= cutoff:
  390. filtered.append(entry)
  391. setattr(self, attr, filtered)
  392. after = sum(len(lst) for lst in [
  393. self._satellite_connections, self._config_connections,
  394. self._training_runs, self._intent_log,
  395. self._error_log, self._conversations,
  396. ])
  397. removed = before - after
  398. if removed > 0:
  399. self._dirty = True
  400. pdebug(f"[ActivityTracker] {removed} alte Eintraege bereinigt")
  401. # =========================================================================
  402. # Hintergrund-Tasks
  403. # =========================================================================
  404. async def _periodic_flush(self) -> None:
  405. """Periodischer Flush der In-Memory-Daten."""
  406. while True:
  407. try:
  408. await asyncio.sleep(FLUSH_INTERVAL)
  409. self._flush_all()
  410. except asyncio.CancelledError:
  411. break
  412. except Exception as e:
  413. perror(f"[ActivityTracker] Flush-Fehler: {e}")
  414. async def _periodic_cleanup(self) -> None:
  415. """Periodische Bereinigung alter Eintraege."""
  416. while True:
  417. try:
  418. await asyncio.sleep(CLEANUP_INTERVAL)
  419. self._cleanup_old_entries()
  420. self._flush_all()
  421. except asyncio.CancelledError:
  422. break
  423. except Exception as e:
  424. perror(f"[ActivityTracker] Cleanup-Fehler: {e}")