| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392 |
- # -*- 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)
- ],
- }
|