3cf8cabaf9
Claude-Session: https://claude.ai/code/session_01BL2Vu1ESDQ9SCN4gVKkdts
227 lines
8.5 KiB
C#
227 lines
8.5 KiB
C#
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<SqliteLocalDb> 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 RecordPeerAck_OutOfOrderAck_DoesNotRegressWatermark()
|
|
{
|
|
using var db = await NewOrdersDb();
|
|
await WriteRows(db, 5);
|
|
var store = new OplogStore(db, new ReplicationOptions());
|
|
var batch = await store.ReadBatchAboveAsync(0, 100, default);
|
|
|
|
await store.RecordPeerAckAsync(batch[4].Seq, default);
|
|
await store.RecordPeerAckAsync(batch[1].Seq, default);
|
|
|
|
var peer = await store.GetPeerStateAsync(default);
|
|
Assert.Equal(batch[4].Seq, peer.LastAckedSeq);
|
|
// A late duplicate/out-of-order ack of an already-acked seq is benign, not a violation.
|
|
Assert.False(peer.NeedsSnapshot);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RecordPeerAck_AckBeyondMaxSeq_ClampsAndFlagsSnapshot()
|
|
{
|
|
using var db = await NewOrdersDb();
|
|
await WriteRows(db, 3);
|
|
var store = new OplogStore(db, new ReplicationOptions());
|
|
var maxSeq = await store.GetMaxSeqAsync(default);
|
|
|
|
await store.RecordPeerAckAsync(maxSeq + 100, default);
|
|
|
|
var peer = await store.GetPeerStateAsync(default);
|
|
Assert.Equal(maxSeq, peer.LastAckedSeq);
|
|
Assert.True(peer.NeedsSnapshot);
|
|
}
|
|
|
|
[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.Equal(3, count[0]);
|
|
Assert.True((await store.GetPeerStateAsync(default)).NeedsSnapshot);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task CapExceeded_ByAge_SetsNeedsSnapshot_AndPrunesExpired()
|
|
{
|
|
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);
|
|
// Age-only trigger must actually prune: both entries are past the age cutoff.
|
|
var count = await db.QueryAsync("SELECT COUNT(*) FROM __localdb_oplog", r => r.GetInt64(0));
|
|
Assert.Equal(0, count[0]);
|
|
}
|
|
|
|
[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));
|
|
}
|
|
}
|