feat(localdb): oplog store (batch read, ack watermark, pruning caps, tombstone retention)

Claude-Session: https://claude.ai/code/session_01BL2Vu1ESDQ9SCN4gVKkdts
This commit is contained in:
Joseph Doherty
2026-07-17 21:45:20 -04:00
parent c061ba733f
commit 529585163c
3 changed files with 344 additions and 0 deletions
@@ -0,0 +1,116 @@
using System.Globalization;
using ZB.MOM.WW.LocalDb.Hlc;
namespace ZB.MOM.WW.LocalDb.Replication.Internal;
/// <summary>A single oplog delta as read for outbound replication. <see cref="RowJson"/> is null for tombstones.</summary>
internal sealed record OplogEntryRecord(
long Seq, string TableName, string PkJson, string? RowJson, long Hlc, string NodeId, bool IsTombstone);
/// <summary>The persisted peer-sync watermarks and flags from <c>__localdb_peer_state</c>.</summary>
internal sealed record PeerState(
string? PeerNodeId, long LastAckedSeq, long LastAppliedRemoteSeq, long LastSeenHlc, bool NeedsSnapshot, string? LastSyncUtc);
/// <summary>
/// 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.
/// </summary>
internal sealed class OplogStore(ILocalDb db, ReplicationOptions options, Func<DateTimeOffset>? 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<DateTimeOffset> _utcNow = utcNow ?? (() => DateTimeOffset.UtcNow);
public Task<IReadOnlyList<OplogEntryRecord>> 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<long> 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<long> 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<bool> 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<PeerState> 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);
}
@@ -0,0 +1,37 @@
namespace ZB.MOM.WW.LocalDb.Replication;
/// <summary>
/// Configuration for the bidirectional 2-node replication engine. All durations are wall-clock UTC.
/// </summary>
public sealed class ReplicationOptions
{
/// <summary>gRPC address of the peer to dial. Required for the initiator role (validated in the sync host).</summary>
public string PeerAddress { get; set; } = "";
/// <summary>API key presented to / expected from the peer; null disables key auth.</summary>
public string? ApiKey { get; set; }
/// <summary>How often the change pump flushes pending oplog deltas to the peer.</summary>
public TimeSpan FlushInterval { get; set; } = TimeSpan.FromMilliseconds(250);
/// <summary>Maximum oplog entries sent per outbound batch.</summary>
public int MaxBatchSize { get; set; } = 500;
/// <summary>Backlog ceiling: exceeding it flags a snapshot resync and prunes the oldest oplog rows.</summary>
public long MaxOplogRows { get; set; } = 1_000_000;
/// <summary>Age ceiling: an unpruned oplog entry older than this flags a snapshot resync.</summary>
public TimeSpan MaxOplogAge { get; set; } = TimeSpan.FromDays(7);
/// <summary>How long deleted-row tombstones are retained before pruning.</summary>
public TimeSpan TombstoneRetention { get; set; } = TimeSpan.FromDays(7);
/// <summary>Upper bound on the exponential reconnect backoff.</summary>
public TimeSpan ReconnectBackoffMax { get; set; } = TimeSpan.FromSeconds(60);
/// <summary>Maximum tolerated HLC drift where a peer's clock runs ahead of ours.</summary>
public TimeSpan MaxHlcDriftAhead { get; set; } = TimeSpan.FromMinutes(5);
/// <summary>When true, drift beyond <see cref="MaxHlcDriftAhead"/> rejects the batch instead of clamping.</summary>
public bool FailClosedOnDrift { get; set; }
}
@@ -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<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 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));
}
}