using System.Net.Http.Json; 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 HttpClient _http; private readonly SyncConfig _config; private readonly Timer _timer; public SyncStatus Status { get; private set; } = new(); public event Action? StatusChanged; public SyncEngine(EventQueue queue, ConflictResolver resolver, HttpClient http, SyncConfig config) { _queue = queue; _resolver = resolver; _http = http; _config = config; _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 (Status.State == SyncState.Syncing) return new() { Skipped = true, Reason = "Sync bereits aktiv" }; SetState(SyncState.Syncing); try { var (pushed, _) = await PushAsync(); var (pulled, conflicts) = await PullAsync(); _queue.SetLastSyncAt(DateTime.UtcNow); SetState(SyncState.Idle); return new() { Success = true, EventsPushed = pushed, EventsPulled = pulled, Conflicts = conflicts }; } catch (HttpRequestException) { SetState(SyncState.Offline); return new() { Reason = "Server nicht erreichbar" }; } catch (Exception ex) { SetState(SyncState.Error, ex.Message); return new() { Reason = ex.Message }; } } private async Task<(int Pushed, int Conflicts)> PushAsync() { var pending = _queue.GetPending(); if (pending.Count == 0) return (0, 0); var resp = await _http.PostAsJsonAsync("/api/sync/push", pending); resp.EnsureSuccessStatusCode(); var result = await resp.Content.ReadFromJsonAsync(); if (result is null) return (0, 0); _queue.Acknowledge(pending .Where(e => !result.ConflictingEventIds.Contains(e.EventId)) .Select(e => e.EventId)); _queue.SetLastServerSeq(result.ServerSequenceNr); return (pending.Count - result.ConflictingEventIds.Count, result.ConflictingEventIds.Count); } private async Task<(int Pulled, int Conflicts)> PullAsync() { var since = _queue.GetLastServerSeq(); var resp = await _http.GetFromJsonAsync( $"/api/sync/pull?since={since}&deviceId={_config.DeviceId}"); if (resp is null || resp.Events.Count == 0) return (0, 0); var conflicts = 0; foreach (var evt in resp.Events) { var c = _resolver.TryResolve(evt, _config.DeviceId); if (c is not null) { _queue.AddConflict(c); conflicts++; } } _queue.SetLastServerSeq(resp.ServerSequenceNr); 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(); _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; } }