85 lines
3.2 KiB
C#
85 lines
3.2 KiB
C#
using LiteDB;
|
|
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();
|
|
|
|
public PushResponse Push(string userId, List<SyncEvent> events)
|
|
{
|
|
var col = GetCol(userId);
|
|
var seq = LastSeq(col);
|
|
var rejects = new List<Guid>();
|
|
foreach (var e in events.OrderBy(e => e.Timestamp))
|
|
{
|
|
var recent = col.FindOne(x =>
|
|
x.EntityType == e.EntityType && x.EntityId == e.EntityId &&
|
|
x.DeviceId != e.DeviceId && x.Timestamp > e.Timestamp.AddSeconds(-30));
|
|
if (recent is not null) { rejects.Add(e.EventId); continue; }
|
|
col.Insert(new ServerEvent { EventId = e.EventId, DeviceId = e.DeviceId,
|
|
DeviceType = e.DeviceType, Timestamp = e.Timestamp,
|
|
ClientSeq = e.SequenceNr, ServerSeq = ++seq,
|
|
EntityType = e.EntityType, EntityId = e.EntityId,
|
|
Operation = e.Operation, Payload = e.Payload });
|
|
}
|
|
return new() { Success = true, ServerSequenceNr = seq, ConflictingEventIds = rejects };
|
|
}
|
|
|
|
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();
|
|
return new() { Events = events, ServerSequenceNr = LastSeq(col) };
|
|
}
|
|
|
|
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; } = "";
|
|
}
|