From fb8f751fac8fac8f4f9519a17600dd8d63789e07 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Fri, 17 Jul 2026 21:00:04 -0400 Subject: [PATCH] feat(localdb): lib-owned schema (meta/oplog/peer_state/row_version/applying/dead_letter) --- .../Internal/LocalDbSchema.cs | 105 +++++++++++++++ .../LocalDbSchemaTests.cs | 123 ++++++++++++++++++ 2 files changed, 228 insertions(+) create mode 100644 ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/LocalDbSchema.cs create mode 100644 ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LocalDbSchemaTests.cs diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/LocalDbSchema.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/LocalDbSchema.cs new file mode 100644 index 0000000..b0b5528 --- /dev/null +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/LocalDbSchema.cs @@ -0,0 +1,105 @@ +using Microsoft.Data.Sqlite; + +namespace ZB.MOM.WW.LocalDb.Internal; + +/// +/// Creates/migrates the library's own bookkeeping tables. Idempotent; runs in one transaction. +/// +internal static class LocalDbSchema +{ + private const string Ddl = """ + CREATE TABLE IF NOT EXISTS __localdb_meta ( + id INTEGER PRIMARY KEY CHECK (id = 1), + node_id TEXT NOT NULL, + hlc_high_water INTEGER NOT NULL DEFAULT 0, + schema_version INTEGER NOT NULL + ); + CREATE TABLE IF NOT EXISTS __localdb_oplog ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + row_json TEXT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + is_tombstone INTEGER NOT NULL DEFAULT 0 + ); + CREATE INDEX IF NOT EXISTS __localdb_oplog_hlc ON __localdb_oplog(table_name, pk_json, hlc); + CREATE TABLE IF NOT EXISTS __localdb_peer_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + peer_node_id TEXT NULL, + last_acked_seq INTEGER NOT NULL DEFAULT 0, + last_applied_remote_seq INTEGER NOT NULL DEFAULT 0, + last_seen_hlc INTEGER NOT NULL DEFAULT 0, + needs_snapshot INTEGER NOT NULL DEFAULT 0, + last_sync_utc TEXT NULL + ); + -- WITHOUT ROWID: the (table_name, pk_json) key IS the lookup path for every LWW conflict compare. + CREATE TABLE IF NOT EXISTS __localdb_row_version ( + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + is_tombstone INTEGER NOT NULL DEFAULT 0, + tombstone_utc TEXT NULL, + PRIMARY KEY (table_name, pk_json) + ) WITHOUT ROWID; + -- Single-row DB-scoped guard the applier toggles inside its txn so capture triggers skip replicated applies. + CREATE TABLE IF NOT EXISTS __localdb_applying ( + id INTEGER PRIMARY KEY CHECK (id = 1), + applying INTEGER NOT NULL DEFAULT 0 + ); + CREATE TABLE IF NOT EXISTS __localdb_dead_letter ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + received_utc TEXT NOT NULL, + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + row_json TEXT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + error TEXT NOT NULL + ); + INSERT INTO __localdb_applying (id, applying) SELECT 1, 0 WHERE NOT EXISTS (SELECT 1 FROM __localdb_applying); + INSERT INTO __localdb_peer_state (id) SELECT 1 WHERE NOT EXISTS (SELECT 1 FROM __localdb_peer_state); + """; + + public static void EnsureCreated(SqliteConnection connection) + { + using var transaction = connection.BeginTransaction(); + + using (var ddl = connection.CreateCommand()) + { + ddl.Transaction = transaction; + ddl.CommandText = Ddl; + ddl.ExecuteNonQuery(); + } + + using (var meta = connection.CreateCommand()) + { + meta.Transaction = transaction; + meta.CommandText = + "INSERT INTO __localdb_meta (id, node_id, schema_version) " + + "SELECT 1, $nodeId, 1 WHERE NOT EXISTS (SELECT 1 FROM __localdb_meta)"; + meta.Parameters.AddWithValue("$nodeId", Guid.NewGuid().ToString("D")); + meta.ExecuteNonQuery(); + } + + transaction.Commit(); + } + + internal static (string NodeId, long HlcHighWater) ReadMeta(SqliteConnection connection) + { + using var cmd = connection.CreateCommand(); + cmd.CommandText = "SELECT node_id, hlc_high_water FROM __localdb_meta WHERE id = 1"; + using var reader = cmd.ExecuteReader(); + reader.Read(); + return (reader.GetString(0), reader.GetInt64(1)); + } + + internal static void FlushHlcHighWater(SqliteConnection connection, long value) + { + using var cmd = connection.CreateCommand(); + cmd.CommandText = "UPDATE __localdb_meta SET hlc_high_water = $value WHERE id = 1"; + cmd.Parameters.AddWithValue("$value", value); + cmd.ExecuteNonQuery(); + } +} diff --git a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LocalDbSchemaTests.cs b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LocalDbSchemaTests.cs new file mode 100644 index 0000000..55598e8 --- /dev/null +++ b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LocalDbSchemaTests.cs @@ -0,0 +1,123 @@ +using Microsoft.Data.Sqlite; +using ZB.MOM.WW.LocalDb.Internal; + +namespace ZB.MOM.WW.LocalDb.Tests; + +public sealed class LocalDbSchemaTests : 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 SqliteConnection Open() + { + var conn = new SqliteConnection($"Data Source={_path}"); + conn.Open(); + return conn; + } + + private static long ScalarLong(SqliteConnection conn, string sql) + { + using var cmd = conn.CreateCommand(); + cmd.CommandText = sql; + return Convert.ToInt64(cmd.ExecuteScalar()); + } + + private static string ScalarString(SqliteConnection conn, string sql) + { + using var cmd = conn.CreateCommand(); + cmd.CommandText = sql; + return (string)cmd.ExecuteScalar()!; + } + + [Fact] + public void EnsureCreated_CreatesAllObjects() + { + using var conn = Open(); + LocalDbSchema.EnsureCreated(conn); + + using var cmd = conn.CreateCommand(); + cmd.CommandText = "SELECT name FROM sqlite_master WHERE type='table'"; + var tables = new List(); + using (var reader = cmd.ExecuteReader()) + while (reader.Read()) + tables.Add(reader.GetString(0)); + + foreach (var expected in new[] + { + "__localdb_meta", "__localdb_oplog", "__localdb_peer_state", + "__localdb_row_version", "__localdb_applying", "__localdb_dead_letter" + }) + Assert.Contains(expected, tables); + } + + [Fact] + public void EnsureCreated_IsIdempotent() + { + using var conn = Open(); + LocalDbSchema.EnsureCreated(conn); + LocalDbSchema.EnsureCreated(conn); + + Assert.Equal(1, ScalarLong(conn, "SELECT COUNT(*) FROM __localdb_meta")); + Assert.Equal(1, ScalarLong(conn, "SELECT COUNT(*) FROM __localdb_applying")); + Assert.Equal(1, ScalarLong(conn, "SELECT COUNT(*) FROM __localdb_peer_state")); + } + + [Fact] + public void Meta_NodeId_MintedOnceAndStable() + { + string first; + using (var conn = Open()) + { + LocalDbSchema.EnsureCreated(conn); + first = ScalarString(conn, "SELECT node_id FROM __localdb_meta WHERE id=1"); + } + + using (var conn = Open()) + { + LocalDbSchema.EnsureCreated(conn); + var second = ScalarString(conn, "SELECT node_id FROM __localdb_meta WHERE id=1"); + Assert.Equal(first, second); + } + + Assert.True(Guid.TryParse(first, out _)); + } + + [Fact] + public void Applying_SeededWithZeroRow() + { + using var conn = Open(); + LocalDbSchema.EnsureCreated(conn); + + Assert.Equal(0, ScalarLong(conn, "SELECT applying FROM __localdb_applying WHERE id=1")); + } + + [Fact] + public void ReadMeta_ReturnsNodeIdAndHighWater() + { + using var conn = Open(); + LocalDbSchema.EnsureCreated(conn); + + var (nodeId, highWater) = LocalDbSchema.ReadMeta(conn); + + Assert.True(Guid.TryParse(nodeId, out _)); + Assert.Equal(0, highWater); + } + + [Fact] + public void FlushHlcHighWater_RoundTrips() + { + using var conn = Open(); + LocalDbSchema.EnsureCreated(conn); + + LocalDbSchema.FlushHlcHighWater(conn, 123_456_789L); + + var (_, highWater) = LocalDbSchema.ReadMeta(conn); + Assert.Equal(123_456_789L, highWater); + } +}