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>
This commit is contained in:
@@ -16,17 +16,7 @@ public class ConflictResolver(EventQueue queue)
|
||||
&& e.DeviceId != remote.DeviceId);
|
||||
if (local is null) return null;
|
||||
|
||||
var winner = (local.DeviceType, remote.DeviceType) switch
|
||||
{
|
||||
(DeviceType.Desktop, DeviceType.Companion) => local,
|
||||
(DeviceType.Companion, DeviceType.Desktop) => remote,
|
||||
// ToUniversalTime(): LiteDB liefert DateTime beim Auslesen aus der Queue als Kind=Local
|
||||
// zurück (Ticks werden dabei um die lokale Zeitzone verschoben). DateTime-Vergleiche
|
||||
// berücksichtigen Kind nicht, sondern vergleichen nur rohe Ticks — ein direkter Vergleich
|
||||
// von local.Timestamp (Local, aus der Queue) mit remote.Timestamp (Utc, vom Server) wäre
|
||||
// daher außerhalb von UTC+0 falsch.
|
||||
_ => local.Timestamp.ToUniversalTime() >= remote.Timestamp.ToUniversalTime() ? local : remote,
|
||||
};
|
||||
var winner = DetermineWinner(local, remote);
|
||||
|
||||
if (winner == remote) queue.Acknowledge([local.EventId]);
|
||||
|
||||
@@ -37,4 +27,23 @@ public class ConflictResolver(EventQueue queue)
|
||||
Resolution = winner == local ? "LocalWon" : "RemoteWon",
|
||||
};
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Desktop schlägt Companion; bei gleichem Gerätetyp gewinnt der spätere Timestamp. Auch von
|
||||
/// SyncEngine für die Ablehnungsbehandlung eines per BasedOnServerSeq abgelehnten Push
|
||||
/// genutzt — dieselbe Politik unabhängig davon, ob der Konflikt beim Pull (gleichzeitig
|
||||
/// eingetroffenes fremdes Ereignis) oder erst durch eine Server-Ablehnung entdeckt wurde.
|
||||
/// </summary>
|
||||
public static SyncEvent DetermineWinner(SyncEvent local, SyncEvent remote) =>
|
||||
(local.DeviceType, remote.DeviceType) switch
|
||||
{
|
||||
(DeviceType.Desktop, DeviceType.Companion) => local,
|
||||
(DeviceType.Companion, DeviceType.Desktop) => remote,
|
||||
// ToUniversalTime(): LiteDB liefert DateTime beim Auslesen aus der Queue als Kind=Local
|
||||
// zurück (Ticks werden dabei um die lokale Zeitzone verschoben). DateTime-Vergleiche
|
||||
// berücksichtigen Kind nicht, sondern vergleichen nur rohe Ticks — ein direkter Vergleich
|
||||
// von local.Timestamp (Local, aus der Queue) mit remote.Timestamp (Utc, vom Server) wäre
|
||||
// daher außerhalb von UTC+0 falsch.
|
||||
_ => local.Timestamp.ToUniversalTime() >= remote.Timestamp.ToUniversalTime() ? local : remote,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -24,7 +24,8 @@ namespace LehrerApp.Sync;
|
||||
/// Pfad bewusst NICHT geprüft (v1-Einschränkung, siehe TODO.md 10.3) — nur harte LiteDB-Unique-
|
||||
/// Constraints greifen noch und führen zum Überspringen des einzelnen Ereignisses.
|
||||
/// </summary>
|
||||
public class EventApplier(LiteDbContext db, byte[] syncKey, HttpClient? http = null, AppLogger? logger = null)
|
||||
public class EventApplier(LiteDbContext db, byte[] syncKey, HttpClient? http = null, AppLogger? logger = null,
|
||||
EventQueue? versions = null)
|
||||
{
|
||||
private static readonly Dictionary<string, EntityHandler> Handlers = BuildHandlers();
|
||||
|
||||
@@ -42,6 +43,10 @@ public class EventApplier(LiteDbContext db, byte[] syncKey, HttpClient? http = n
|
||||
handler(db, evt.Operation, evt.EntityId, json);
|
||||
if (evt.EntityType == nameof(Documentation) && evt.Operation != "Delete" && http is not null)
|
||||
await DownloadMissingAttachmentsAsync(json);
|
||||
// evt.SequenceNr trägt bei einem vom Server empfangenen Ereignis immer dessen
|
||||
// ServerSeq (siehe EventStore.Pull) — Grundlage für BasedOnServerSeq beim nächsten
|
||||
// eigenen Push dieser Entität (optimistische Nebenläufigkeitskontrolle, TODO 10.3.4).
|
||||
versions?.SetKnownServerSeq(evt.EntityType, evt.EntityId, evt.SequenceNr);
|
||||
logger?.Info($"Sync: Ereignis angewendet - {evt.EntityType} {evt.Operation} EntityId={evt.EntityId}");
|
||||
}
|
||||
catch (LiteException ex)
|
||||
|
||||
@@ -14,6 +14,7 @@ public class EventQueue : IDisposable
|
||||
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)
|
||||
@@ -23,6 +24,7 @@ public class EventQueue : IDisposable
|
||||
_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;
|
||||
@@ -67,6 +69,20 @@ public class EventQueue : IDisposable
|
||||
_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)
|
||||
{
|
||||
@@ -103,3 +119,10 @@ 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; }
|
||||
}
|
||||
|
||||
@@ -21,6 +21,14 @@ public class SyncEvent
|
||||
public string Operation { get; init; } = "";
|
||||
/// <summary>Verschlüsselt (Desktop) oder Klartext (Companion/WebApp).</summary>
|
||||
public string Payload { get; init; } = "";
|
||||
/// <summary>
|
||||
/// ServerSeq, auf der die lokale Änderung aufbaut (null = Entität wurde hier noch nie
|
||||
/// synchronisiert, z.B. Neuanlage). Wird von SyncEngine.PushAsync erst unmittelbar vor dem
|
||||
/// Senden aus der lokalen Versionsverfolgung (EventQueue) gesetzt, nicht beim Einreihen —
|
||||
/// so verwendet ein zweiter Push desselben Ereignisses (falls der erste abgelehnt wurde)
|
||||
/// automatisch den inzwischen aktualisierten Stand.
|
||||
/// </summary>
|
||||
public long? BasedOnServerSeq { get; set; }
|
||||
}
|
||||
|
||||
/// <summary>WebApp/Companion-Event: Payload ist Klartext-JSON.</summary>
|
||||
@@ -43,6 +51,10 @@ public class PushResponse
|
||||
public bool Success { get; init; }
|
||||
public long ServerSequenceNr { get; init; }
|
||||
public List<Guid> ConflictingEventIds { get; init; } = [];
|
||||
/// <summary>Je akzeptiertem Ereignis die tatsächlich vergebene ServerSeq — der Client braucht
|
||||
/// das, um seine lokale Versionsverfolgung je Entität (EventQueue) auf den neuen Stand zu
|
||||
/// bringen.</summary>
|
||||
public Dictionary<Guid, long> AssignedServerSeqs { get; init; } = [];
|
||||
}
|
||||
public class PullResponse
|
||||
{
|
||||
|
||||
@@ -76,8 +76,13 @@ public class SyncEngine : IDisposable
|
||||
|
||||
private async Task<(int Pushed, int Conflicts)> PushAsync()
|
||||
{
|
||||
var pending = _queue.GetPending();
|
||||
var pending = DeduplicatePending();
|
||||
if (pending.Count == 0) return (0, 0);
|
||||
// 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).
|
||||
foreach (var evt in pending)
|
||||
evt.BasedOnServerSeq = _queue.GetKnownServerSeq(evt.EntityType, evt.EntityId);
|
||||
_logger?.Info($"Sync: Push - {pending.Count} Ereignis(se) ausstehend: " +
|
||||
string.Join(", ", pending.Select(e => $"{e.EntityType}/{e.Operation}")));
|
||||
var resp = await _http.PostAsJsonAsync("/api/sync/push", pending);
|
||||
@@ -87,13 +92,95 @@ public class SyncEngine : IDisposable
|
||||
_queue.Acknowledge(pending
|
||||
.Where(e => !result.ConflictingEventIds.Contains(e.EventId))
|
||||
.Select(e => e.EventId));
|
||||
foreach (var evt in pending)
|
||||
if (result.AssignedServerSeqs.TryGetValue(evt.EventId, out var seq))
|
||||
_queue.SetKnownServerSeq(evt.EntityType, evt.EntityId, seq);
|
||||
_queue.SetLastServerSeq(result.ServerSequenceNr);
|
||||
_logger?.Info($"Sync: Push - vom Server bestätigt bis ServerSequenceNr={result.ServerSequenceNr}, " +
|
||||
$"{result.ConflictingEventIds.Count} vom Server 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);
|
||||
}
|
||||
|
||||
// Payload ist immer ein vollständiges Entitäts-Snapshot (nie ein Delta, siehe
|
||||
// SyncEventPublisher) - mehrere ausstehende lokale Änderungen derselben Entität lassen sich
|
||||
// deshalb gefahrlos auf das jüngste zusammenfassen, bevor gepusht wird. Wichtig auch für die
|
||||
// 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<SyncEvent> DeduplicatePending()
|
||||
{
|
||||
var pending = _queue.GetPending();
|
||||
if (pending.Count == 0) return pending;
|
||||
var latest = pending
|
||||
.GroupBy(e => (e.EntityType, e.EntityId))
|
||||
.Select(g => g.OrderBy(e => e.SequenceNr).Last())
|
||||
.ToList();
|
||||
var superseded = pending.Except(latest).Select(e => e.EventId).ToList();
|
||||
if (superseded.Count > 0)
|
||||
{
|
||||
_queue.Acknowledge(superseded);
|
||||
_logger?.Info($"Sync: Push - {superseded.Count} veraltete Ereignis(se) derselben " +
|
||||
"Entität lokal zusammengefasst (Full-Snapshot).");
|
||||
}
|
||||
return latest;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Ein Push wurde abgelehnt, weil der Server für diese Entität bereits einen neueren Stand
|
||||
/// hat (BasedOnServerSeq-Mismatch). Lädt den aktuellen Server-Stand sofort nach (statt auf
|
||||
/// den nächsten regulären Pull zu warten), löst den Konflikt nach derselben Politik wie
|
||||
/// <see cref="ConflictResolver"/> auf und legt in jedem Fall einen ConflictEntry an, damit der
|
||||
/// Nutzer sieht, dass hier bereits neuere Daten vorlagen - unabhängig davon, ob die lokale
|
||||
/// Änderung verworfen wird oder nicht.
|
||||
/// </summary>
|
||||
private async Task HandleRejectedAsync(IEnumerable<SyncEvent> rejected)
|
||||
{
|
||||
foreach (var local in rejected)
|
||||
{
|
||||
SyncEvent? remote;
|
||||
try
|
||||
{
|
||||
remote = await _http.GetFromJsonAsync<SyncEvent>(
|
||||
$"/api/sync/entity/{local.EntityType}/{local.EntityId}");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger?.Error($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} - " +
|
||||
"aktueller Server-Stand konnte nicht nachgeladen werden.", ex);
|
||||
continue;
|
||||
}
|
||||
if (remote is null)
|
||||
{
|
||||
_logger?.Warn($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId}, aber " +
|
||||
"kein Server-Stand gefunden - übersprungen.");
|
||||
continue;
|
||||
}
|
||||
|
||||
_queue.SetKnownServerSeq(local.EntityType, local.EntityId, remote.SequenceNr);
|
||||
var winner = ConflictResolver.DetermineWinner(local, remote);
|
||||
_queue.AddConflict(new ConflictEntry
|
||||
{
|
||||
LocalEvent = local,
|
||||
RemoteEvent = remote,
|
||||
Resolution = winner == local ? "LocalWon" : "RemoteWon",
|
||||
});
|
||||
_logger?.Info($"Sync: Push-Konflikt bei {local.EntityType}/{local.EntityId} aufgelöst - " +
|
||||
$"Server hatte bereits neueren Stand (ServerSeq={remote.SequenceNr}), " +
|
||||
$"Auflösung={(winner == local ? "LocalWon" : "RemoteWon")}.");
|
||||
|
||||
if (winner == remote)
|
||||
{
|
||||
await _applier.ApplyAsync(remote);
|
||||
_queue.Acknowledge([local.EventId]);
|
||||
}
|
||||
// LocalWon: Ereignis bleibt unbestätigt in der Queue - der nächste PushAsync-Lauf
|
||||
// versucht es erneut, jetzt mit dem soeben aktualisierten BasedOnServerSeq.
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<(int Pulled, int Conflicts)> PullAsync()
|
||||
{
|
||||
var since = _queue.GetLastServerSeq();
|
||||
|
||||
Reference in New Issue
Block a user