import 'dart:async'; import 'dart:convert'; import 'dart:io'; import 'package:flutter/foundation.dart'; import 'package:shared_preferences/shared_preferences.dart'; import '../config/app_config.dart'; import '../database/db_helper.dart'; import '../services/auth_service.dart'; import '../services/favoriten_service.dart'; import '../services/sync_service.dart'; import '../services/melo_logger.dart'; /// Sync-Modi für die Cloud-Synchronisation. enum SyncModus { echtzeit, // SSE-Push, sofort bei Änderungen manuell, // Kein Auto-Sync, nur auf Knopfdruck intervall, // Periodisches Polling (2T/3T/5T/1W/2W) } /// Echtzeit-Cloud-Sync via SSE (Server-Sent-Events). /// /// Hält eine persistente HTTP-Verbindung zum Server offen und empfängt /// Änderungen (Favoriten-Toggle, Song-Löschung, Titel-Update) in Echtzeit. /// /// **3 Sync-Modi** (via SharedPreferences `sync_modus`): /// 1. ⚡ **Echtzeit** (default): SSE-Push, sofortige Verarbeitung. /// Bei Fehlern: Reconnect-Backoff (max 3 Versuche), /// dann Fallback auf Intervall-Modus. /// 2. 👆 **Manuell**: Kein Auto-Sync. Nur Button „Jetzt synchronisieren". /// 3. 🗓️ **Intervall**: Periodisches Polling mit `cloud_interval_stunden` /// (48/72/120/168/336 Stunden = 2T/3T/5T/1W/2W). /// /// **Lifecycle** (via WidgetsBindingObserver in main.dart): /// - resumed → fortsetzen() (SSE neu verbinden) /// - paused/hidden → pausiere() (Batterie sparen) class RealtimeSyncService { RealtimeSyncService() : _db = DbHelper(), _auth = AuthService(), _favoriten = FavoritenService(); final DbHelper _db; final AuthService _auth; final FavoritenService _favoriten; // ─── Verbindungs-Zustand ─── StreamSubscription? _sseSubscription; HttpClientResponse? _sseResponse; HttpClient? _httpClient; bool _laeuft = false; bool _pausiert = false; /// Aktueller SSE-Event-Typ (wird über `event:`-Zeile gesetzt). String? _aktuellerEventTyp; // Reconnection-Backoff int _reconnectVersuche = 0; static const int _maxReconnectVersuche = 3; static const Duration _backoffStart = Duration(seconds: 2); // ─── Notifier für UI ─── final ValueNotifier istVerbundenNotifier = ValueNotifier(false); final ValueNotifier fehlerNotifier = ValueNotifier(null); // ─── Konfiguration ─── /// Liest den aktuellen Sync-Modus aus SharedPreferences. static Future aktuellerModus() async { final p = await SharedPreferences.getInstance(); final modus = p.getString('sync_modus') ?? 'realtime'; switch (modus) { case 'manual': return SyncModus.manuell; case 'interval': return SyncModus.intervall; default: return SyncModus.echtzeit; } } /// Liest das Cloud-Intervall in Stunden (für Modus 3). /// Migriert alte `cloud_interval`-Werte (1/3/6/12h) auf /// den neuen Default 48h (2 Tage). static Future cloudIntervallStunden() async { final p = await SharedPreferences.getInstance(); // Neuer Key zuerst prüfen final neu = p.getInt('cloud_interval_stunden'); if (neu != null && neu > 0) return neu; // Migration vom alten Key final alt = p.getInt('cloud_interval') ?? 0; if (alt == 1 || alt == 3 || alt == 6 || alt == 12) { // Alte Stunden-Werte → auf neuen Default migrieren await p.setInt('cloud_interval_stunden', 48); MeloLogger().aktion('cloud_interval_migriert', {'alt': alt, 'neu': 48}); return 48; } // Default: 48h (2 Tage) await p.setInt('cloud_interval_stunden', 48); return 48; } // ─── Öffentliche API ─── /// Startet den Realtime-Sync, sofern der Modus 'realtime' aktiv ist /// und der Benutzer angemeldet ist. Future starteWennAktiviert() async { final modus = await aktuellerModus(); final user = _auth.benutzer; if (modus != SyncModus.echtzeit || user.isEmpty) { stoppe(); return; } // Prüfe, ob der Server SSE unterstützt (HEAD auf /subscribe). // 404 = kein SSE-Endpoint → Fallback auf Intervall. final kannSse = await _pruefeSseUnterstuetzung(); if (!kannSse) { MeloLogger().aktion('realtime_sse_nicht_verfuegbar', {}); // Fallback: sync_modus auf interval setzen + Auto-Sync-Timer starten final p = await SharedPreferences.getInstance(); await p.setString('sync_modus', 'interval'); await SyncService.starteAutoSyncTimer(); stoppe(); return; } await _verbinde(); } /// Prüft via HEAD-Request, ob der SSE-Endpoint existiert. Future _pruefeSseUnterstuetzung() async { try { final token = _auth.token; if (token == null || token.isEmpty) return false; final uri = Uri.parse('${AppConfig.cloudUrl}/api/v1/cloud/subscribe'); final client = HttpClient(); client.connectionTimeout = const Duration(seconds: 5); try { final request = await client.openUrl('HEAD', uri); request.headers.add('Authorization', 'Bearer $token'); final response = await request.close(); final ok = response.statusCode != 404; client.close(); return ok; } finally { client.close(); } } catch (_) { return false; } } /// Stellt die SSE-Verbindung her und beginnt den Event-Stream. Future _verbinde() async { if (_laeuft || _pausiert) return; _laeuft = true; try { final token = _auth.token; if (token == null || token.isEmpty) { _fehlerBericht('Kein Auth-Token'); return; } _httpClient?.close(); _httpClient = HttpClient(); _httpClient!.connectionTimeout = const Duration(seconds: 30); final uri = Uri.parse('${AppConfig.cloudUrl}/api/v1/cloud/subscribe'); final request = await _httpClient!.getUrl(uri); request.headers.add('Authorization', 'Bearer $token'); request.headers.add('Accept', 'text/event-stream'); request.headers.add('Cache-Control', 'no-cache'); _sseResponse = await request.close(); if (_sseResponse!.statusCode != 200) { _fehlerBericht('HTTP ${_sseResponse!.statusCode}'); await _reconnecteOderFallback(); return; } _reconnectVersuche = 0; istVerbundenNotifier.value = true; fehlerNotifier.value = null; MeloLogger().aktion('realtime_verbunden', {}); // Stream in UTF-8-Zeilen parsen _sseSubscription = _sseResponse! .transform(utf8.decoder) .transform(const LineSplitter()) .listen( verarbeiteSseZeile, onError: (e) { _fehlerBericht('Stream-Fehler: $e'); _reconnecteOderFallback(); }, onDone: () => _reconnecteOderFallback(), cancelOnError: false, ); } catch (e) { _fehlerBericht('Verbindungsfehler: $e'); await _reconnecteOderFallback(); } } /// SSE-Zeilen-Parser: akkumuliert `event:`- und `data:`-Zeilen /// und dispatched bei Leerzeile. @visibleForTesting void verarbeiteSseZeile(String zeile) { if (zeile.isEmpty) { // Leerzeile = Event-Ende → dispatch _aktuelleEventVerarbeiten(); _aktuellerEventTyp = null; return; } if (zeile.startsWith('event: ')) { _aktuellerEventTyp = zeile.substring(7).trim(); } else if (zeile.startsWith('data: ')) { final datenStr = zeile.substring(6).trim(); if (datenStr.isEmpty) return; try { final daten = jsonDecode(datenStr) as Map; _verarbeiteSseEvent(_aktuellerEventTyp, daten); } catch (e) { MeloLogger().fehler('realtime_parse', e); } } // Andere Zeilen (id:, retry:, Kommentare) ignorieren } /// Dispatch nach Leerzeile (falls unverarbeitete Daten übrig). void _aktuelleEventVerarbeiten() { // Keine Aktion nötig — data:-Zeilen werden sofort verarbeitet, // da sie direkt auf eine event:-Zeile folgen. } /// Verarbeitet ein einzelnes SSE-Event. void _verarbeiteSseEvent(String? eventTyp, Map daten) { // Heartbeat: ignorieren (hält nur die Verbindung offen) if (eventTyp == 'heartbeat') return; MeloLogger().aktion('realtime_event', { 'type': eventTyp, 'data': daten, }); switch (eventTyp) { case 'song_favorite': _handleFavoritEvent(daten); break; case 'song_delete': _handleDeleteEvent(daten); break; case 'song_update': _handleUpdateEvent(daten); break; case 'playlist_update': // Ignorieren — Playlisten werden beim nächsten vollständigen // Sync aktualisiert (zu bandwidth-intensiv für Echtzeit). break; } } /// Song-Favorit-Event: Server hat einen Favoriten-Status geändert. Future _handleFavoritEvent(Map daten) async { final songId = daten['song_id'] as String?; final isFavorite = daten['is_favorite'] as bool? ?? false; if (songId == null || songId.isEmpty) return; try { final song = await _db.songNachCloudId(songId); if (song == null || song.id == null) { // Song nicht lokal → beim nächsten Sync nachziehen return; } final istLokalFavorit = await _favoriten.istFavorit(song.id!); // Nur toggeln, wenn der lokale Zustand vom Server abweicht if (isFavorite != istLokalFavorit) { await _favoriten.umschalten(song.id!); } MeloLogger().aktion('realtime_favorit_angewendet', { 'songId': songId, 'isFavorite': isFavorite, }); } catch (e) { MeloLogger().fehler('realtime_favorit_error', e); } } /// Song-Lösch-Event: Server hat einen Song gelöscht (Tombstone). Future _handleDeleteEvent(Map daten) async { final songId = daten['song_id'] as String?; final deletedAt = daten['deleted_at'] as String?; if (songId == null || songId.isEmpty) return; // Tombstone-Check: nur löschen, wenn die Server-Löschung // neuer als der letzte lokale Sync ist. final letzterSync = await _letzterSync(); final geloeschtAm = deletedAt != null ? DateTime.tryParse(deletedAt) : null; // Alt-Tombstone (kein Zeitstempel) oder Erst-Sync → anwenden if (letzterSync != null && geloeschtAm != null && !geloeschtAm.isAfter(letzterSync)) { return; } try { final song = await _db.songNachCloudId(songId); if (song == null || song.id == null) return; // Datei löschen (falls lokal vorhanden) try { final datei = File(song.dateiPfad); if (song.dateiPfad.isNotEmpty && await datei.exists()) { await datei.delete(); } } catch (_) { // Datei-Fehler dürfen den Sync nicht abbrechen } await _db.loeschSong(song.id!); MeloLogger().aktion('realtime_song_geloescht', {'songId': songId}); } catch (e) { MeloLogger().fehler('realtime_delete_error', e); } } /// Song-Update-Event: Titel oder Künstler wurden auf dem Server geändert. Future _handleUpdateEvent(Map daten) async { final songId = daten['song_id'] as String?; final title = daten['title'] as String?; final artist = daten['artist'] as String?; if (songId == null || songId.isEmpty) return; if (title == null && artist == null) return; try { final song = await _db.songNachCloudId(songId); if (song == null || song.id == null) return; // Metadaten lokal aktualisieren await _db.metadatenAktualisieren( song.id!, titel: title, kuenstler: artist, ); MeloLogger().aktion('realtime_update', { 'songId': songId, 'title': title, 'artist': artist, }); } catch (e) { MeloLogger().fehler('realtime_update_error', e); } } /// Reconnect mit exponentiellem Backoff: 2s → 4s → 8s. /// Nach [maxReconnectVersuche] Fehlversuchen: Fallback auf Intervall-Modus. Future _reconnecteOderFallback() async { istVerbundenNotifier.value = false; _laeuft = false; _reconnectVersuche++; if (_reconnectVersuche <= _maxReconnectVersuche) { final backoff = backoffBerechnen(_reconnectVersuche); MeloLogger().aktion('realtime_reconnect', { 'attempt': _reconnectVersuche, 'delaySeconds': backoff.inSeconds, }); await Future.delayed(backoff); if (!_pausiert) { await _verbinde(); } } else { // Max Reconnects überschritten → Fallback auf Intervall-Modus _fehlerBericht('Max Reconnects erreicht — Fallback auf Intervall'); MeloLogger().aktion('realtime_fallback_to_interval', {}); final p = await SharedPreferences.getInstance(); await p.setString('sync_modus', 'interval'); await SyncService.starteAutoSyncTimer(); } } /// Berechnet das Backoff-Delay: 2^attempt * 2 Sekunden. /// Clamped auf max 8 Sekunden. Versuch <1 → 0 Sekunden (kein Shift-Fehler). @visibleForTesting Duration backoffBerechnen(int versuch) { if (versuch < 1) return Duration.zero; final sekunden = (_backoffStart.inSeconds * (1 << (versuch - 1))).clamp(0, 8); return Duration(seconds: sekunden); } /// Liest den letzten Sync-Zeitstempel aus SharedPreferences. Future _letzterSync() async { final p = await SharedPreferences.getInstance(); final ts = p.getString('cloud_last_sync_ts'); return ts != null ? DateTime.tryParse(ts) : null; } void _fehlerBericht(String nachricht) { MeloLogger().fehler('realtime_error', Exception(nachricht)); fehlerNotifier.value = nachricht; } /// Pausiert die SSE-Verbindung (App im Hintergrund → Batterie sparen). void pausiere() { if (_pausiert) return; _pausiert = true; _sseSubscription?.cancel(); _sseSubscription = null; _sseResponse = null; _httpClient?.close(force: true); _httpClient = null; _laeuft = false; istVerbundenNotifier.value = false; MeloLogger().aktion('realtime_pausiert', {}); } /// Setzt die SSE-Verbindung nach einer Pause fort. Future fortsetzen() async { if (!_pausiert) return; _pausiert = false; await starteWennAktiviert(); } /// Stoppt die SSE-Verbindung endgültig (z. B. Logout, Modus-Wechsel). void stoppe() { _sseSubscription?.cancel(); _sseSubscription = null; _sseResponse = null; _httpClient?.close(force: true); _httpClient = null; _laeuft = false; _pausiert = false; istVerbundenNotifier.value = false; fehlerNotifier.value = null; MeloLogger().aktion('realtime_gestoppt', {}); } }