diff --git a/LehrerApp.Api/EventStore.cs b/LehrerApp.Api/EventStore.cs index 3e1c317..fa9060d 100644 --- a/LehrerApp.Api/EventStore.cs +++ b/LehrerApp.Api/EventStore.cs @@ -1,4 +1,5 @@ using LiteDB; +using LehrerApp.Sync; using LehrerApp.Sync.Models; namespace LehrerApp.Api; @@ -67,7 +68,7 @@ public class EventStore(string dataPath) : IDisposable { var col = GetCol(userId); var events = col.Find(e => e.ServerSeq > since && e.DeviceId != requestingDeviceId) - .OrderBy(e => e.ServerSeq).Take(500) + .OrderBy(e => e.ServerSeq).Take(SyncProtocol.PullBatchSize) .Select(e => new SyncEvent { EventId = e.EventId, DeviceId = e.DeviceId, DeviceType = e.DeviceType, Timestamp = e.Timestamp, SequenceNr = e.ServerSeq, EntityType = e.EntityType, diff --git a/LehrerApp.Sync.Tests/SyncEngineTests.cs b/LehrerApp.Sync.Tests/SyncEngineTests.cs index 74c25f0..180ca4c 100644 --- a/LehrerApp.Sync.Tests/SyncEngineTests.cs +++ b/LehrerApp.Sync.Tests/SyncEngineTests.cs @@ -229,6 +229,95 @@ public sealed class SyncEngineTests 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() { diff --git a/LehrerApp.Sync/EventQueue.cs b/LehrerApp.Sync/EventQueue.cs index f05a42c..058f14b 100644 --- a/LehrerApp.Sync/EventQueue.cs +++ b/LehrerApp.Sync/EventQueue.cs @@ -49,8 +49,14 @@ public class EventQueue : IDisposable return evt; } - public List GetPending(int max = 200) => - _queue.Find(Query.All(nameof(SyncEvent.SequenceNr))).Take(max).ToList(); + public List GetPending(int max = SyncProtocol.PushBatchSize, + IReadOnlySet? excludedEventIds = null) + { + var pending = _queue.Find(Query.All(nameof(SyncEvent.SequenceNr))).AsEnumerable(); + if (excludedEventIds is not null) + pending = pending.Where(e => !excludedEventIds.Contains(e.EventId)); + return pending.Take(max).ToList(); + } public int PendingCount() => _queue.Count(); public void Acknowledge(IEnumerable ids) { foreach (var id in ids) _queue.Delete(id); } public long GetLastServerSeq() => _meta.FindById("serverSeq")?.Value ?? 0; diff --git a/LehrerApp.Sync/SyncEngine.cs b/LehrerApp.Sync/SyncEngine.cs index ce56581..7f8aec4 100644 --- a/LehrerApp.Sync/SyncEngine.cs +++ b/LehrerApp.Sync/SyncEngine.cs @@ -86,14 +86,37 @@ public class SyncEngine : IDisposable private async Task RunSyncAsync(bool isAutomatic) { + var pulledDataChanged = false; SetState(SyncState.Syncing); _logger?.Info($"Sync: Start ({(isAutomatic ? "automatisch" : "manuell")}), Gerät={_config.DeviceId}, " + $"{_queue.PendingCount()} lokal ausstehend"); try { - var (pushed, pushConflicts) = await PushAsync(); + var pushed = 0; + var pushConflicts = 0; + var deferredPushEvents = new HashSet(); + PushBatchResult pushBatch; + do + { + pushBatch = await PushAsync(deferredPushEvents); + pushed += pushBatch.Pushed; + pushConflicts += pushBatch.Conflicts; + deferredPushEvents.UnionWith(pushBatch.DeferredEventIds); + } while (pushBatch.SourceEventCount == SyncProtocol.PushBatchSize); + await _attachments.UploadPendingAsync(_queue); - var (pulled, conflicts) = await PullAsync(); + + var pulled = 0; + var conflicts = 0; + PullBatchResult pullBatch; + do + { + pullBatch = await PullAsync(); + pulled += pullBatch.Pulled; + conflicts += pullBatch.Conflicts; + pulledDataChanged |= pullBatch.Pulled > 0; + } while (pullBatch.Pulled == SyncProtocol.PullBatchSize); + _queue.SetLastSyncAt(DateTime.UtcNow); SetState(SyncState.Idle); _logger?.Info($"Sync: Fertig - {pushed} gepusht ({pushConflicts} Push-Konflikte), " + @@ -121,12 +144,19 @@ public class SyncEngine : IDisposable SetState(SyncState.Error, ex.Message); return new() { Reason = ex.Message }; } + finally + { + // Auch wenn ein späterer Batch fehlschlägt, wurden frühere Batches bereits dauerhaft + // angewendet. Die Oberfläche muss diesen erfolgreich übernommenen Zwischenstand dann + // trotzdem neu laden; bei einem erfolgreichen Drain feuert der Hook genau einmal. + if (pulledDataChanged) DataChanged?.Invoke(); + } } - private async Task<(int Pushed, int Conflicts)> PushAsync() + private async Task PushAsync(IReadOnlySet excludedEventIds) { - var pending = DeduplicatePending(); - if (pending.Count == 0) return (0, 0); + var pending = DeduplicatePending(excludedEventIds, out var sourceEventCount); + if (pending.Count == 0) return new(0, 0, sourceEventCount, []); // 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). @@ -139,7 +169,11 @@ public class SyncEngine : IDisposable using var resp = await _http.SendAsync(request); await SyncProtocol.EnsureCompatibleSuccessAsync(resp); var result = await resp.Content.ReadFromJsonAsync(); - if (result is null) { _logger?.Warn("Sync: Push - leere Server-Antwort."); return (0, 0); } + if (result is null) + { + _logger?.Warn("Sync: Push - leere Server-Antwort."); + return new(0, 0, sourceEventCount, pending.Select(e => e.EventId).ToList()); + } _queue.Acknowledge(pending .Where(e => !result.ConflictingEventIds.Contains(e.EventId)) .Select(e => e.EventId)); @@ -161,8 +195,8 @@ public class SyncEngine : IDisposable $"{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); + return new(pending.Count - result.ConflictingEventIds.Count, + result.ConflictingEventIds.Count, sourceEventCount, result.ConflictingEventIds); } // Payload ist immer ein vollständiges Entitäts-Snapshot (nie ein Delta, siehe @@ -171,9 +205,11 @@ public class SyncEngine : IDisposable // 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 DeduplicatePending() + private List DeduplicatePending(IReadOnlySet excludedEventIds, + out int sourceEventCount) { - var pending = _queue.GetPending(); + var pending = _queue.GetPending(SyncProtocol.PushBatchSize, excludedEventIds); + sourceEventCount = pending.Count; if (pending.Count == 0) return pending; var latest = pending .GroupBy(e => (e.EntityType, e.EntityId)) @@ -259,7 +295,7 @@ public class SyncEngine : IDisposable } } - private async Task<(int Pulled, int Conflicts)> PullAsync() + private async Task PullAsync() { var since = _queue.GetLastServerSeq(); _logger?.Info($"Sync: Pull - frage Server nach Ereignissen seit ServerSequenceNr={since}."); @@ -271,7 +307,7 @@ public class SyncEngine : IDisposable if (resp is null || resp.Events.Count == 0) { _logger?.Info("Sync: Pull - keine neuen Ereignisse vom Server."); - return (0, 0); + return new(0, 0); } _logger?.Info($"Sync: Pull - {resp.Events.Count} Ereignis(se) vom Server erhalten: " + string.Join(", ", resp.Events.Select(e => $"{e.EntityType}/{e.Operation}"))); @@ -287,10 +323,13 @@ public class SyncEngine : IDisposable if (c.Resolution == "RemoteWon") await _applier.ApplyAsync(evt); } _queue.SetLastServerSeq(resp.ServerSequenceNr); - DataChanged?.Invoke(); - return (resp.Events.Count, conflicts); + return new(resp.Events.Count, conflicts); } + private readonly record struct PushBatchResult(int Pushed, int Conflicts, + int SourceEventCount, IReadOnlyCollection DeferredEventIds); + private readonly record struct PullBatchResult(int Pulled, int Conflicts); + private void SetState(SyncState state, string? error = null) { Status = new SyncStatus diff --git a/LehrerApp.Sync/SyncProtocol.cs b/LehrerApp.Sync/SyncProtocol.cs index 6b20135..b18fc7b 100644 --- a/LehrerApp.Sync/SyncProtocol.cs +++ b/LehrerApp.Sync/SyncProtocol.cs @@ -14,6 +14,8 @@ public static class SyncProtocol { public const string CurrentVersion = "1"; public const string VersionHeaderName = "X-LehrerApp-Sync-Version"; + public const int PushBatchSize = 200; + public const int PullBatchSize = 500; public static HttpRequestMessage CreateRequest(HttpMethod method, string requestUri) {