feat(localdb): lib-owned schema (meta/oplog/peer_state/row_version/applying/dead_letter)

This commit is contained in:
Joseph Doherty
2026-07-17 21:00:04 -04:00
parent cf58077398
commit fb8f751fac
2 changed files with 228 additions and 0 deletions
@@ -0,0 +1,105 @@
using Microsoft.Data.Sqlite;
namespace ZB.MOM.WW.LocalDb.Internal;
/// <summary>
/// Creates/migrates the library's own bookkeeping tables. Idempotent; runs in one transaction.
/// </summary>
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();
}
}
@@ -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<string>();
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);
}
}