# -*- coding: utf-8 -*- """ Datei-basierter Event Store. Speichert Events in Dateien für Persistenz. """ from __future__ import annotations import asyncio import json import os from datetime import datetime from pathlib import Path from typing import TYPE_CHECKING, Any from trixy_core.events.persistence.store import ( EventStore, StoredEvent, EventQuery, EventQueryResult, ) if TYPE_CHECKING: from trixy_core.events.event_data.base import EventData class FileEventStore(EventStore): """ Event Store mit Datei-Persistierung. Speichert Events als JSON-Lines in Dateien. Unterstützt automatische Rotation. Example: store = FileEventStore( base_path="./events", max_file_size_mb=10, rotate_daily=True, ) await store.append("click", event_data) # Events laden result = await store.query(EventQuery(limit=100)) """ def __init__( self, base_path: str | Path, max_file_size_mb: float = 10.0, rotate_daily: bool = True, compress_old: bool = False, max_files: int = 100, ) -> None: """ Initialisiert den Store. Args: base_path: Basisverzeichnis für Events. max_file_size_mb: Maximale Dateigröße in MB. rotate_daily: Täglich neue Datei. compress_old: Alte Dateien komprimieren. max_files: Maximale Anzahl Dateien. """ self._base_path = Path(base_path) self._max_file_size = max_file_size_mb * 1024 * 1024 # Bytes self._rotate_daily = rotate_daily self._compress_old = compress_old self._max_files = max_files self._current_file: Path | None = None self._current_sequence = 0 self._lock = asyncio.Lock() # Verzeichnis erstellen self._base_path.mkdir(parents=True, exist_ok=True) # Letzte Sequenz ermitteln self._load_last_sequence() def _load_last_sequence(self) -> None: """Lädt die letzte Sequenznummer.""" max_seq = 0 for file_path in self._base_path.glob("events_*.jsonl"): try: with open(file_path, "r", encoding="utf-8") as f: for line in f: try: data = json.loads(line.strip()) seq = data.get("sequence", 0) max_seq = max(max_seq, seq) except json.JSONDecodeError: continue except IOError: continue self._current_sequence = max_seq def _get_current_file(self) -> Path: """Gibt die aktuelle Event-Datei zurück.""" if self._rotate_daily: date_str = datetime.now().strftime("%Y%m%d") filename = f"events_{date_str}.jsonl" else: filename = "events.jsonl" file_path = self._base_path / filename # Größenprüfung if file_path.exists() and file_path.stat().st_size >= self._max_file_size: # Rotation counter = 1 while True: rotated = self._base_path / f"{file_path.stem}_{counter}.jsonl" if not rotated.exists(): file_path.rename(rotated) break counter += 1 return file_path async def append( self, event_name: str, event_data: "EventData", ) -> StoredEvent: """ Speichert ein Event. Args: event_name: Name des Events. event_data: Die Event-Daten. Returns: Das gespeicherte Event. """ async with self._lock: self._current_sequence += 1 stored = StoredEvent.from_event( event_name=event_name, event_data=event_data, sequence=self._current_sequence, ) file_path = self._get_current_file() # Event als JSON-Line schreiben json_line = stored.to_json() + "\n" with open(file_path, "a", encoding="utf-8") as f: f.write(json_line) self._current_file = file_path # Cleanup alter Dateien await self._cleanup_old_files() return stored async def get(self, event_id: str) -> StoredEvent | None: """ Gibt ein Event nach ID zurück. Args: event_id: Die Event-ID. Returns: StoredEvent oder None. """ for file_path in sorted(self._base_path.glob("events_*.jsonl"), reverse=True): try: with open(file_path, "r", encoding="utf-8") as f: for line in f: try: data = json.loads(line.strip()) if data.get("id") == event_id: return StoredEvent.from_dict(data) except json.JSONDecodeError: continue except IOError: continue return None async def query(self, query: EventQuery) -> EventQueryResult: """ Führt eine Abfrage aus. Args: query: Die Abfrage. Returns: EventQueryResult. """ all_events: list[StoredEvent] = [] # Dateien durchsuchen files = sorted(self._base_path.glob("events_*.jsonl")) if not query.order_ascending: files = reversed(files) for file_path in files: try: with open(file_path, "r", encoding="utf-8") as f: for line in f: try: data = json.loads(line.strip()) event = StoredEvent.from_dict(data) # Filter anwenden if not self._matches_query(event, query): continue all_events.append(event) # Früh abbrechen wenn genug if len(all_events) >= query.offset + query.limit + 1: break except json.JSONDecodeError: continue except IOError: continue # Sortieren all_events.sort( key=lambda e: e.sequence, reverse=not query.order_ascending, ) total_count = len(all_events) # Paginierung start = query.offset end = start + query.limit paginated = all_events[start:end] return EventQueryResult( events=paginated, total_count=total_count, has_more=end < total_count, first_sequence=paginated[0].sequence if paginated else None, last_sequence=paginated[-1].sequence if paginated else None, ) def _matches_query(self, event: StoredEvent, query: EventQuery) -> bool: """ Prüft, ob ein Event zur Query passt. Args: event: Das Event. query: Die Query. Returns: True wenn passend. """ if query.event_names and event.event_name not in query.event_names: return False if query.start_time and event.timestamp < query.start_time: return False if query.end_time and event.timestamp > query.end_time: return False if query.source and event.source != query.source: return False if query.after_sequence is not None and event.sequence <= query.after_sequence: return False if query.before_sequence is not None and event.sequence >= query.before_sequence: return False if query.metadata_filter: for key, value in query.metadata_filter.items(): if event.metadata.get(key) != value: return False return True async def count(self, query: EventQuery | None = None) -> int: """ Zählt Events. Args: query: Optionale Abfrage. Returns: Anzahl der Events. """ if query is None: count = 0 for file_path in self._base_path.glob("events_*.jsonl"): try: with open(file_path, "r", encoding="utf-8") as f: count += sum(1 for _ in f) except IOError: continue return count result = await self.query(EventQuery( event_names=query.event_names, start_time=query.start_time, end_time=query.end_time, source=query.source, metadata_filter=query.metadata_filter, limit=1000000, )) return result.total_count async def delete(self, event_id: str) -> bool: """ Löscht ein Event. Args: event_id: Die Event-ID. Returns: True wenn gelöscht. """ # Für Datei-Store nicht direkt unterstützt # Wäre aufwändig - besser: TTL-basiertes Cleanup return False async def clear(self) -> int: """ Löscht alle Events. Returns: Anzahl gelöschter Events. """ count = await self.count() for file_path in self._base_path.glob("events_*.jsonl"): try: file_path.unlink() except IOError: continue self._current_sequence = 0 return count async def _cleanup_old_files(self) -> int: """ Entfernt alte Dateien. Returns: Anzahl entfernter Dateien. """ files = sorted(self._base_path.glob("events_*.jsonl")) if len(files) <= self._max_files: return 0 to_remove = len(files) - self._max_files removed = 0 for file_path in files[:to_remove]: try: file_path.unlink() removed += 1 except IOError: continue return removed def get_file_stats(self) -> dict[str, Any]: """ Gibt Datei-Statistiken zurück. Returns: Statistik-Dictionary. """ files = list(self._base_path.glob("events_*.jsonl")) total_size = sum(f.stat().st_size for f in files if f.exists()) return { "file_count": len(files), "total_size_mb": total_size / (1024 * 1024), "max_file_size_mb": self._max_file_size / (1024 * 1024), "max_files": self._max_files, "current_sequence": self._current_sequence, "files": [ { "name": f.name, "size_mb": f.stat().st_size / (1024 * 1024), } for f in sorted(files) ], }