| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493 |
- # -*- coding: utf-8 -*-
- """
- ActivityTracker — Zentrales Activity-Tracking als IService.
- Protokolliert Verbindungen, Intents, Fehler, Trainings und
- Conversations. Daten werden in JSON-Dateien unter data/activity/
- persistiert mit konfigurierbarer Retention.
- """
- from __future__ import annotations
- import asyncio
- import json
- import time
- from dataclasses import asdict
- from datetime import datetime, timedelta, timezone
- from pathlib import Path
- from typing import TYPE_CHECKING, Any
- 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.activity.models import (
- SatelliteConnectionEntry,
- ConfigConnectionEntry,
- TrainingRunEntry,
- IntentLogEntry,
- ErrorLogEntry,
- ConversationEntry,
- )
- from trixy_core.utils.debug import pinfo, pdebug, perror
- if TYPE_CHECKING:
- from trixy_core.application import IApplication
- # Maximale Eintraege pro Kategorie
- MAX_ENTRIES = 10000
- # Flush-Intervall in Sekunden
- FLUSH_INTERVAL = 60
- # Cleanup-Intervall in Sekunden (1 Stunde)
- CLEANUP_INTERVAL = 3600
- def _now_iso() -> str:
- """Gibt den aktuellen Zeitstempel als ISO-8601 String zurueck."""
- return datetime.now(timezone.utc).isoformat()
- def _parse_iso(iso_str: str) -> datetime | None:
- """Parst einen ISO-8601 String, gibt None bei Fehler zurueck."""
- if not iso_str:
- return None
- try:
- return datetime.fromisoformat(iso_str)
- except (ValueError, TypeError):
- return None
- class ActivityTracker(IService):
- """
- Zentraler Activity-Tracker fuer Server-Aktivitaeten.
- Lauscht auf Events und bietet Logging-Methoden fuer externe
- Aufrufer. Persistiert Daten in JSON-Dateien und bereinigt
- alte Eintraege basierend auf retention_days.
- """
- NAME = "ActivityTracker"
- PRIORITY = ServicePriority.CORE
- GROUP = ServiceGroup.UTILITY
- DEPENDENCIES: list[str] = []
- # Dateinamen fuer JSON-Persistenz
- _FILES = {
- "satellite_connections": "satellite_connections.json",
- "config_connections": "config_connections.json",
- "training_runs": "training_runs.json",
- "intent_log": "intent_log.json",
- "error_log": "error_log.json",
- "conversations": "conversations.json",
- }
- def __init__(
- self,
- application: "IApplication",
- retention_days: int = 7,
- data_directory: str = "data/activity",
- ) -> None:
- super().__init__(application)
- self._retention_days = retention_days
- self._data_dir = Path(data_directory)
- # In-Memory-Listen
- self._satellite_connections: list[dict[str, Any]] = []
- self._config_connections: list[dict[str, Any]] = []
- self._training_runs: list[dict[str, Any]] = []
- self._intent_log: list[dict[str, Any]] = []
- self._error_log: list[dict[str, Any]] = []
- self._conversations: list[dict[str, Any]] = []
- # Hintergrund-Tasks
- self._flush_task: asyncio.Task | None = None
- self._cleanup_task: asyncio.Task | None = None
- self._dirty = False
- async def start(self) -> None:
- """Startet den Tracker: Daten laden, Events registrieren, Hintergrund-Tasks starten."""
- self._data_dir.mkdir(parents=True, exist_ok=True)
- self._load_all()
- self._cleanup_old_entries()
- # Event-Handler registrieren
- self._application.events.register("satellite_connected", self._on_satellite_connected)
- self._application.events.register("satellite_disconnected", self._on_satellite_disconnected)
- self._application.events.register("intent_handled", self._on_intent_handled)
- self._application.events.register("conversation_started", self._on_conversation_started)
- self._application.events.register("conversation_ended", self._on_conversation_ended)
- # Periodische Tasks starten
- self._flush_task = asyncio.create_task(self._periodic_flush())
- self._cleanup_task = asyncio.create_task(self._periodic_cleanup())
- pinfo(f"ActivityTracker gestartet (Retention: {self._retention_days} Tage, "
- f"Pfad: {self._data_dir})")
- async def stop(self) -> None:
- """Stoppt den Tracker und flusht ausstehende Daten."""
- if self._flush_task:
- self._flush_task.cancel()
- self._flush_task = None
- if self._cleanup_task:
- self._cleanup_task.cancel()
- self._cleanup_task = None
- self._flush_all()
- pinfo("ActivityTracker gestoppt")
- # =========================================================================
- # Event-Handler
- # =========================================================================
- async def _on_satellite_connected(self, event_name: str, event_data: Any) -> None:
- """Neuen Satellite-Verbindungseintrag erstellen."""
- entry = SatelliteConnectionEntry(
- satellite_id=getattr(event_data, "satellite_id", ""),
- alias=getattr(event_data, "alias", ""),
- room=getattr(event_data, "room", ""),
- connected_at=_now_iso(),
- ip_address=getattr(event_data, "ip_address", ""),
- )
- self._satellite_connections.append(asdict(entry))
- self._trim_list(self._satellite_connections)
- self._dirty = True
- pdebug(f"[ActivityTracker] Satellite verbunden: {entry.alias}")
- async def _on_satellite_disconnected(self, event_name: str, event_data: Any) -> None:
- """Disconnected-Zeitpunkt im letzten passenden Eintrag setzen."""
- sid = getattr(event_data, "satellite_id", "")
- now = _now_iso()
- # Letzten offenen Eintrag fuer diesen Satellite finden
- for entry in reversed(self._satellite_connections):
- if entry.get("satellite_id") == sid and not entry.get("disconnected_at"):
- entry["disconnected_at"] = now
- self._dirty = True
- pdebug(f"[ActivityTracker] Satellite getrennt: {sid}")
- return
- async def _on_intent_handled(self, event_name: str, event_data: Any) -> None:
- """Neuen Intent-Log-Eintrag erstellen."""
- # Plugin-Name aus IntentRegistry ermitteln
- plugin_name = ""
- try:
- from trixy_core.nlp.intent_registry import IntentRegistry
- registry = IntentRegistry.get_instance()
- intent_name = getattr(event_data, "intent", "")
- definition = registry.get(intent_name)
- if definition:
- plugin_name = definition.plugin_name
- except Exception:
- pass
- entry = IntentLogEntry(
- timestamp=_now_iso(),
- intent=getattr(event_data, "intent", ""),
- plugin=plugin_name,
- satellite_id=getattr(event_data, "satellite_id", ""),
- room_id=getattr(event_data, "room_id", ""),
- original_text=getattr(event_data, "original_text", ""),
- success=getattr(event_data, "success", True),
- )
- self._intent_log.append(asdict(entry))
- self._trim_list(self._intent_log)
- self._dirty = True
- async def _on_conversation_started(self, event_name: str, event_data: Any) -> None:
- """Neuen Conversation-Eintrag erstellen."""
- entry = ConversationEntry(
- session_id=getattr(event_data, "conversation_id", ""),
- satellite_id=getattr(event_data, "satellite_id", ""),
- started_at=_now_iso(),
- )
- self._conversations.append(asdict(entry))
- self._trim_list(self._conversations)
- self._dirty = True
- async def _on_conversation_ended(self, event_name: str, event_data: Any) -> None:
- """Completed-Zeitpunkt im passenden Conversation-Eintrag setzen."""
- session_id = getattr(event_data, "conversation_id", "")
- now = _now_iso()
- for entry in reversed(self._conversations):
- if entry.get("session_id") == session_id and not entry.get("completed_at"):
- entry["completed_at"] = now
- self._dirty = True
- return
- # =========================================================================
- # Oeffentliche Logging-Methoden (fuer externe Aufrufer)
- # =========================================================================
- def log_config_connection(self, peer: str) -> None:
- """Protokolliert eine neue Config-Tool-Verbindung."""
- entry = ConfigConnectionEntry(
- peer_address=peer,
- connected_at=_now_iso(),
- )
- self._config_connections.append(asdict(entry))
- self._trim_list(self._config_connections)
- self._dirty = True
- def log_config_disconnection(self, peer: str) -> None:
- """Setzt den Disconnected-Zeitpunkt fuer eine Config-Tool-Verbindung."""
- now = _now_iso()
- for entry in reversed(self._config_connections):
- if entry.get("peer_address") == peer and not entry.get("disconnected_at"):
- entry["disconnected_at"] = now
- self._dirty = True
- return
- def log_training_start(self, trainer_id: str, model_name: str) -> None:
- """Protokolliert den Start eines Trainings."""
- entry = TrainingRunEntry(
- trainer_id=trainer_id,
- model_name=model_name,
- started_at=_now_iso(),
- )
- self._training_runs.append(asdict(entry))
- self._trim_list(self._training_runs)
- self._dirty = True
- def log_training_complete(self, trainer_id: str, success: bool, error: str = "") -> None:
- """Protokolliert den Abschluss eines Trainings."""
- now = _now_iso()
- for entry in reversed(self._training_runs):
- if entry.get("trainer_id") == trainer_id and not entry.get("completed_at"):
- entry["completed_at"] = now
- entry["success"] = success
- entry["error"] = error
- # Dauer berechnen
- started = _parse_iso(entry.get("started_at", ""))
- if started:
- completed = _parse_iso(now)
- if completed:
- entry["duration_seconds"] = (completed - started).total_seconds()
- self._dirty = True
- return
- def log_error(self, source: str, message: str, satellite_id: str = "") -> None:
- """Protokolliert einen Fehler."""
- entry = ErrorLogEntry(
- timestamp=_now_iso(),
- source=source,
- message=message,
- satellite_id=satellite_id,
- )
- self._error_log.append(asdict(entry))
- self._trim_list(self._error_log)
- self._dirty = True
- # =========================================================================
- # Query-Methoden (fuer AdminCommands Plugin)
- # =========================================================================
- def get_last_intents(self, count: int = 5) -> list[dict[str, Any]]:
- """Liefert die letzten N Intents."""
- return self._intent_log[-count:]
- def get_last_intents_by_room(self, room: str, count: int = 5) -> list[dict[str, Any]]:
- """Liefert die letzten N Intents fuer einen bestimmten Raum."""
- filtered = [e for e in self._intent_log if e.get("room_id") == room]
- return filtered[-count:]
- def get_last_conversations(self, count: int = 5) -> list[dict[str, Any]]:
- """Liefert die letzten N Conversations."""
- return self._conversations[-count:]
- def get_last_errors(self, count: int = 5) -> list[dict[str, Any]]:
- """Liefert die letzten N Fehler."""
- return self._error_log[-count:]
- def get_error_count_since(self, hours: int = 24) -> int:
- """Zaehlt Fehler der letzten N Stunden."""
- cutoff = datetime.now(timezone.utc) - timedelta(hours=hours)
- count = 0
- for entry in self._error_log:
- ts = _parse_iso(entry.get("timestamp", ""))
- if ts and ts >= cutoff:
- count += 1
- return count
- def get_last_training(self, trainer_id: str = "") -> dict[str, Any] | None:
- """Liefert den letzten Trainings-Lauf (optional gefiltert nach trainer_id)."""
- for entry in reversed(self._training_runs):
- if not trainer_id or entry.get("trainer_id") == trainer_id:
- return entry
- return None
- def get_last_config_connection(self) -> dict[str, Any] | None:
- """Liefert die letzte Config-Tool-Verbindung."""
- if self._config_connections:
- return self._config_connections[-1]
- return None
- def get_satellite_connection_history(self, alias: str) -> list[dict[str, Any]]:
- """Liefert die Verbindungshistorie fuer einen Satellite (nach Alias)."""
- return [
- e for e in self._satellite_connections
- if e.get("alias", "").lower() == alias.lower()
- ]
- def get_intent_stats(self, days: int = 7) -> dict[str, int]:
- """Liefert Intent-Statistiken der letzten N Tage."""
- cutoff = datetime.now(timezone.utc) - timedelta(days=days)
- stats: dict[str, int] = {}
- for entry in self._intent_log:
- ts = _parse_iso(entry.get("timestamp", ""))
- if ts and ts >= cutoff:
- intent_name = entry.get("intent", "unknown")
- stats[intent_name] = stats.get(intent_name, 0) + 1
- return stats
- def clear_all(self) -> None:
- """Loescht alle Activity-Daten."""
- self._satellite_connections.clear()
- self._config_connections.clear()
- self._training_runs.clear()
- self._intent_log.clear()
- self._error_log.clear()
- self._conversations.clear()
- self._dirty = True
- self._flush_all()
- pinfo("[ActivityTracker] Alle Daten geloescht")
- # =========================================================================
- # Persistenz
- # =========================================================================
- def _load_all(self) -> None:
- """Laedt alle JSON-Dateien in den Speicher."""
- mapping = {
- "satellite_connections": "_satellite_connections",
- "config_connections": "_config_connections",
- "training_runs": "_training_runs",
- "intent_log": "_intent_log",
- "error_log": "_error_log",
- "conversations": "_conversations",
- }
- for key, attr in mapping.items():
- file_path = self._data_dir / self._FILES[key]
- data = self._load_json(file_path)
- setattr(self, attr, data)
- total = sum(
- len(getattr(self, attr)) for attr in mapping.values()
- )
- pdebug(f"[ActivityTracker] {total} Eintraege geladen")
- def _flush_all(self) -> None:
- """Schreibt alle In-Memory-Daten in JSON-Dateien."""
- if not self._dirty:
- return
- mapping = {
- "satellite_connections": self._satellite_connections,
- "config_connections": self._config_connections,
- "training_runs": self._training_runs,
- "intent_log": self._intent_log,
- "error_log": self._error_log,
- "conversations": self._conversations,
- }
- for key, data in mapping.items():
- file_path = self._data_dir / self._FILES[key]
- self._save_json(file_path, data)
- self._dirty = False
- def _load_json(self, path: Path) -> list[dict[str, Any]]:
- """Laedt eine JSON-Datei, gibt leere Liste bei Fehler zurueck."""
- if not path.exists():
- return []
- try:
- with open(path, "r", encoding="utf-8") as f:
- data = json.load(f)
- if isinstance(data, list):
- return data
- except (json.JSONDecodeError, OSError) as e:
- perror(f"[ActivityTracker] Fehler beim Laden von {path}: {e}")
- return []
- def _save_json(self, path: Path, data: list[dict[str, Any]]) -> None:
- """Speichert Daten als JSON-Datei."""
- try:
- with open(path, "w", encoding="utf-8") as f:
- json.dump(data, f, ensure_ascii=False, indent=2)
- except OSError as e:
- perror(f"[ActivityTracker] Fehler beim Speichern von {path}: {e}")
- def _trim_list(self, lst: list) -> None:
- """Begrenzt eine Liste auf MAX_ENTRIES Eintraege."""
- if len(lst) > MAX_ENTRIES:
- del lst[: len(lst) - MAX_ENTRIES]
- # =========================================================================
- # Retention / Cleanup
- # =========================================================================
- def _cleanup_old_entries(self) -> None:
- """Entfernt Eintraege die aelter als retention_days sind."""
- cutoff = datetime.now(timezone.utc) - timedelta(days=self._retention_days)
- before = sum(len(lst) for lst in [
- self._satellite_connections, self._config_connections,
- self._training_runs, self._intent_log,
- self._error_log, self._conversations,
- ])
- # Timestamp-Feld pro Kategorie
- ts_fields = {
- "_satellite_connections": "connected_at",
- "_config_connections": "connected_at",
- "_training_runs": "started_at",
- "_intent_log": "timestamp",
- "_error_log": "timestamp",
- "_conversations": "started_at",
- }
- for attr, ts_field in ts_fields.items():
- lst: list[dict] = getattr(self, attr)
- filtered = []
- for entry in lst:
- ts = _parse_iso(entry.get(ts_field, ""))
- if ts is None or ts >= cutoff:
- filtered.append(entry)
- setattr(self, attr, filtered)
- after = sum(len(lst) for lst in [
- self._satellite_connections, self._config_connections,
- self._training_runs, self._intent_log,
- self._error_log, self._conversations,
- ])
- removed = before - after
- if removed > 0:
- self._dirty = True
- pdebug(f"[ActivityTracker] {removed} alte Eintraege bereinigt")
- # =========================================================================
- # Hintergrund-Tasks
- # =========================================================================
- async def _periodic_flush(self) -> None:
- """Periodischer Flush der In-Memory-Daten."""
- while True:
- try:
- await asyncio.sleep(FLUSH_INTERVAL)
- self._flush_all()
- except asyncio.CancelledError:
- break
- except Exception as e:
- perror(f"[ActivityTracker] Flush-Fehler: {e}")
- async def _periodic_cleanup(self) -> None:
- """Periodische Bereinigung alter Eintraege."""
- while True:
- try:
- await asyncio.sleep(CLEANUP_INTERVAL)
- self._cleanup_old_entries()
- self._flush_all()
- except asyncio.CancelledError:
- break
- except Exception as e:
- perror(f"[ActivityTracker] Cleanup-Fehler: {e}")
|