Files
LehrerApp/LehrerApp.Sync/SyncEngine.cs
T
adminandClaude Sonnet 5 e7e5faeba8 fix: eigener Push ließ Client seinen Pull-Cursor über fremde Ereignisse springen
SyncEngine.PushAsync setzte den lokalen Pull-Cursor bisher aus PushResponse.ServerSequenceNr -
dem globalen Zähler über alle Geräte nach dem eigenen Push, nicht dem tatsächlich zugestellten
Stand. War beim Server zu diesem Zeitpunkt bereits ein noch nicht abgeholtes Ereignis eines
anderen Geräts mit niedrigerer ServerSeq vorhanden, sprang der Cursor darüber hinweg und der
direkt folgende Pull bekam 0 Ereignisse, ohne es je angewendet zu haben - ein zweiter,
unabhängiger Cursor-Bug mit demselben Symptom wie der vorherige Wasserzeichen-Fix, diesmal
client- statt serverseitig.

Ergänzt außerdem einen "Vollständigen Sync erzwingen"-Button in den Sync-Einstellungen, damit
bereits durch diesen Bug zu weit vorgerückte Geräte ihren Fortschritt manuell zurücksetzen und
alle Ereignisse erneut laden können.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-18 23:38:46 +02:00

255 lines
12 KiB
C#

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;
public SyncStatus Status { get; private set; } = new();
public event Action<SyncStatus>? StatusChanged;
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 (Status.State == SyncState.Syncing)
return new() { Skipped = true, Reason = "Sync bereits aktiv" };
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 (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}")));
var resp = await _http.PostAsJsonAsync("/api/sync/push", pending);
resp.EnsureSuccessStatusCode();
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)
{
SyncEvent? remote;
try
{
remote = await _http.GetFromJsonAsync<SyncEvent>(
$"/api/sync/entity/{local.EntityType}/{local.EntityId}");
}
catch (Exception ex)
{
_logger?.Error($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} - " +
"aktueller Server-Stand konnte nicht nachgeladen werden.", ex);
continue;
}
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]);
}
// 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}.");
var resp = await _http.GetFromJsonAsync<PullResponse>(
$"/api/sync/pull?since={since}&deviceId={_config.DeviceId}");
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);
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; }
}