using LiteDB; using LehrerApp.Sync.Models; namespace LehrerApp.Api; /// /// Append-only Event-Log pro User. Server versteht Payload nicht. /// public class EventStore(string dataPath) : IDisposable { private readonly Dictionary _dbs = new(); private readonly Lock _lock = new(); /// /// Nimmt Ereignisse an, wenn ihre exakt der aktuellen /// ServerSeq der jeweiligen Entität entspricht (null == Entität hier noch nie gesehen, z.B. /// Neuanlage) — echte optimistische Nebenläufigkeitskontrolle statt der früheren 30-Sekunden- /// Heuristik ("hat ein anderes Gerät kürzlich dieselbe Entität angefasst"), die sowohl falsch- /// positiv (zwei Geräte bearbeiten zufällig kurz hintereinander verschiedene Felder) als auch /// falsch-negativ (echter Konflikt liegt außerhalb des 30s-Fensters) sein konnte. /// Mehrere Ereignisse derselben Entität IM SELBEN Aufruf bauen bewusst aufeinander auf (das /// zweite prüft gegen den vom ersten gerade neu vergebenen Stand) — der Client schickt ohnehin /// nur noch das jüngste ausstehende Ereignis je Entität (siehe SyncEngine.PushAsync). /// public PushResponse Push(string userId, List events) { var col = GetCol(userId); var seq = LastSeq(col); var rejects = new List(); var assigned = new Dictionary(); foreach (var e in events.OrderBy(e => e.Timestamp)) { var currentSeq = LatestForEntity(col, e.EntityType, e.EntityId)?.ServerSeq; if (currentSeq != e.BasedOnServerSeq) { rejects.Add(e.EventId); continue; } var newSeq = ++seq; col.Insert(new ServerEvent { EventId = e.EventId, DeviceId = e.DeviceId, DeviceType = e.DeviceType, Timestamp = e.Timestamp, ClientSeq = e.SequenceNr, ServerSeq = newSeq, EntityType = e.EntityType, EntityId = e.EntityId, Operation = e.Operation, Payload = e.Payload }); assigned[e.EventId] = newSeq; } return new() { Success = true, ServerSequenceNr = seq, ConflictingEventIds = rejects, AssignedServerSeqs = assigned }; } /// Aktuellstes Ereignis einer einzelnen Entität — für Clients, deren Push wegen eines /// neueren Server-Stands abgelehnt wurde (siehe ), um sofort den aktuellen /// Stand nachzuladen, statt auf den nächsten regulären Pull zu warten. public SyncEvent? GetLatestForEntity(string userId, string entityType, string entityId) { var col = GetCol(userId); var e = LatestForEntity(col, entityType, entityId); // SequenceNr trägt hier (wie bei Pull) die ServerSeq — dieselbe Konvention wie sonst im // Sync-Protokoll: bei vom Server stammenden Ereignissen ist SequenceNr immer die ServerSeq. return e is null ? null : new SyncEvent { EventId = e.EventId, DeviceId = e.DeviceId, DeviceType = e.DeviceType, Timestamp = e.Timestamp, SequenceNr = e.ServerSeq, EntityType = e.EntityType, EntityId = e.EntityId, Operation = e.Operation, Payload = e.Payload }; } private static ServerEvent? LatestForEntity(ILiteCollection col, string entityType, string entityId) => col.Find(x => x.EntityType == entityType && x.EntityId == entityId) .OrderByDescending(x => x.ServerSeq).FirstOrDefault(); public PullResponse Pull(string userId, long since, string requestingDeviceId) { var col = GetCol(userId); var events = col.Find(e => e.ServerSeq > since && e.DeviceId != requestingDeviceId) .OrderBy(e => e.ServerSeq).Take(500) .Select(e => new SyncEvent { EventId = e.EventId, DeviceId = e.DeviceId, DeviceType = e.DeviceType, Timestamp = e.Timestamp, SequenceNr = e.ServerSeq, EntityType = e.EntityType, EntityId = e.EntityId, Operation = e.Operation, Payload = e.Payload }) .ToList(); // ServerSequenceNr MUSS die höchste ServerSeq unter den tatsächlich zurückgegebenen // Ereignissen sein, NICHT der globale Höchststand (LastSeq(col)) — der schließt auch // Ereignisse ANDERER Geräte ein, die z.B. gerade erst (nach dem obigen Find-Aufruf, aber // vor dieser Zeile) eingetroffen sind, oder — der Bug, der hier tatsächlich beobachtet // wurde — Ereignisse des anfragenden Geräts selbst, die oben bewusst per // "DeviceId != requestingDeviceId" herausgefiltert wurden. SyncEngine.PullAsync übernimmt // ServerSequenceNr 1:1 als neuen "since"-Cursor für den nächsten Pull; mit dem globalen // Höchststand würde der Client seinen Cursor über Ereignisse hinweg vorrücken, die er nie // erhalten hat, und sie dauerhaft verpassen — genau das vom Nutzer beobachtete Symptom // (Push meldet Erfolg, Pull liefert 0 Ereignisse, obwohl welche ausstehen). var newWatermark = events.Count > 0 ? events.Max(e => e.SequenceNr) : since; return new() { Events = events, ServerSequenceNr = newWatermark }; } private ILiteCollection GetCol(string userId) { lock (_lock) { if (!_dbs.TryGetValue(userId, out var db)) { var safe = string.Concat(userId.Where(c => char.IsLetterOrDigit(c) || c == '-')); db = new LiteDatabase(Path.Combine(dataPath, $"{safe}.db")); _dbs[userId] = db; } var col = db.GetCollection("events"); col.EnsureIndex(x => x.ServerSeq); return col; } } private static long LastSeq(ILiteCollection col) { var last = col.FindOne(Query.All(nameof(ServerEvent.ServerSeq), Query.Descending)); return last?.ServerSeq ?? 0; } public void Dispose() { foreach (var db in _dbs.Values) db.Dispose(); } } internal class ServerEvent { public Guid EventId { get; set; } public string DeviceId { get; set; } = ""; public DeviceType DeviceType { get; set; } public DateTime Timestamp { get; set; } public long ClientSeq { get; set; } public long ServerSeq { get; set; } public string EntityType { get; set; } = ""; public string EntityId { get; set; } = ""; public string Operation { get; set; } = ""; public string Payload { get; set; } = ""; }