Files
LehrerApp/LehrerApp.Sync/EventQueue.cs
T
adminandClaude Sonnet 5 5bd6967421 fix: Pull-Wasserzeichen sprang über fremde Ereignisse + schärfere Push-Kollisionskontrolle
EventStore.Pull gab bisher den globalen ServerSeq-Höchststand als neuen Cursor zurück statt
den höchsten unter den tatsächlich gelieferten Ereignissen - hatte ein Gerät selbst kurz zuvor
etwas gepusht, sprang sein Pull-Cursor über noch nicht abgeholte Ereignisse anderer Geräte
hinweg und verpasste sie dauerhaft, ohne jeden Fehler.

Ersetzt außerdem die bisherige 30-Sekunden-Heuristik zur Konflikterkennung beim Push durch
exakte BasedOnServerSeq-Prüfung: jedes SyncEvent trägt die ServerSeq, auf der es aufbaut: der
Server lehnt ab, wenn der aktuelle Stand nicht mehr passt. Bei Ablehnung lädt der Client sofort
den neuen Server-Stand nach, löst den Konflikt nach der bestehenden Desktop-vs-Companion/
Timestamp-Politik auf und macht ihn immer in der Konflikt-Review-UI sichtbar, statt die
verworfene Änderung stillschweigend zu verlieren.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-18 23:25:04 +02:00

129 lines
5.2 KiB
C#

using LiteDB;
using LehrerApp.Sync.Models;
namespace LehrerApp.Sync;
/// <summary>
/// Lokale Event-Queue in LiteDB. Puffert Events bis sie
/// erfolgreich zum Server gepusht wurden.
/// </summary>
public class EventQueue : IDisposable
{
private readonly LiteDatabase _db;
private readonly ILiteCollection<SyncEvent> _queue;
private readonly ILiteCollection<SyncMeta> _meta;
private readonly ILiteCollection<ConflictEntry> _conflicts;
private readonly ILiteCollection<PendingAttachmentUpload> _attachmentUploads;
private readonly ILiteCollection<EntityVersion> _entityVersions;
private long _currentSeq;
public EventQueue(string path)
{
_db = new LiteDatabase(path);
_queue = _db.GetCollection<SyncEvent>("queue");
_meta = _db.GetCollection<SyncMeta>("meta");
_conflicts = _db.GetCollection<ConflictEntry>("conflicts");
_attachmentUploads = _db.GetCollection<PendingAttachmentUpload>("attachment_uploads");
_entityVersions = _db.GetCollection<EntityVersion>("entity_versions");
_attachmentUploads.EnsureIndex(x => x.StorageId, unique: true);
_queue.EnsureIndex(x => x.SequenceNr);
_currentSeq = _meta.FindById("seq")?.Value ?? 0;
}
public SyncEvent Enqueue(string deviceId, DeviceType deviceType,
string entityType, string entityId, string operation, string payload)
{
var evt = new SyncEvent
{
DeviceId = deviceId,
DeviceType = deviceType,
Timestamp = DateTime.UtcNow,
SequenceNr = ++_currentSeq,
EntityType = entityType,
EntityId = entityId,
Operation = operation,
Payload = payload,
};
_queue.Insert(evt);
_meta.Upsert(new SyncMeta { Id = "seq", Value = _currentSeq });
return evt;
}
public List<SyncEvent> GetPending(int max = 200) =>
_queue.Find(Query.All(nameof(SyncEvent.SequenceNr))).Take(max).ToList();
public int PendingCount() => _queue.Count();
public void Acknowledge(IEnumerable<Guid> ids) { foreach (var id in ids) _queue.Delete(id); }
public long GetLastServerSeq() => _meta.FindById("serverSeq")?.Value ?? 0;
public void SetLastServerSeq(long nr) => _meta.Upsert(new SyncMeta { Id = "serverSeq", Value = nr });
public DateTime? GetLastSyncAt() => _meta.FindById("lastSync")?.Timestamp;
public void SetLastSyncAt(DateTime dt) =>
_meta.Upsert(new SyncMeta { Id = "lastSync", Value = 0, Timestamp = dt });
public void AddConflict(ConflictEntry c) => _conflicts.Insert(c);
public List<ConflictEntry> GetUnreviewed() => _conflicts.Find(c => !c.Reviewed).ToList();
public int ConflictCount() => _conflicts.Count(c => !c.Reviewed);
public void MarkReviewed(Guid id)
{
var conflict = _conflicts.FindById(id);
if (conflict is null) return;
conflict.Reviewed = true;
_conflicts.Update(conflict);
}
// ── Lokale Versionsverfolgung je Entität (optimistische Nebenläufigkeitskontrolle) ──────
// Merkt sich pro Entität die zuletzt bekannte ServerSeq — Grundlage für SyncEvent.
// BasedOnServerSeq beim Push (siehe SyncEngine.PushAsync) und dafür, wie ein abgelehnter
// Push nach dem Nachladen des aktuellen Server-Stands aufgelöst wird.
public long? GetKnownServerSeq(string entityType, string entityId) =>
_entityVersions.FindById(EntityVersionKey(entityType, entityId))?.ServerSeq;
public void SetKnownServerSeq(string entityType, string entityId, long serverSeq) =>
_entityVersions.Upsert(new EntityVersion
{ Key = EntityVersionKey(entityType, entityId), ServerSeq = serverSeq });
private static string EntityVersionKey(string entityType, string entityId) => $"{entityType}:{entityId}";
// ── Anhang-Warteliste (getrennt von der JSON-Ereignis-Outbox, siehe AttachmentSyncer) ────
public void QueueAttachmentUpload(string storageId)
{
if (!_attachmentUploads.Exists(a => a.StorageId == storageId))
_attachmentUploads.Insert(new PendingAttachmentUpload { StorageId = storageId });
}
public List<string> GetPendingAttachmentUploads() =>
_attachmentUploads.FindAll().Select(a => a.StorageId).ToList();
public void MarkAttachmentUploaded(string storageId) =>
_attachmentUploads.DeleteMany(a => a.StorageId == storageId);
public void Dispose() => _db.Dispose();
}
public class ConflictEntry
{
public Guid Id { get; init; } = Guid.NewGuid();
public DateTime DetectedAt { get; init; } = DateTime.UtcNow;
public SyncEvent LocalEvent { get; init; } = null!;
public SyncEvent RemoteEvent { get; init; } = null!;
public string Resolution { get; init; } = "";
public bool Reviewed { get; set; }
}
internal class SyncMeta
{
public string Id { get; set; } = "";
public long Value { get; set; }
public DateTime? Timestamp { get; set; }
}
internal class PendingAttachmentUpload
{
public ObjectId Id { get; set; } = ObjectId.NewObjectId();
public string StorageId { get; set; } = "";
}
internal class EntityVersion
{
[BsonId]
public string Key { get; set; } = "";
public long ServerSeq { get; set; }
}