327 lines
15 KiB
C#
327 lines
15 KiB
C#
using System.Net;
|
|
using System.Net.Http.Json;
|
|
using LehrerApp.Core.Services;
|
|
using LehrerApp.Sync.Models;
|
|
|
|
namespace LehrerApp.Sync;
|
|
|
|
/// <summary>
|
|
/// Push/Pull Orchestrierung.
|
|
/// Automatisch alle N Minuten + manuell per SyncNowAsync().
|
|
/// </summary>
|
|
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<SyncStatus>? StatusChanged;
|
|
/// <summary>
|
|
/// 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).
|
|
/// </summary>
|
|
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<SyncResult> 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();
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Wartet kurz auf noch laufende lokale Speicheroperationen und führt anschließend garantiert
|
|
/// einen Sync aus. Anders als <see cref="SyncNowAsync"/> 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.
|
|
/// </summary>
|
|
public async Task<SyncResult> 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<SyncResult> 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<PushResponse>();
|
|
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<SyncEvent> 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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// 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
|
|
/// <see cref="ConflictResolver"/> 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.
|
|
/// </summary>
|
|
private async Task HandleRejectedAsync(IEnumerable<SyncEvent> 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<SyncEvent>();
|
|
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<PullResponse>();
|
|
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; }
|
|
}
|