Files
LehrerApp/LehrerApp.Api/EventStore.cs
2026-08-25 21:30:00 +02:00

129 lines
6.5 KiB
C#

using LiteDB;
using LehrerApp.Sync;
using LehrerApp.Sync.Models;
namespace LehrerApp.Api;
/// <summary>
/// Append-only Event-Log pro User. Server versteht Payload nicht.
/// </summary>
public class EventStore(string dataPath) : IDisposable
{
private readonly Dictionary<string, LiteDatabase> _dbs = new();
private readonly Lock _lock = new();
/// <summary>
/// Nimmt Ereignisse an, wenn ihre <see cref="SyncEvent.BasedOnServerSeq"/> 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).
/// </summary>
public PushResponse Push(string userId, List<SyncEvent> events)
{
var col = GetCol(userId);
var seq = LastSeq(col);
var rejects = new List<Guid>();
var assigned = new Dictionary<Guid, long>();
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 };
}
/// <summary>Aktuellstes Ereignis einer einzelnen Entität — für Clients, deren Push wegen eines
/// neueren Server-Stands abgelehnt wurde (siehe <see cref="Push"/>), um sofort den aktuellen
/// Stand nachzuladen, statt auf den nächsten regulären Pull zu warten.</summary>
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<ServerEvent> 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(SyncProtocol.PullBatchSize)
.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<ServerEvent> 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<ServerEvent>("events");
col.EnsureIndex(x => x.ServerSeq);
return col;
}
}
private static long LastSeq(ILiteCollection<ServerEvent> 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; } = "";
}