Files
LehrerApp/LehrerApp.Sync/SyncEngine.cs
T
adminandClaude Sonnet 5 ce4dfb0197 Baustein 5: EventApplier (Inbound Apply) + Loop-Prevention (Kapitel 10)
Die zentrale, bisher komplett fehlende Luecke: selbst mit Baustein
1-4 haette SyncEngine.PullAsync empfangene Ereignisse nur zur
Konflikterkennung genutzt, nie in die lokale LiteDB geschrieben -
ankommende Aenderungen von anderen Geraeten waeren nirgends sichtbar
geworden.

Neu EventApplier: entschluesselt, dispatcht ueber eine explizite
EntityType-Tabelle, schreibt IMMER direkt auf die rohe LiteDB-
Collection, nie ueber eine Repository-Save/Delete-Methode - sonst
wuerde der OnChange-Hook (Baustein 2) die gerade angewendete Aenderung
als neues ausgehendes Ereignis re-enqueuen (Sync-Ping-Pong). Ein
gemeinsames Suppress-Flag wurde geprueft und verworfen: SyncEngine
laeuft per Timer nebenlaeufig zum UI-Thread, ein Flag koennte einen
echten Nutzer-Save waehrenddessen verschlucken. Der direkte Collection-
Zugriff ist zustandslos und dadurch korrekt. Kaskaden-Faelle nutzen
dieselben internen LiteDbContext-Hilfsmethoden wie die Repositories
(Baustein 4).

Neu SyncEventPublisher, der den OnChange-Hook in ein verschluesseltes
EventQueue.Enqueue uebersetzt (an LiteDbContext.OnChange gehaengt).

Mit dediziertem Loop-Prevention-Test abgesichert: belegt mit echtem
LiteDbContext + OnChange-Zaehler, dass Apply keinen neuen Hook-Aufruf
ausloest.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-17 11:27:12 +02:00

122 lines
4.4 KiB
C#

using System.Net.Http.Json;
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 HttpClient _http;
private readonly SyncConfig _config;
private readonly Timer _timer;
public SyncStatus Status { get; private set; } = new();
public event Action<SyncStatus>? StatusChanged;
public SyncEngine(EventQueue queue, ConflictResolver resolver, EventApplier applier,
HttpClient http, SyncConfig config)
{
_queue = queue;
_resolver = resolver;
_applier = applier;
_http = http;
_config = config;
_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);
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<PushResponse>();
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<PullResponse>(
$"/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 null) { _applier.Apply(evt); continue; }
_queue.AddConflict(c);
conflicts++;
if (c.Resolution == "RemoteWon") _applier.Apply(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; }
}