using System.Net; using System.Net.Http.Json; using LehrerApp.Core.Models; using LehrerApp.Data; using LehrerApp.Sync.Crypto; using LehrerApp.Sync.Models; using Xunit; namespace LehrerApp.Sync.Tests; public sealed class SyncEngineTests { private static LiteDbContext NewInMemoryContext() => new(new MemoryStream()); private static readonly byte[] Key = SyncCrypto.GenerateKey(); [Fact] public async Task SyncNowAsync_InkompatibleServerVersion_BrichtVorWeiterenRequestsAb() { using var temp = new TempEventQueue(); var pending = temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", Guid.NewGuid().ToString(), "Save", "lokale-daten"); var handler = new FakeHttpMessageHandler(req => { var response = new HttpResponseMessage(HttpStatusCode.UpgradeRequired); response.Headers.Add(SyncProtocol.VersionHeaderName, "2"); return response; }); var engine = MakeEngine(temp, handler); var result = await engine.SyncNowAsync(); Assert.False(result.Success); Assert.Contains("Sync-Version 1", result.Reason); Assert.Contains("Server Version 2", result.Reason); Assert.Equal(SyncState.IncompatibleVersion, engine.Status.State); Assert.Equal(pending.EventId, Assert.Single(temp.Queue.GetPending()).EventId); var request = Assert.Single(handler.Requests); Assert.Equal("/api/sync/push", request.RequestUri!.AbsolutePath); Assert.Equal(SyncProtocol.CurrentVersion, request.Headers.GetValues(SyncProtocol.VersionHeaderName).Single()); } /// Regression: EventApplier.ApplyAsync fing früher nur LiteException ab. Jede andere Ausnahme /// (z.B. eine ungültige/korrupte Payload eines einzelnen Ereignisses) fiel unbehandelt aus /// SyncEngine.PullAsync heraus, BEVOR _queue.SetLastServerSeq() erreicht wurde — der nächste /// Sync-Versuch lud denselben Batch erneut und scheiterte am selben Ereignis erneut: ein /// dauerhaft blockierter Sync, bei dem selbst bereits im selben Batch erfolgreich angewendete /// Ereignisse (wie hier "goodStudent") nie durchkamen. [Fact] public async Task SyncNowAsync_KorruptesEreignisImPullBatch_BlockiertNachfolgendeEreignisseNicht() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var applier = new EventApplier(db, Key); var resolver = new ConflictResolver(temp.Queue); var attachments = new AttachmentSyncer(db, new HttpClient(), Key); var goodStudent = new Student { FirstName = "Anna", LastName = "Beispiel" }; var badEvent = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = Guid.NewGuid().ToString(), Operation = "Save", Payload = "offensichtlich-keine-gueltige-verschluesselte-payload", SequenceNr = 1, }; var goodEvent = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = goodStudent.Id.ToString(), Operation = "Save", Payload = SyncCrypto.EncryptObject(goodStudent, Key), SequenceNr = 2, }; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/pull") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse { Events = [badEvent, goodEvent], ServerSequenceNr = 2 }), }; return new HttpResponseMessage(HttpStatusCode.OK); }); var http = new HttpClient(handler) { BaseAddress = new Uri("https://example.invalid") }; var engine = new SyncEngine(temp.Queue, resolver, applier, attachments, http, new SyncConfig { DeviceId = "this-device" }); var result = await engine.SyncNowAsync(); Assert.True(result.Success); Assert.NotNull(db.Students.FindById(goodStudent.Id)); // Nicht bei 0 stecken geblieben - der Batch gilt als vollständig verarbeitet, auch wenn ein // einzelnes Ereignis darin übersprungen werden musste. Assert.Equal(2, temp.Queue.GetLastServerSeq()); } // ── DataChanged: Hook für den Desktop-Client, die sichtbare Seite neu zu laden (TODO 10.1.11) ─ [Fact] public async Task DataChanged_PullMitEreignissen_Feuert() { using var temp = new TempEventQueue(); var goodStudent = new Student { FirstName = "Anna", LastName = "Beispiel" }; var evt = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = goodStudent.Id.ToString(), Operation = "Save", Payload = SyncCrypto.EncryptObject(goodStudent, Key), SequenceNr = 1, }; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/pull") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse { Events = [evt], ServerSequenceNr = 1 }) }; return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse()) }; }); var engine = MakeEngine(temp, handler); var fired = false; engine.DataChanged += () => fired = true; await engine.SyncNowAsync(); Assert.True(fired); } [Fact] public async Task DataChanged_PullOhneEreignisse_FeuertNicht() { using var temp = new TempEventQueue(); var handler = new FakeHttpMessageHandler(req => new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(req.RequestUri!.AbsolutePath == "/api/sync/pull" ? new PullResponse() : (object)new PushResponse()), }); var engine = MakeEngine(temp, handler); var fired = false; engine.DataChanged += () => fired = true; await engine.SyncNowAsync(); Assert.False(fired); } /// Regression: PushAsync setzte den lokalen Pull-Cursor bisher direkt aus PushResponse. /// ServerSequenceNr - dem GLOBALEN Zähler über alle Geräte NACH diesem Push, nicht dem Stand, /// den DIESES Gerät tatsächlich per Pull erhalten hat. Hatte der Server zum Push-Zeitpunkt /// bereits ein noch nicht abgeholtes Ereignis eines ANDEREN Geräts mit niedrigerer ServerSeq, /// sprang der Cursor beim eigenen Push darüber hinweg - der direkt anschließende PullAsync /// fragte dann schon mit einem "since" danach und bekam 0 Ereignisse, obwohl das fremde /// Ereignis nie angewendet wurde. Genau das vom Nutzer beobachtete Symptom ("Client 1 konnte /// übermitteln, Client 2 bekommt weiterhin 'keine Änderungen'"), diesmal über einen anderen /// Pfad als die bereits behobene EventStore.Pull-Wasserzeichen-Berechnung. [Fact] public async Task SyncNowAsync_EigenerPushWaehrendFremdesEreignisNochAussteht_LiefertFremdesEreignisTrotzdem() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var applier = new EventApplier(db, Key, versions: temp.Queue); var goodStudent = new Student { FirstName = "Anna", LastName = "Beispiel" }; var foreignEvent = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = goodStudent.Id.ToString(), Operation = "Save", Payload = SyncCrypto.EncryptObject(goodStudent, Key), SequenceNr = 5, }; temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", Guid.NewGuid().ToString(), "Save", "x"); var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") // Globaler Höchststand (10) schließt das fremde, von DIESEM Gerät noch nicht // abgeholte Ereignis (Seq 5) bereits mit ein - genau das durfte SyncEngine NICHT // als eigenen Pull-Cursor übernehmen. return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 10 }) }; if (req.RequestUri!.AbsolutePath == "/api/sync/pull") { var query = req.RequestUri.Query.TrimStart('?').Split('&') .Select(p => p.Split('=')).ToDictionary(p => p[0], p => p[1]); var since = long.Parse(query["since"]); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(since < 5 ? new PullResponse { Events = [foreignEvent], ServerSequenceNr = 5 } : new PullResponse { ServerSequenceNr = since }), }; } return new HttpResponseMessage(HttpStatusCode.OK); }); var engine = MakeEngine(temp, handler, applier); var result = await engine.SyncNowAsync(); Assert.True(result.Success); Assert.NotNull(db.Students.FindById(goodStudent.Id)); } // ── PushAsync: Dedup, BasedOnServerSeq, AssignedServerSeqs (TODO 10.3.4) ──────────────── [Fact] public async Task PushAsync_MehrereAusstehendeEreignisseDerselbenEntitaet_SendetNurDasJuengste() { using var temp = new TempEventQueue(); var entityId = Guid.NewGuid().ToString(); temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", entityId, "Save", "alt"); var newest = temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", entityId, "Save", "neu"); List? pushed = null; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") { pushed = req.Content!.ReadFromJsonAsync>().GetAwaiter().GetResult(); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 1 }) }; } return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var engine = MakeEngine(temp, handler); var result = await engine.SyncNowAsync(); Assert.True(result.Success); var evt = Assert.Single(pushed!); Assert.Equal(newest.EventId, evt.EventId); Assert.Equal("neu", evt.Payload); Assert.Equal(0, temp.Queue.PendingCount()); } [Fact] public async Task SyncNowAsync_MehrAlsEinPushBatch_LeertQueueInEinemSync() { using var temp = new TempEventQueue(); for (var i = 0; i < SyncProtocol.PushBatchSize + 5; i++) temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", Guid.NewGuid().ToString(), "Save", $"payload-{i}"); var pushedBatchSizes = new List(); long serverSeq = 0; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") { var events = req.Content!.ReadFromJsonAsync>().GetAwaiter().GetResult()!; pushedBatchSizes.Add(events.Count); var assigned = events.ToDictionary(e => e.EventId, _ => ++serverSeq); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = serverSeq, AssignedServerSeqs = assigned }), }; } return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var engine = MakeEngine(temp, handler); var result = await engine.SyncNowAsync(); Assert.True(result.Success); Assert.Equal(SyncProtocol.PushBatchSize + 5, result.EventsPushed); Assert.Equal([SyncProtocol.PushBatchSize, 5], pushedBatchSizes); Assert.Equal(0, temp.Queue.PendingCount()); } [Fact] public async Task SyncNowAsync_VollerPullBatch_LaedtFolgebatchImSelbenSync() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var allEvents = Enumerable.Range(1, SyncProtocol.PullBatchSize + 3) .Select(sequence => { var student = new Student { FirstName = $"Vorname-{sequence}", LastName = "Beispiel" }; return new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = student.Id.ToString(), Operation = "Save", Payload = SyncCrypto.EncryptObject(student, Key), SequenceNr = sequence, }; }).ToList(); var pullRequests = 0; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/pull") { pullRequests++; var since = long.Parse(req.RequestUri.Query.TrimStart('?').Split('&') .Select(part => part.Split('=')).Single(part => part[0] == "since")[1]); var events = allEvents.Where(e => e.SequenceNr > since) .Take(SyncProtocol.PullBatchSize).ToList(); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse { Events = events, ServerSequenceNr = events.Count == 0 ? since : events[^1].SequenceNr, }), }; } return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse()) }; }); var engine = MakeEngine(temp, handler, new EventApplier(db, Key, versions: temp.Queue)); var dataChangedCount = 0; engine.DataChanged += () => dataChangedCount++; var result = await engine.SyncNowAsync(); Assert.True(result.Success); Assert.Equal(SyncProtocol.PullBatchSize + 3, result.EventsPulled); Assert.Equal(2, pullRequests); Assert.Equal(SyncProtocol.PullBatchSize + 3, db.Students.Count()); Assert.Equal(SyncProtocol.PullBatchSize + 3, temp.Queue.GetLastServerSeq()); Assert.Equal(1, dataChangedCount); } [Fact] public async Task PushAsync_SetztBasedOnServerSeqAusLokalerVersionsverfolgung() { using var temp = new TempEventQueue(); var entityId = Guid.NewGuid().ToString(); temp.Queue.SetKnownServerSeq("Lesson", entityId, 5); temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", entityId, "Save", "x"); List? pushed = null; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") { pushed = req.Content!.ReadFromJsonAsync>().GetAwaiter().GetResult(); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 6 }) }; } return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var engine = MakeEngine(temp, handler); await engine.SyncNowAsync(); Assert.Equal(5, Assert.Single(pushed!).BasedOnServerSeq); } [Fact] public async Task PushAsync_ErfolgreicherPush_AktualisiertLokaleVersionsverfolgung() { using var temp = new TempEventQueue(); var entityId = Guid.NewGuid().ToString(); var evt = temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Lesson", entityId, "Save", "x"); var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 7, AssignedServerSeqs = new() { [evt.EventId] = 7 } }), }; return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var engine = MakeEngine(temp, handler); await engine.SyncNowAsync(); Assert.Equal(7, temp.Queue.GetKnownServerSeq("Lesson", entityId)); } /// Regression: ein wegen neuerem Server-Stand abgelehnter Push blieb bisher unsichtbar - das /// Ereignis verschwand einfach nicht aus der Queue, ohne dass der Nutzer je erfuhr, dass es /// bereits neuere Daten gab (siehe TODO 10.3.4). RemoteWon-Fall: der Server-Stand gewinnt, wird /// sofort angewendet, und die lokale Änderung wird verworfen. [Fact] public async Task PushAsync_AbgelehnterPushRemoteWon_WendetServerStandAnUndVerwirftLokal() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var entityId = Guid.NewGuid(); // Companion verliert immer gegen Desktop (siehe ConflictResolver.DetermineWinner). var local = temp.Queue.Enqueue("companion-device", DeviceType.Companion, nameof(Student), entityId.ToString(), "Save", "veraltete-lokale-payload"); var remoteStudent = new Student { Id = entityId, FirstName = "Anna", LastName = "Beispiel" }; var remoteEvent = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Desktop, EntityType = nameof(Student), EntityId = entityId.ToString(), Operation = "Save", Payload = SyncCrypto.EncryptObject(remoteStudent, Key), SequenceNr = 99, }; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 99, ConflictingEventIds = [local.EventId] }), }; if (req.RequestUri!.AbsolutePath == $"/api/sync/entity/{nameof(Student)}/{entityId}") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(remoteEvent) }; return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var applier = new EventApplier(db, Key, versions: temp.Queue); var engine = MakeEngine(temp, handler, applier); var dataChanged = false; engine.DataChanged += () => dataChanged = true; await engine.SyncNowAsync(); var conflict = Assert.Single(temp.Queue.GetUnreviewed()); Assert.Equal("RemoteWon", conflict.Resolution); Assert.NotNull(db.Students.FindById(entityId)); Assert.Equal(0, temp.Queue.PendingCount()); // DataChanged muss auch bei einem RemoteWon-Konflikt feuern - dabei wird lokal genauso // Daten angewendet wie bei einem regulären Pull (TODO 10.1.11). Assert.True(dataChanged); Assert.Equal(99, temp.Queue.GetKnownServerSeq(nameof(Student), entityId.ToString())); } /// Gegenstück: LocalWon-Fall - die lokale Änderung bleibt (unbestätigt) in der Queue, damit der /// nächste Sync-Versuch sie mit dem soeben aktualisierten BasedOnServerSeq erneut versucht. [Fact] public async Task PushAsync_AbgelehnterPushLocalWon_BleibtUnbestaetigtInDerQueue() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var entityId = Guid.NewGuid(); // Desktop gewinnt immer gegen Companion (siehe ConflictResolver.DetermineWinner). var local = temp.Queue.Enqueue("this-device", DeviceType.Desktop, nameof(Student), entityId.ToString(), "Save", SyncCrypto.EncryptObject( new Student { Id = entityId, FirstName = "Lokal", LastName = "Beispiel" }, Key)); var remoteEvent = new SyncEvent { DeviceId = "other-device", DeviceType = DeviceType.Companion, EntityType = nameof(Student), EntityId = entityId.ToString(), Operation = "Save", Payload = "irrelevant-verliert-ohnehin", SequenceNr = 42, }; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PushResponse { ServerSequenceNr = 42, ConflictingEventIds = [local.EventId] }), }; if (req.RequestUri!.AbsolutePath == $"/api/sync/entity/{nameof(Student)}/{entityId}") return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(remoteEvent) }; return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var applier = new EventApplier(db, Key, versions: temp.Queue); var engine = MakeEngine(temp, handler, applier); await engine.SyncNowAsync(); var conflict = Assert.Single(temp.Queue.GetUnreviewed()); Assert.Equal("LocalWon", conflict.Resolution); Assert.Equal(1, temp.Queue.PendingCount()); Assert.Equal(42, temp.Queue.GetKnownServerSeq(nameof(Student), entityId.ToString())); } /// Regression (TODO 10.3.5): nach einem Kontowechsel (z.B. Login-Korrektur der Groß-/ /// Kleinschreibung, TODO 10.2.5) referenziert die lokale Versionsverfolgung ServerSeq-Werte /// eines FREMDEN Kontos. Der Server lehnt den Push ab (BasedOnServerSeq-Mismatch), kennt die /// Entität unter der aktuellen userId aber selbst gar nicht (404 bei GetLatestForEntity) - vor /// diesem Fix blieb das Ereignis dadurch dauerhaft und ohne jede Selbstheilung stecken (derselbe /// 404 bei jedem weiteren Sync-Versuch). Jetzt wird der stale Cache-Eintrag gelöscht, sodass /// der NÄCHSTE Push die Entität korrekt als neu behandelt und vom (für sie leeren) Server-Konto /// angenommen wird. [Fact] public async Task PushAsync_AbgelehnterPushServerKenntEntitaetNicht_LoeschtStaleCacheUndErholtSichSelbst() { using var temp = new TempEventQueue(); using var db = NewInMemoryContext(); var entityId = Guid.NewGuid().ToString(); // Stale Cache-Eintrag aus einem früheren (fremden) Konto - der Server unter der jetzigen // userId hat davon nie etwas gehört. temp.Queue.SetKnownServerSeq("Unit", entityId, 422); var local = temp.Queue.Enqueue("this-device", DeviceType.Desktop, "Unit", entityId, "Save", "x"); var pushAttempts = 0; var handler = new FakeHttpMessageHandler(req => { if (req.RequestUri!.AbsolutePath == "/api/sync/push") { pushAttempts++; // Erster Versuch: BasedOnServerSeq=422 (stale) passt nicht zum leeren Server-Konto // -> abgelehnt. Zweiter Versuch (nach Cache-Löschung): BasedOnServerSeq=null passt // zur ebenfalls unbekannten Entität -> angenommen. return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(pushAttempts == 1 ? new PushResponse { ServerSequenceNr = 0, ConflictingEventIds = [local.EventId] } : new PushResponse { ServerSequenceNr = 1, AssignedServerSeqs = new() { [local.EventId] = 1 }, }), }; } if (req.RequestUri!.AbsolutePath == $"/api/sync/entity/Unit/{entityId}") return new HttpResponseMessage(HttpStatusCode.NotFound); return new HttpResponseMessage(HttpStatusCode.OK) { Content = JsonContent.Create(new PullResponse()) }; }); var engine = MakeEngine(temp, handler); // SyncResult.Conflicts spiegelt nur PULL-Konflikte wider (siehe SyncEngine.SyncNowAsync) - // Push-Konflikte werden ausschließlich geloggt, daher hier über EventsPushed==0 geprüft. var first = await engine.SyncNowAsync(); Assert.Equal(0, first.EventsPushed); Assert.Empty(temp.Queue.GetUnreviewed()); Assert.Null(temp.Queue.GetKnownServerSeq("Unit", entityId)); Assert.Equal(1, temp.Queue.PendingCount()); var second = await engine.SyncNowAsync(); Assert.Equal(1, second.EventsPushed); Assert.Equal(0, temp.Queue.PendingCount()); } private static SyncEngine MakeEngine(TempEventQueue temp, FakeHttpMessageHandler handler, EventApplier? applier = null) { var db = NewInMemoryContext(); var http = new HttpClient(handler) { BaseAddress = new Uri("https://example.invalid") }; return new SyncEngine(temp.Queue, new ConflictResolver(temp.Queue), applier ?? new EventApplier(db, Key), new AttachmentSyncer(db, http, Key), http, new SyncConfig { DeviceId = "this-device" }); } private sealed class TempEventQueue : IDisposable { private readonly string _directory = Path.Combine( Path.GetTempPath(), $"lehrerapp-sync-tests-engine-{Guid.NewGuid():N}"); public EventQueue Queue { get; } public TempEventQueue() { Directory.CreateDirectory(_directory); Queue = new EventQueue(Path.Combine(_directory, "queue.db")); } public void Dispose() { Queue.Dispose(); if (Directory.Exists(_directory)) Directory.Delete(_directory, recursive: true); } } }