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)); } }