using System.Net; using System.Net.Http.Json; using LehrerApp.Core.Services; using LehrerApp.Sync.Models; namespace LehrerApp.Sync; /// /// Push/Pull Orchestrierung. /// Automatisch alle N Minuten + manuell per SyncNowAsync(). /// public class SyncEngine : IDisposable { private readonly EventQueue _queue; private readonly ConflictResolver _resolver; private readonly EventApplier _applier; private readonly AttachmentSyncer _attachments; private readonly HttpClient _http; private readonly SyncConfig _config; private readonly AppLogger? _logger; private readonly Timer _timer; private readonly SemaphoreSlim _syncGate = new(1, 1); public SyncStatus Status { get; private set; } = new(); public event Action? StatusChanged; /// /// Feuert, wenn ein Pull tatsächlich Ereignisse angewendet hat (siehe TODO 10.1.11) - der /// Desktop-Client abonniert das, um die gerade sichtbare Seite neu zu laden, da /// EventApplier absichtlich an jedem ViewModel vorbei direkt auf die LiteDB schreibt (siehe /// EventApplier-Klassenkommentar). /// public event Action? DataChanged; public SyncEngine(EventQueue queue, ConflictResolver resolver, EventApplier applier, AttachmentSyncer attachments, HttpClient http, SyncConfig config, AppLogger? logger = null) { _queue = queue; _resolver = resolver; _applier = applier; _attachments = attachments; _http = http; _config = config; _logger = logger; _timer = new Timer( async _ => await SyncNowAsync(true), null, TimeSpan.FromMinutes(config.AutoSyncIntervalMinutes), TimeSpan.FromMinutes(config.AutoSyncIntervalMinutes)); UpdateStatus(); } public async Task SyncNowAsync(bool isAutomatic = false) { if (!await _syncGate.WaitAsync(0)) return new() { Skipped = true, Reason = "Sync bereits aktiv" }; try { return await RunSyncAsync(isAutomatic); } finally { _syncGate.Release(); } } /// /// Wartet kurz auf noch laufende lokale Speicheroperationen und führt anschließend garantiert /// einen Sync aus. Anders als wird bei einem bereits laufenden Sync /// nicht abgebrochen, sondern gewartet; so bleiben Änderungen, die während dieses Laufs in die /// Outbox gelangen, beim Beenden nicht zurück. /// public async Task SyncBeforeShutdownAsync(TimeSpan delay) { if (delay > TimeSpan.Zero) await Task.Delay(delay); await _syncGate.WaitAsync(); try { return await RunSyncAsync(isAutomatic: true); } finally { _syncGate.Release(); } } private async Task RunSyncAsync(bool isAutomatic) { SetState(SyncState.Syncing); _logger?.Info($"Sync: Start ({(isAutomatic ? "automatisch" : "manuell")}), Gerät={_config.DeviceId}, " + $"{_queue.PendingCount()} lokal ausstehend"); try { var (pushed, pushConflicts) = await PushAsync(); await _attachments.UploadPendingAsync(_queue); var (pulled, conflicts) = await PullAsync(); _queue.SetLastSyncAt(DateTime.UtcNow); SetState(SyncState.Idle); _logger?.Info($"Sync: Fertig - {pushed} gepusht ({pushConflicts} Push-Konflikte), " + $"{pulled} gepullt ({conflicts} Pull-Konflikte)"); return new() { Success = true, EventsPushed = pushed, EventsPulled = pulled, Conflicts = conflicts }; } catch (SyncProtocolMismatchException ex) { _logger?.Error("Sync wegen inkompatibler Protokollversion abgebrochen", ex); SetState(SyncState.IncompatibleVersion, ex.Message); return new() { Reason = ex.Message }; } catch (HttpRequestException ex) { // Sammelt sowohl echte Netzwerkfehler als auch nicht-erfolgreiche HTTP-Antworten // (EnsureSuccessStatusCode() in Push/PullAsync) unter demselben "Offline"-Status - // ohne Log wäre ein z.B. 401/500 vom Server nicht von "kein Internet" unterscheidbar. _logger?.Error("Sync fehlgeschlagen (HTTP)", ex); SetState(SyncState.Offline); return new() { Reason = "Server nicht erreichbar" }; } catch (Exception ex) { _logger?.Error("Sync fehlgeschlagen", ex); SetState(SyncState.Error, ex.Message); return new() { Reason = ex.Message }; } } private async Task<(int Pushed, int Conflicts)> PushAsync() { var pending = DeduplicatePending(); if (pending.Count == 0) return (0, 0); // BasedOnServerSeq erst unmittelbar vor dem Senden setzen (nicht beim Enqueue) - zwischen // Enqueue und Push kann ein Pull den lokal bekannten Stand dieser Entität aktualisiert // haben (siehe EventApplier.ApplyAsync). foreach (var evt in pending) evt.BasedOnServerSeq = _queue.GetKnownServerSeq(evt.EntityType, evt.EntityId); _logger?.Info($"Sync: Push - {pending.Count} Ereignis(se) ausstehend: " + string.Join(", ", pending.Select(e => $"{e.EntityType}/{e.Operation}"))); using var request = SyncProtocol.CreateRequest(HttpMethod.Post, "/api/sync/push"); request.Content = JsonContent.Create(pending); using var resp = await _http.SendAsync(request); await SyncProtocol.EnsureCompatibleSuccessAsync(resp); var result = await resp.Content.ReadFromJsonAsync(); if (result is null) { _logger?.Warn("Sync: Push - leere Server-Antwort."); return (0, 0); } _queue.Acknowledge(pending .Where(e => !result.ConflictingEventIds.Contains(e.EventId)) .Select(e => e.EventId)); foreach (var evt in pending) if (result.AssignedServerSeqs.TryGetValue(evt.EventId, out var seq)) _queue.SetKnownServerSeq(evt.EntityType, evt.EntityId, seq); // ANDERS ALS FRÜHER wird der lokale Pull-Cursor (_queue.SetLastServerSeq) hier NICHT aus // result.ServerSequenceNr gesetzt: das ist der GLOBALE Zähler über ALLE Geräte NACH diesem // Push, nicht der Stand, den DIESES Gerät tatsächlich per Pull erhalten hat. Hatte der // Server zu diesem Zeitpunkt bereits ein noch nicht abgeholtes Ereignis eines ANDEREN // Geräts mit niedrigerer ServerSeq, würde der Cursor hier darüber hinwegspringen - der // direkt anschließende PullAsync würde dann schon mit einem "since" danach fragen und 0 // Ereignisse zurückbekommen, ohne das fremde Ereignis je angewendet zu haben (exakt das // Symptom, das die Ereignis-Anzahl-Korrektur in EventStore.Pull eigentlich beheben sollte - // hier aber über einen anderen Pfad wieder hereinkam). PullAsync pflegt den Cursor bereits // korrekt selbst, ausschließlich anhand tatsächlich zugestellter Ereignisse. _logger?.Info($"Sync: Push - {pending.Count - result.ConflictingEventIds.Count} vom Server " + $"angenommen (aktueller globaler Stand ServerSequenceNr={result.ServerSequenceNr}), " + $"{result.ConflictingEventIds.Count} abgelehnt (Konflikt)."); if (result.ConflictingEventIds.Count > 0) await HandleRejectedAsync(pending.Where(e => result.ConflictingEventIds.Contains(e.EventId))); return (pending.Count - result.ConflictingEventIds.Count, result.ConflictingEventIds.Count); } // Payload ist immer ein vollständiges Entitäts-Snapshot (nie ein Delta, siehe // SyncEventPublisher) - mehrere ausstehende lokale Änderungen derselben Entität lassen sich // deshalb gefahrlos auf das jüngste zusammenfassen, bevor gepusht wird. Wichtig auch für die // exakte BasedOnServerSeq-Prüfung: ohne Dedup könnten zwei Ereignisse derselben Entität im // selben Batch mit demselben (veralteten) BasedOnServerSeq ankommen und sich gegenseitig ins // Aus laufen. private List DeduplicatePending() { var pending = _queue.GetPending(); if (pending.Count == 0) return pending; var latest = pending .GroupBy(e => (e.EntityType, e.EntityId)) .Select(g => g.OrderBy(e => e.SequenceNr).Last()) .ToList(); var superseded = pending.Except(latest).Select(e => e.EventId).ToList(); if (superseded.Count > 0) { _queue.Acknowledge(superseded); _logger?.Info($"Sync: Push - {superseded.Count} veraltete Ereignis(se) derselben " + "Entität lokal zusammengefasst (Full-Snapshot)."); } return latest; } /// /// Ein Push wurde abgelehnt, weil der Server für diese Entität bereits einen neueren Stand /// hat (BasedOnServerSeq-Mismatch). Lädt den aktuellen Server-Stand sofort nach (statt auf /// den nächsten regulären Pull zu warten), löst den Konflikt nach derselben Politik wie /// auf und legt in jedem Fall einen ConflictEntry an, damit der /// Nutzer sieht, dass hier bereits neuere Daten vorlagen - unabhängig davon, ob die lokale /// Änderung verworfen wird oder nicht. /// private async Task HandleRejectedAsync(IEnumerable rejected) { foreach (var local in rejected) { HttpResponseMessage resp; try { using var request = SyncProtocol.CreateRequest(HttpMethod.Get, $"/api/sync/entity/{local.EntityType}/{local.EntityId}"); resp = await _http.SendAsync(request); } catch (Exception ex) { _logger?.Error($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} - " + "aktueller Server-Stand konnte nicht nachgeladen werden.", ex); continue; } if (resp.StatusCode == HttpStatusCode.NotFound) { // Der Server kennt diese Entität unter der aktuellen userId gar nicht - der lokal // zwischengespeicherte BasedOnServerSeq war stale (z.B. nach einem Kontowechsel, // siehe TODO 10.3.5, oder wenn dieses Gerät die Entität nie zuvor unter diesem // Konto gepusht hat). Cache löschen, statt endlos mit demselben 404 zu scheitern - // der nächste Push behandelt die Entität dann korrekt als neu für dieses Konto. _queue.ClearKnownServerSeq(local.EntityType, local.EntityId); _logger?.Warn($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} - Server " + "kennt diese Entität nicht. Lokale Versionsverfolgung zurückgesetzt, " + "nächster Push behandelt sie als neu."); continue; } await SyncProtocol.EnsureCompatibleSuccessAsync(resp); var remote = await resp.Content.ReadFromJsonAsync(); if (remote is null) { _logger?.Warn($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId}, aber " + "kein Server-Stand gefunden - übersprungen."); continue; } _queue.SetKnownServerSeq(local.EntityType, local.EntityId, remote.SequenceNr); var winner = ConflictResolver.DetermineWinner(local, remote); _queue.AddConflict(new ConflictEntry { LocalEvent = local, RemoteEvent = remote, Resolution = winner == local ? "LocalWon" : "RemoteWon", }); _logger?.Info($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} aufgelöst - " + $"Server hatte bereits neueren Stand (ServerSeq={remote.SequenceNr}), " + $"Auflösung={(winner == local ? "LocalWon" : "RemoteWon")}."); if (winner == remote) { await _applier.ApplyAsync(remote); _queue.Acknowledge([local.EventId]); DataChanged?.Invoke(); } // LocalWon: Ereignis bleibt unbestätigt in der Queue - der nächste PushAsync-Lauf // versucht es erneut, jetzt mit dem soeben aktualisierten BasedOnServerSeq. } } private async Task<(int Pulled, int Conflicts)> PullAsync() { var since = _queue.GetLastServerSeq(); _logger?.Info($"Sync: Pull - frage Server nach Ereignissen seit ServerSequenceNr={since}."); using var request = SyncProtocol.CreateRequest(HttpMethod.Get, $"/api/sync/pull?since={since}&deviceId={_config.DeviceId}"); using var response = await _http.SendAsync(request); await SyncProtocol.EnsureCompatibleSuccessAsync(response); var resp = await response.Content.ReadFromJsonAsync(); if (resp is null || resp.Events.Count == 0) { _logger?.Info("Sync: Pull - keine neuen Ereignisse vom Server."); return (0, 0); } _logger?.Info($"Sync: Pull - {resp.Events.Count} Ereignis(se) vom Server erhalten: " + string.Join(", ", resp.Events.Select(e => $"{e.EntityType}/{e.Operation}"))); var conflicts = 0; foreach (var evt in resp.Events) { var c = _resolver.TryResolve(evt, _config.DeviceId); if (c is null) { await _applier.ApplyAsync(evt); continue; } _queue.AddConflict(c); conflicts++; _logger?.Info($"Sync: Pull - Konflikt bei {evt.EntityType}/{evt.EntityId}, " + $"Auflösung={c.Resolution}"); if (c.Resolution == "RemoteWon") await _applier.ApplyAsync(evt); } _queue.SetLastServerSeq(resp.ServerSequenceNr); DataChanged?.Invoke(); return (resp.Events.Count, conflicts); } private void SetState(SyncState state, string? error = null) { Status = new SyncStatus { State = state, LastSyncAt = _queue.GetLastSyncAt(), PendingEvents = _queue.PendingCount(), ConflictCount = _queue.ConflictCount(), ErrorMessage = error, }; StatusChanged?.Invoke(Status); } private void UpdateStatus() => SetState(Status.State); public void Dispose() { _timer.Dispose(); _syncGate.Dispose(); _queue.Dispose(); } } public class SyncConfig { public string ServerUrl { get; set; } = ""; public string DeviceId { get; set; } = ""; public DeviceType DeviceType { get; set; } = DeviceType.Desktop; public int AutoSyncIntervalMinutes { get; set; } = 5; } public class SyncResult { public bool Success { get; set; } public bool Skipped { get; set; } public string? Reason { get; set; } public int EventsPushed { get; set; } public int EventsPulled { get; set; } public int Conflicts { get; set; } }