diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/OplogStore.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/OplogStore.cs new file mode 100644 index 0000000..76749b9 --- /dev/null +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/OplogStore.cs @@ -0,0 +1,116 @@ +using System.Globalization; +using ZB.MOM.WW.LocalDb.Hlc; + +namespace ZB.MOM.WW.LocalDb.Replication.Internal; + +/// A single oplog delta as read for outbound replication. is null for tombstones. +internal sealed record OplogEntryRecord( + long Seq, string TableName, string PkJson, string? RowJson, long Hlc, string NodeId, bool IsTombstone); + +/// The persisted peer-sync watermarks and flags from __localdb_peer_state. +internal sealed record PeerState( + string? PeerNodeId, long LastAckedSeq, long LastAppliedRemoteSeq, long LastSeenHlc, bool NeedsSnapshot, string? LastSyncUtc); + +/// +/// DB-side operations over the replication bookkeeping tables: batched oplog reads above a +/// watermark, peer ack + prune, backlog/age cap enforcement, tombstone retention, and peer state. +/// +internal sealed class OplogStore(ILocalDb db, ReplicationOptions options, Func? utcNow = null) +{ + // Must match the delete trigger's strftime('%Y-%m-%dT%H:%M:%fZ') output so string comparison is chronological. + private const string TombstoneUtcFormat = "yyyy-MM-ddTHH:mm:ss.fffZ"; + + private readonly Func _utcNow = utcNow ?? (() => DateTimeOffset.UtcNow); + + public Task> ReadBatchAboveAsync(long afterSeq, int maxBatch, CancellationToken ct = default) => + db.QueryAsync( + "SELECT seq, table_name, pk_json, row_json, hlc, node_id, is_tombstone FROM __localdb_oplog " + + "WHERE seq > @afterSeq ORDER BY seq LIMIT @maxBatch", + static r => new OplogEntryRecord( + r.GetInt64(0), r.GetString(1), r.GetString(2), r.IsDBNull(3) ? null : r.GetString(3), + r.GetInt64(4), r.GetString(5), r.GetInt64(6) != 0), + new { afterSeq, maxBatch }, ct); + + public async Task RecordPeerAckAsync(long ackedSeq, CancellationToken ct = default) + { + await using var tx = await db.BeginTransactionAsync(ct); + await tx.ExecuteAsync( + "UPDATE __localdb_peer_state SET last_acked_seq = MAX(last_acked_seq, @acked), last_sync_utc = @now WHERE id = 1", + new { acked = ackedSeq, now = FormatUtc(_utcNow()) }, ct); + await tx.ExecuteAsync( + "DELETE FROM __localdb_oplog WHERE seq <= (SELECT last_acked_seq FROM __localdb_peer_state WHERE id = 1)", + null, ct); + await tx.CommitAsync(ct); + } + + public Task PruneTombstonesAsync(CancellationToken ct = default) => + db.ExecuteAsync( + "DELETE FROM __localdb_row_version WHERE is_tombstone = 1 AND tombstone_utc < @cutoff", + new { cutoff = FormatUtc(_utcNow() - options.TombstoneRetention) }, ct); + + public async Task GetOplogDepthAsync(CancellationToken ct = default) + { + var rows = await db.QueryAsync( + "SELECT COUNT(*) FROM __localdb_oplog WHERE seq > (SELECT last_acked_seq FROM __localdb_peer_state WHERE id = 1)", + static r => r.GetInt64(0), null, ct); + return rows[0]; + } + + public async Task GetMaxSeqAsync(CancellationToken ct = default) + { + var rows = await db.QueryAsync( + "SELECT COALESCE(MAX(seq), 0) FROM __localdb_oplog", static r => r.GetInt64(0), null, ct); + return rows[0]; + } + + public async Task EnforceCapsAsync(CancellationToken ct = default) + { + var overRows = await GetOplogDepthAsync(ct) > options.MaxOplogRows; + + var overAge = false; + var oldest = await db.QueryAsync( + "SELECT hlc FROM __localdb_oplog ORDER BY seq LIMIT 1", static r => r.GetInt64(0), null, ct); + if (oldest.Count > 0) + { + var cutoffMs = (_utcNow() - options.MaxOplogAge).ToUnixTimeMilliseconds(); + overAge = HybridLogicalClock.PhysicalMs(oldest[0]) < cutoffMs; + } + + if (!overRows && !overAge) + return false; + + await using var tx = await db.BeginTransactionAsync(ct); + await tx.ExecuteAsync("UPDATE __localdb_peer_state SET needs_snapshot = 1 WHERE id = 1", null, ct); + await tx.ExecuteAsync( + "DELETE FROM __localdb_oplog WHERE seq < ((SELECT MAX(seq) FROM __localdb_oplog) - @keep + 1)", + new { keep = options.MaxOplogRows }, ct); + await tx.CommitAsync(ct); + return true; + } + + public async Task GetPeerStateAsync(CancellationToken ct = default) + { + var rows = await db.QueryAsync( + "SELECT peer_node_id, last_acked_seq, last_applied_remote_seq, last_seen_hlc, needs_snapshot, last_sync_utc " + + "FROM __localdb_peer_state WHERE id = 1", + static r => new PeerState( + r.IsDBNull(0) ? null : r.GetString(0), r.GetInt64(1), r.GetInt64(2), r.GetInt64(3), + r.GetInt64(4) != 0, r.IsDBNull(5) ? null : r.GetString(5)), + null, ct); + return rows[0]; + } + + public Task SetPeerNodeIdAsync(string peerNodeId, CancellationToken ct = default) => + db.ExecuteAsync("UPDATE __localdb_peer_state SET peer_node_id = @v WHERE id = 1", new { v = peerNodeId }, ct); + + public Task SetNeedsSnapshotAsync(bool value, CancellationToken ct = default) => + db.ExecuteAsync("UPDATE __localdb_peer_state SET needs_snapshot = @v WHERE id = 1", new { v = value ? 1 : 0 }, ct); + + public Task SetLastAppliedRemoteSeqAsync(long seq, long lastSeenHlc, CancellationToken ct = default) => + db.ExecuteAsync( + "UPDATE __localdb_peer_state SET last_applied_remote_seq = @seq, last_seen_hlc = @hlc, last_sync_utc = @now WHERE id = 1", + new { seq, hlc = lastSeenHlc, now = FormatUtc(_utcNow()) }, ct); + + private static string FormatUtc(DateTimeOffset value) => + value.UtcDateTime.ToString(TombstoneUtcFormat, CultureInfo.InvariantCulture); +} diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/ReplicationOptions.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/ReplicationOptions.cs new file mode 100644 index 0000000..b1a6a8e --- /dev/null +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/ReplicationOptions.cs @@ -0,0 +1,37 @@ +namespace ZB.MOM.WW.LocalDb.Replication; + +/// +/// Configuration for the bidirectional 2-node replication engine. All durations are wall-clock UTC. +/// +public sealed class ReplicationOptions +{ + /// gRPC address of the peer to dial. Required for the initiator role (validated in the sync host). + public string PeerAddress { get; set; } = ""; + + /// API key presented to / expected from the peer; null disables key auth. + public string? ApiKey { get; set; } + + /// How often the change pump flushes pending oplog deltas to the peer. + public TimeSpan FlushInterval { get; set; } = TimeSpan.FromMilliseconds(250); + + /// Maximum oplog entries sent per outbound batch. + public int MaxBatchSize { get; set; } = 500; + + /// Backlog ceiling: exceeding it flags a snapshot resync and prunes the oldest oplog rows. + public long MaxOplogRows { get; set; } = 1_000_000; + + /// Age ceiling: an unpruned oplog entry older than this flags a snapshot resync. + public TimeSpan MaxOplogAge { get; set; } = TimeSpan.FromDays(7); + + /// How long deleted-row tombstones are retained before pruning. + public TimeSpan TombstoneRetention { get; set; } = TimeSpan.FromDays(7); + + /// Upper bound on the exponential reconnect backoff. + public TimeSpan ReconnectBackoffMax { get; set; } = TimeSpan.FromSeconds(60); + + /// Maximum tolerated HLC drift where a peer's clock runs ahead of ours. + public TimeSpan MaxHlcDriftAhead { get; set; } = TimeSpan.FromMinutes(5); + + /// When true, drift beyond rejects the batch instead of clamping. + public bool FailClosedOnDrift { get; set; } +} diff --git a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/OplogStoreTests.cs b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/OplogStoreTests.cs new file mode 100644 index 0000000..fa4fd82 --- /dev/null +++ b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/OplogStoreTests.cs @@ -0,0 +1,191 @@ +using Microsoft.Data.Sqlite; +using ZB.MOM.WW.LocalDb.Internal; +using ZB.MOM.WW.LocalDb.Replication; +using ZB.MOM.WW.LocalDb.Replication.Internal; + +namespace ZB.MOM.WW.LocalDb.Tests; + +public sealed class OplogStoreTests : IDisposable +{ + private readonly string _path = Path.Combine(Path.GetTempPath(), Guid.NewGuid() + ".db"); + + public void Dispose() + { + SqliteConnection.ClearAllPools(); + if (File.Exists(_path)) + File.Delete(_path); + } + + private async Task NewOrdersDb() + { + var db = new SqliteLocalDb(new LocalDbOptions { Path = _path }); + await db.ExecuteAsync("CREATE TABLE orders (id INTEGER PRIMARY KEY, sku TEXT, qty INTEGER)"); + db.RegisterReplicated("orders"); + return db; + } + + private async Task WriteRows(SqliteLocalDb db, int count) + { + for (var i = 1; i <= count; i++) + await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (@id, 'S', @id)", new { id = i }); + } + + [Fact] + public async Task ReadBatchAbove_ReturnsInSeqOrder_RespectsMaxBatch() + { + using var db = await NewOrdersDb(); + await WriteRows(db, 5); + var store = new OplogStore(db, new ReplicationOptions()); + + var first = await store.ReadBatchAboveAsync(0, 3, default); + Assert.Equal(3, first.Count); + Assert.True(first.Select(e => e.Seq).SequenceEqual(first.Select(e => e.Seq).OrderBy(x => x))); + + var rest = await store.ReadBatchAboveAsync(first[^1].Seq, 100, default); + Assert.Equal(2, rest.Count); + Assert.All(rest, e => Assert.True(e.Seq > first[^1].Seq)); + } + + [Fact] + public async Task ReadBatch_MapsAllFields() + { + using var db = await NewOrdersDb(); + await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)"); + await db.ExecuteAsync("DELETE FROM orders WHERE id = 1"); + var store = new OplogStore(db, new ReplicationOptions()); + + var batch = await store.ReadBatchAboveAsync(0, 100, default); + Assert.Equal(2, batch.Count); + + var ins = batch[0]; + Assert.Equal("orders", ins.TableName); + Assert.Equal("{\"id\":1}", ins.PkJson); + Assert.Equal("{\"id\":1,\"sku\":\"ABC\",\"qty\":5}", ins.RowJson); + Assert.Equal(db.NodeId, ins.NodeId); + Assert.False(ins.IsTombstone); + Assert.True(ins.Hlc > 0); + + var tomb = batch[1]; + Assert.Equal("{\"id\":1}", tomb.PkJson); + Assert.Null(tomb.RowJson); + Assert.True(tomb.IsTombstone); + } + + [Fact] + public async Task RecordPeerAck_PrunesBelowWatermark() + { + using var db = await NewOrdersDb(); + await WriteRows(db, 5); + var store = new OplogStore(db, new ReplicationOptions()); + var batch = await store.ReadBatchAboveAsync(0, 100, default); + var third = batch[2]; + + await store.RecordPeerAckAsync(third.Seq, default); + + var remaining = await store.ReadBatchAboveAsync(0, 100, default); + Assert.Equal(2, remaining.Count); + Assert.All(remaining, e => Assert.True(e.Seq > third.Seq)); + + var peer = await store.GetPeerStateAsync(default); + Assert.Equal(third.Seq, peer.LastAckedSeq); + } + + [Fact] + public async Task Prune_DropsExpiredTombstoneRowVersions_KeepsLive() + { + using var db = await NewOrdersDb(); + await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1,'A',1),(2,'B',2),(3,'C',3)"); + await db.ExecuteAsync("DELETE FROM orders WHERE id = 1"); // will be backdated -> expired + await db.ExecuteAsync("DELETE FROM orders WHERE id = 3"); // fresh tombstone -> kept + await db.ExecuteAsync( + "UPDATE __localdb_row_version SET tombstone_utc = '2000-01-01T00:00:00.000Z' WHERE pk_json = '{\"id\":1}'"); + + var store = new OplogStore(db, new ReplicationOptions()); + await store.PruneTombstonesAsync(default); + + var rv = await db.QueryAsync( + "SELECT pk_json, is_tombstone FROM __localdb_row_version WHERE table_name='orders'", + r => (Pk: r.GetString(0), Tomb: r.GetInt64(1))); + + Assert.DoesNotContain(rv, x => x.Pk == "{\"id\":1}"); + Assert.Contains(rv, x => x.Pk == "{\"id\":2}" && x.Tomb == 0); + Assert.Contains(rv, x => x.Pk == "{\"id\":3}" && x.Tomb == 1); + } + + [Fact] + public async Task OplogDepth_ReportsBacklog() + { + using var db = await NewOrdersDb(); + await WriteRows(db, 5); + var store = new OplogStore(db, new ReplicationOptions()); + Assert.Equal(5, await store.GetOplogDepthAsync(default)); + + var batch = await store.ReadBatchAboveAsync(0, 100, default); + await store.RecordPeerAckAsync(batch[1].Seq, default); // ack 2nd -> prunes 2, backlog 3 + + Assert.Equal(3, await store.GetOplogDepthAsync(default)); + } + + [Fact] + public async Task CapExceeded_SetsNeedsSnapshotAndPrunes() + { + using var db = await NewOrdersDb(); + await WriteRows(db, 5); + var store = new OplogStore(db, new ReplicationOptions { MaxOplogRows = 3 }); + + var flagged = await store.EnforceCapsAsync(default); + + Assert.True(flagged); + var count = await db.QueryAsync("SELECT COUNT(*) FROM __localdb_oplog", r => r.GetInt64(0)); + Assert.True(count[0] <= 3, $"oplog should be pruned to <= 3, was {count[0]}"); + Assert.True((await store.GetPeerStateAsync(default)).NeedsSnapshot); + } + + [Fact] + public async Task CapExceeded_ByAge_SetsNeedsSnapshot() + { + using var db = await NewOrdersDb(); + await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1,'A',1),(2,'B',2)"); + // Inject a clock 30 days ahead so the fresh entries exceed the 7d MaxOplogAge. + var future = DateTimeOffset.UtcNow.AddDays(30); + var store = new OplogStore(db, new ReplicationOptions(), () => future); + + var flagged = await store.EnforceCapsAsync(default); + + Assert.True(flagged); + Assert.True((await store.GetPeerStateAsync(default)).NeedsSnapshot); + } + + [Fact] + public async Task PeerState_RoundTrips() + { + using var db = await NewOrdersDb(); + var store = new OplogStore(db, new ReplicationOptions()); + + await store.SetPeerNodeIdAsync("node-xyz", default); + await store.SetNeedsSnapshotAsync(true, default); + await store.SetLastAppliedRemoteSeqAsync(42, 987654, default); + + var p = await store.GetPeerStateAsync(default); + Assert.Equal("node-xyz", p.PeerNodeId); + Assert.True(p.NeedsSnapshot); + Assert.Equal(42, p.LastAppliedRemoteSeq); + Assert.Equal(987654, p.LastSeenHlc); + Assert.NotNull(p.LastSyncUtc); + + await store.SetNeedsSnapshotAsync(false, default); + Assert.False((await store.GetPeerStateAsync(default)).NeedsSnapshot); + } + + [Fact] + public async Task GetMaxSeq_ReturnsHighestSeq_ZeroWhenEmpty() + { + using var db = await NewOrdersDb(); + var store = new OplogStore(db, new ReplicationOptions()); + Assert.Equal(0, await store.GetMaxSeqAsync(default)); + + await WriteRows(db, 4); + var batch = await store.ReadBatchAboveAsync(0, 100, default); + Assert.Equal(batch[^1].Seq, await store.GetMaxSeqAsync(default)); + } +}