CRITICAL (C1): RealtimeSyncService jetzt Singleton (Factory-Pattern) - Verhindert doppelte SSE-Verbindungen durch home_screen + settings - UI-Notifier (istVerbundenNotifier) app-weit konsistent HIGH (H1): _reconnecteOderFallback + starteWennAktiviert setzen cloud_interval vor starteAutoSyncTimer() — sonst bleibt Timer bei Fallback auf 0 (kein Sync) HIGH (H2): Race-Condition in _verbinde() nach await request.close() — _pausiert-Flag re-check; bei true → clean exit mit _laeuft=false HIGH (H3): fortsetzen() re-checkt _pausiert nach starteWennAktiviert() — verhindert Leak bei schnellem App-Umschalten MEDIUM (M1): _syncModusSetzen stoppt SSE bei Wechsel auf manual/interval (vorher lief die Verbindung weiter) MEDIUM (M2): Event-Logging loggt nur noch Event-Typ, nicht rohe Daten (Privacy: Song-IDs/Titel/Künstler nicht im Log) LOW (L1): _verbinde() setzt _laeuft=false bei Token-Fehler (kein Deadlock)
482 lines
16 KiB
Dart
482 lines
16 KiB
Dart
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)
|
|
///
|
|
/// **Singleton**: Nur EINE Instanz pro App-Lebenszyklus, damit
|
|
/// home_screen.dart und settings_screen.dart dieselbe Verbindung teilen.
|
|
class RealtimeSyncService {
|
|
static final RealtimeSyncService _instanz = RealtimeSyncService._();
|
|
factory RealtimeSyncService() => _instanz;
|
|
RealtimeSyncService._()
|
|
: _db = DbHelper(),
|
|
_auth = AuthService(),
|
|
_favoriten = FavoritenService();
|
|
|
|
final DbHelper _db;
|
|
final AuthService _auth;
|
|
final FavoritenService _favoriten;
|
|
|
|
// ─── Verbindungs-Zustand ───
|
|
|
|
StreamSubscription<String>? _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<bool> istVerbundenNotifier = ValueNotifier(false);
|
|
final ValueNotifier<String?> fehlerNotifier = ValueNotifier(null);
|
|
|
|
// ─── Konfiguration ───
|
|
|
|
/// Liest den aktuellen Sync-Modus aus SharedPreferences.
|
|
static Future<SyncModus> 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<int> 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<void> 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');
|
|
// cloud_interval setzen, damit starteAutoSyncTimer() den Timer startet
|
|
final intervallStunden = await cloudIntervallStunden();
|
|
await p.setInt('cloud_interval', intervallStunden);
|
|
await SyncService.starteAutoSyncTimer();
|
|
stoppe();
|
|
return;
|
|
}
|
|
|
|
await _verbinde();
|
|
}
|
|
|
|
/// Prüft via HEAD-Request, ob der SSE-Endpoint existiert.
|
|
Future<bool> _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<void> _verbinde() async {
|
|
if (_laeuft || _pausiert) return;
|
|
_laeuft = true;
|
|
|
|
try {
|
|
final token = _auth.token;
|
|
if (token == null || token.isEmpty) {
|
|
_fehlerBericht('Kein Auth-Token');
|
|
_laeuft = false;
|
|
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();
|
|
|
|
// Nach await prüfen: pausiere() könnte während des Connects
|
|
// aufgerufen worden sein (WidgetsBindingObserver-Race).
|
|
if (_pausiert) {
|
|
_sseResponse = null;
|
|
_laeuft = false;
|
|
return;
|
|
}
|
|
|
|
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<String, dynamic>;
|
|
_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<String, dynamic> daten) {
|
|
// Heartbeat: ignorieren (hält nur die Verbindung offen)
|
|
if (eventTyp == 'heartbeat') return;
|
|
|
|
MeloLogger().aktion('realtime_event', {
|
|
'type': eventTyp,
|
|
// Daten nicht loggen — könnten sensitive User-Informationen enthalten
|
|
// (Song-IDs, Titel, Künstler). Nur Event-Typ wird geloggt.
|
|
});
|
|
|
|
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<void> _handleFavoritEvent(Map<String, dynamic> 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<void> _handleDeleteEvent(Map<String, dynamic> 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<void> _handleUpdateEvent(Map<String, dynamic> 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<void> _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');
|
|
// cloud_interval (alter Key) auf Default 48h setzen, damit
|
|
// SyncService.starteAutoSyncTimer() den Timer auch wirklich startet.
|
|
final intervallStunden = await cloudIntervallStunden();
|
|
await p.setInt('cloud_interval', intervallStunden);
|
|
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<DateTime?> _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<void> fortsetzen() async {
|
|
if (!_pausiert) return;
|
|
_pausiert = false;
|
|
await starteWennAktiviert();
|
|
// Nach await prüfen: pausiere() könnte während starteWennAktiviert()
|
|
// aufgerufen worden sein (z.B. schnelles App-Umschalten).
|
|
// Falls ja, Verbindung sofort wieder trennen.
|
|
if (_pausiert) {
|
|
stoppe();
|
|
}
|
|
}
|
|
|
|
/// 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', {});
|
|
}
|
|
}
|