file.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392
  1. # -*- coding: utf-8 -*-
  2. """
  3. Datei-basierter Event Store.
  4. Speichert Events in Dateien für Persistenz.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. import json
  9. import os
  10. from datetime import datetime
  11. from pathlib import Path
  12. from typing import TYPE_CHECKING, Any
  13. from trixy_core.events.persistence.store import (
  14. EventStore,
  15. StoredEvent,
  16. EventQuery,
  17. EventQueryResult,
  18. )
  19. if TYPE_CHECKING:
  20. from trixy_core.events.event_data.base import EventData
  21. class FileEventStore(EventStore):
  22. """
  23. Event Store mit Datei-Persistierung.
  24. Speichert Events als JSON-Lines in Dateien.
  25. Unterstützt automatische Rotation.
  26. Example:
  27. store = FileEventStore(
  28. base_path="./events",
  29. max_file_size_mb=10,
  30. rotate_daily=True,
  31. )
  32. await store.append("click", event_data)
  33. # Events laden
  34. result = await store.query(EventQuery(limit=100))
  35. """
  36. def __init__(
  37. self,
  38. base_path: str | Path,
  39. max_file_size_mb: float = 10.0,
  40. rotate_daily: bool = True,
  41. compress_old: bool = False,
  42. max_files: int = 100,
  43. ) -> None:
  44. """
  45. Initialisiert den Store.
  46. Args:
  47. base_path: Basisverzeichnis für Events.
  48. max_file_size_mb: Maximale Dateigröße in MB.
  49. rotate_daily: Täglich neue Datei.
  50. compress_old: Alte Dateien komprimieren.
  51. max_files: Maximale Anzahl Dateien.
  52. """
  53. self._base_path = Path(base_path)
  54. self._max_file_size = max_file_size_mb * 1024 * 1024 # Bytes
  55. self._rotate_daily = rotate_daily
  56. self._compress_old = compress_old
  57. self._max_files = max_files
  58. self._current_file: Path | None = None
  59. self._current_sequence = 0
  60. self._lock = asyncio.Lock()
  61. # Verzeichnis erstellen
  62. self._base_path.mkdir(parents=True, exist_ok=True)
  63. # Letzte Sequenz ermitteln
  64. self._load_last_sequence()
  65. def _load_last_sequence(self) -> None:
  66. """Lädt die letzte Sequenznummer."""
  67. max_seq = 0
  68. for file_path in self._base_path.glob("events_*.jsonl"):
  69. try:
  70. with open(file_path, "r", encoding="utf-8") as f:
  71. for line in f:
  72. try:
  73. data = json.loads(line.strip())
  74. seq = data.get("sequence", 0)
  75. max_seq = max(max_seq, seq)
  76. except json.JSONDecodeError:
  77. continue
  78. except IOError:
  79. continue
  80. self._current_sequence = max_seq
  81. def _get_current_file(self) -> Path:
  82. """Gibt die aktuelle Event-Datei zurück."""
  83. if self._rotate_daily:
  84. date_str = datetime.now().strftime("%Y%m%d")
  85. filename = f"events_{date_str}.jsonl"
  86. else:
  87. filename = "events.jsonl"
  88. file_path = self._base_path / filename
  89. # Größenprüfung
  90. if file_path.exists() and file_path.stat().st_size >= self._max_file_size:
  91. # Rotation
  92. counter = 1
  93. while True:
  94. rotated = self._base_path / f"{file_path.stem}_{counter}.jsonl"
  95. if not rotated.exists():
  96. file_path.rename(rotated)
  97. break
  98. counter += 1
  99. return file_path
  100. async def append(
  101. self,
  102. event_name: str,
  103. event_data: "EventData",
  104. ) -> StoredEvent:
  105. """
  106. Speichert ein Event.
  107. Args:
  108. event_name: Name des Events.
  109. event_data: Die Event-Daten.
  110. Returns:
  111. Das gespeicherte Event.
  112. """
  113. async with self._lock:
  114. self._current_sequence += 1
  115. stored = StoredEvent.from_event(
  116. event_name=event_name,
  117. event_data=event_data,
  118. sequence=self._current_sequence,
  119. )
  120. file_path = self._get_current_file()
  121. # Event als JSON-Line schreiben
  122. json_line = stored.to_json() + "\n"
  123. with open(file_path, "a", encoding="utf-8") as f:
  124. f.write(json_line)
  125. self._current_file = file_path
  126. # Cleanup alter Dateien
  127. await self._cleanup_old_files()
  128. return stored
  129. async def get(self, event_id: str) -> StoredEvent | None:
  130. """
  131. Gibt ein Event nach ID zurück.
  132. Args:
  133. event_id: Die Event-ID.
  134. Returns:
  135. StoredEvent oder None.
  136. """
  137. for file_path in sorted(self._base_path.glob("events_*.jsonl"), reverse=True):
  138. try:
  139. with open(file_path, "r", encoding="utf-8") as f:
  140. for line in f:
  141. try:
  142. data = json.loads(line.strip())
  143. if data.get("id") == event_id:
  144. return StoredEvent.from_dict(data)
  145. except json.JSONDecodeError:
  146. continue
  147. except IOError:
  148. continue
  149. return None
  150. async def query(self, query: EventQuery) -> EventQueryResult:
  151. """
  152. Führt eine Abfrage aus.
  153. Args:
  154. query: Die Abfrage.
  155. Returns:
  156. EventQueryResult.
  157. """
  158. all_events: list[StoredEvent] = []
  159. # Dateien durchsuchen
  160. files = sorted(self._base_path.glob("events_*.jsonl"))
  161. if not query.order_ascending:
  162. files = reversed(files)
  163. for file_path in files:
  164. try:
  165. with open(file_path, "r", encoding="utf-8") as f:
  166. for line in f:
  167. try:
  168. data = json.loads(line.strip())
  169. event = StoredEvent.from_dict(data)
  170. # Filter anwenden
  171. if not self._matches_query(event, query):
  172. continue
  173. all_events.append(event)
  174. # Früh abbrechen wenn genug
  175. if len(all_events) >= query.offset + query.limit + 1:
  176. break
  177. except json.JSONDecodeError:
  178. continue
  179. except IOError:
  180. continue
  181. # Sortieren
  182. all_events.sort(
  183. key=lambda e: e.sequence,
  184. reverse=not query.order_ascending,
  185. )
  186. total_count = len(all_events)
  187. # Paginierung
  188. start = query.offset
  189. end = start + query.limit
  190. paginated = all_events[start:end]
  191. return EventQueryResult(
  192. events=paginated,
  193. total_count=total_count,
  194. has_more=end < total_count,
  195. first_sequence=paginated[0].sequence if paginated else None,
  196. last_sequence=paginated[-1].sequence if paginated else None,
  197. )
  198. def _matches_query(self, event: StoredEvent, query: EventQuery) -> bool:
  199. """
  200. Prüft, ob ein Event zur Query passt.
  201. Args:
  202. event: Das Event.
  203. query: Die Query.
  204. Returns:
  205. True wenn passend.
  206. """
  207. if query.event_names and event.event_name not in query.event_names:
  208. return False
  209. if query.start_time and event.timestamp < query.start_time:
  210. return False
  211. if query.end_time and event.timestamp > query.end_time:
  212. return False
  213. if query.source and event.source != query.source:
  214. return False
  215. if query.after_sequence is not None and event.sequence <= query.after_sequence:
  216. return False
  217. if query.before_sequence is not None and event.sequence >= query.before_sequence:
  218. return False
  219. if query.metadata_filter:
  220. for key, value in query.metadata_filter.items():
  221. if event.metadata.get(key) != value:
  222. return False
  223. return True
  224. async def count(self, query: EventQuery | None = None) -> int:
  225. """
  226. Zählt Events.
  227. Args:
  228. query: Optionale Abfrage.
  229. Returns:
  230. Anzahl der Events.
  231. """
  232. if query is None:
  233. count = 0
  234. for file_path in self._base_path.glob("events_*.jsonl"):
  235. try:
  236. with open(file_path, "r", encoding="utf-8") as f:
  237. count += sum(1 for _ in f)
  238. except IOError:
  239. continue
  240. return count
  241. result = await self.query(EventQuery(
  242. event_names=query.event_names,
  243. start_time=query.start_time,
  244. end_time=query.end_time,
  245. source=query.source,
  246. metadata_filter=query.metadata_filter,
  247. limit=1000000,
  248. ))
  249. return result.total_count
  250. async def delete(self, event_id: str) -> bool:
  251. """
  252. Löscht ein Event.
  253. Args:
  254. event_id: Die Event-ID.
  255. Returns:
  256. True wenn gelöscht.
  257. """
  258. # Für Datei-Store nicht direkt unterstützt
  259. # Wäre aufwändig - besser: TTL-basiertes Cleanup
  260. return False
  261. async def clear(self) -> int:
  262. """
  263. Löscht alle Events.
  264. Returns:
  265. Anzahl gelöschter Events.
  266. """
  267. count = await self.count()
  268. for file_path in self._base_path.glob("events_*.jsonl"):
  269. try:
  270. file_path.unlink()
  271. except IOError:
  272. continue
  273. self._current_sequence = 0
  274. return count
  275. async def _cleanup_old_files(self) -> int:
  276. """
  277. Entfernt alte Dateien.
  278. Returns:
  279. Anzahl entfernter Dateien.
  280. """
  281. files = sorted(self._base_path.glob("events_*.jsonl"))
  282. if len(files) <= self._max_files:
  283. return 0
  284. to_remove = len(files) - self._max_files
  285. removed = 0
  286. for file_path in files[:to_remove]:
  287. try:
  288. file_path.unlink()
  289. removed += 1
  290. except IOError:
  291. continue
  292. return removed
  293. def get_file_stats(self) -> dict[str, Any]:
  294. """
  295. Gibt Datei-Statistiken zurück.
  296. Returns:
  297. Statistik-Dictionary.
  298. """
  299. files = list(self._base_path.glob("events_*.jsonl"))
  300. total_size = sum(f.stat().st_size for f in files if f.exists())
  301. return {
  302. "file_count": len(files),
  303. "total_size_mb": total_size / (1024 * 1024),
  304. "max_file_size_mb": self._max_file_size / (1024 * 1024),
  305. "max_files": self._max_files,
  306. "current_sequence": self._current_sequence,
  307. "files": [
  308. {
  309. "name": f.name,
  310. "size_mb": f.stat().st_size / (1024 * 1024),
  311. }
  312. for f in sorted(files)
  313. ],
  314. }