feat(localdb): opt-in table registration + capture triggers (oplog + row_version, applying guard)
This commit is contained in:
@@ -0,0 +1,163 @@
|
||||
using Microsoft.Data.Sqlite;
|
||||
using ZB.MOM.WW.LocalDb.Internal;
|
||||
|
||||
namespace ZB.MOM.WW.LocalDb.Tests;
|
||||
|
||||
public sealed class CaptureTriggerTests : 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 SqliteLocalDb Create() => new(new LocalDbOptions { Path = _path });
|
||||
|
||||
private async Task<SqliteLocalDb> NewOrdersDb()
|
||||
{
|
||||
var db = Create();
|
||||
await db.ExecuteAsync("CREATE TABLE orders (id INTEGER PRIMARY KEY, sku TEXT, qty INTEGER)");
|
||||
db.RegisterReplicated("orders");
|
||||
return db;
|
||||
}
|
||||
|
||||
private sealed record OplogRow(long Seq, string PkJson, string? RowJson, long Hlc, string NodeId, long IsTombstone);
|
||||
|
||||
private Task<IReadOnlyList<OplogRow>> Oplog(SqliteLocalDb db) =>
|
||||
db.QueryAsync(
|
||||
"SELECT seq, pk_json, row_json, hlc, node_id, is_tombstone FROM __localdb_oplog WHERE table_name='orders' ORDER BY seq",
|
||||
r => new OplogRow(
|
||||
r.GetInt64(0), r.GetString(1), r.IsDBNull(2) ? null : r.GetString(2),
|
||||
r.GetInt64(3), r.GetString(4), r.GetInt64(5)));
|
||||
|
||||
[Fact]
|
||||
public async Task Insert_WritesOplogAndRowVersion_SameHlc()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
|
||||
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)");
|
||||
|
||||
var oplog = await Oplog(db);
|
||||
var row = Assert.Single(oplog);
|
||||
Assert.Equal("{\"id\":1}", row.PkJson);
|
||||
Assert.Equal("{\"id\":1,\"sku\":\"ABC\",\"qty\":5}", row.RowJson);
|
||||
Assert.Equal(0, row.IsTombstone);
|
||||
Assert.Equal(db.NodeId, row.NodeId);
|
||||
|
||||
var rvHlc = await db.QueryAsync(
|
||||
"SELECT hlc FROM __localdb_row_version WHERE table_name='orders' AND pk_json='{\"id\":1}'",
|
||||
r => r.GetInt64(0));
|
||||
Assert.Equal(row.Hlc, Assert.Single(rvHlc));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Update_WritesNewFullRow()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)");
|
||||
var insertHlc = (await Oplog(db)).Single().Hlc;
|
||||
|
||||
await db.ExecuteAsync("UPDATE orders SET qty = 9 WHERE id = 1");
|
||||
|
||||
var oplog = await Oplog(db);
|
||||
Assert.Equal(2, oplog.Count);
|
||||
var newest = oplog[^1];
|
||||
Assert.Equal("{\"id\":1,\"sku\":\"ABC\",\"qty\":9}", newest.RowJson);
|
||||
|
||||
var rvHlc = await db.QueryAsync(
|
||||
"SELECT hlc FROM __localdb_row_version WHERE table_name='orders' AND pk_json='{\"id\":1}'",
|
||||
r => r.GetInt64(0));
|
||||
Assert.Equal(newest.Hlc, Assert.Single(rvHlc));
|
||||
Assert.True(newest.Hlc > insertHlc);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Delete_WritesTombstone_NullRowJson_RowVersionTombstoned()
|
||||
{
|
||||
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 oplog = await Oplog(db);
|
||||
var tomb = oplog[^1];
|
||||
Assert.Null(tomb.RowJson);
|
||||
Assert.Equal(1, tomb.IsTombstone);
|
||||
|
||||
var rv = await db.QueryAsync(
|
||||
"SELECT is_tombstone, tombstone_utc FROM __localdb_row_version WHERE table_name='orders' AND pk_json='{\"id\":1}'",
|
||||
r => (Tomb: r.GetInt64(0), Utc: r.IsDBNull(1) ? null : r.GetString(1)));
|
||||
var v = Assert.Single(rv);
|
||||
Assert.Equal(1, v.Tomb);
|
||||
Assert.NotNull(v.Utc);
|
||||
Assert.EndsWith("Z", v.Utc);
|
||||
Assert.True(DateTimeOffset.TryParse(v.Utc, out _), $"tombstone_utc '{v.Utc}' should be ISO-8601");
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task MultiRowStatement_OneOplogEntryPerRow_DistinctHlcs()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
|
||||
await db.ExecuteAsync(
|
||||
"INSERT INTO orders (id, sku, qty) VALUES (1, 'A', 1), (2, 'B', 2), (3, 'C', 3)");
|
||||
|
||||
var oplog = await Oplog(db);
|
||||
Assert.Equal(3, oplog.Count);
|
||||
Assert.Equal(3, oplog.Select(o => o.Hlc).Distinct().Count());
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ApplyingFlag_SuppressesCapture()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
|
||||
await using (var tx = await db.BeginTransactionAsync())
|
||||
{
|
||||
await tx.ExecuteAsync("UPDATE __localdb_applying SET applying = 1 WHERE id = 1");
|
||||
await tx.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'A', 1)");
|
||||
await tx.ExecuteAsync("UPDATE __localdb_applying SET applying = 0 WHERE id = 1");
|
||||
await tx.CommitAsync(default);
|
||||
}
|
||||
|
||||
var oplog = await Oplog(db);
|
||||
Assert.Empty(oplog);
|
||||
|
||||
var rv = await db.QueryAsync("SELECT COUNT(*) FROM __localdb_row_version", r => r.GetInt64(0));
|
||||
Assert.Equal(0L, rv[0]);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task UnregisteredTable_NoCapture()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
await db.ExecuteAsync("CREATE TABLE other (id INTEGER PRIMARY KEY, v TEXT)");
|
||||
|
||||
await db.ExecuteAsync("INSERT INTO other (id, v) VALUES (1, 'x')");
|
||||
|
||||
var oplog = await db.QueryAsync(
|
||||
"SELECT COUNT(*) FROM __localdb_oplog WHERE table_name='other'", r => r.GetInt64(0));
|
||||
Assert.Equal(0L, oplog[0]);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Capture_InsideConsumerTransaction_AtomicWithRollback()
|
||||
{
|
||||
using var db = await NewOrdersDb();
|
||||
|
||||
await using (var tx = await db.BeginTransactionAsync())
|
||||
{
|
||||
await tx.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'A', 1)");
|
||||
await tx.RollbackAsync(default);
|
||||
}
|
||||
|
||||
var oplog = await Oplog(db);
|
||||
Assert.Empty(oplog);
|
||||
|
||||
var orders = await db.QueryAsync("SELECT COUNT(*) FROM orders", r => r.GetInt64(0));
|
||||
Assert.Equal(0L, orders[0]);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
using Microsoft.Data.Sqlite;
|
||||
using ZB.MOM.WW.LocalDb.Internal;
|
||||
using ZB.MOM.WW.LocalDb.Registration;
|
||||
|
||||
namespace ZB.MOM.WW.LocalDb.Tests;
|
||||
|
||||
public sealed class RegistrationTests : 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 SqliteLocalDb Create() => new(new LocalDbOptions { Path = _path });
|
||||
|
||||
private static async Task CreateOrders(SqliteLocalDb db) =>
|
||||
await db.ExecuteAsync("CREATE TABLE orders (id INTEGER PRIMARY KEY, sku TEXT, qty INTEGER)");
|
||||
|
||||
private Task<IReadOnlyList<string>> Triggers(SqliteLocalDb db) =>
|
||||
db.QueryAsync(
|
||||
"SELECT name FROM sqlite_master WHERE type='trigger' AND tbl_name='orders' ORDER BY name",
|
||||
r => r.GetString(0));
|
||||
|
||||
[Fact]
|
||||
public async Task Register_TableWithPk_InstallsThreeTriggers()
|
||||
{
|
||||
using var db = Create();
|
||||
await CreateOrders(db);
|
||||
|
||||
db.RegisterReplicated("orders");
|
||||
|
||||
var triggers = await Triggers(db);
|
||||
Assert.Equal(
|
||||
new[] { "__localdb_orders_ad", "__localdb_orders_ai", "__localdb_orders_au" },
|
||||
triggers);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Register_NoPk_Throws()
|
||||
{
|
||||
using var db = Create();
|
||||
await db.ExecuteAsync("CREATE TABLE nopk (a TEXT, b TEXT)");
|
||||
|
||||
var ex = Assert.Throws<LocalDbRegistrationException>(() => db.RegisterReplicated("nopk"));
|
||||
Assert.Contains("nopk", ex.Message);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void Register_MissingTable_Throws()
|
||||
{
|
||||
using var db = Create();
|
||||
|
||||
var ex = Assert.Throws<LocalDbRegistrationException>(() => db.RegisterReplicated("ghost"));
|
||||
Assert.Contains("ghost", ex.Message);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Register_IsIdempotent_RegeneratesTriggers()
|
||||
{
|
||||
using var db = Create();
|
||||
await CreateOrders(db);
|
||||
|
||||
db.RegisterReplicated("orders");
|
||||
db.RegisterReplicated("orders");
|
||||
|
||||
var triggers = await Triggers(db);
|
||||
Assert.Equal(3, triggers.Count);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Register_CompositePk_Works()
|
||||
{
|
||||
using var db = Create();
|
||||
await db.ExecuteAsync("CREATE TABLE kv (a INTEGER, b INTEGER, v TEXT, PRIMARY KEY (a, b))");
|
||||
|
||||
db.RegisterReplicated("kv");
|
||||
await db.ExecuteAsync("INSERT INTO kv (a, b, v) VALUES (1, 2, 'x')");
|
||||
|
||||
var pk = await db.QueryAsync(
|
||||
"SELECT pk_json FROM __localdb_oplog WHERE table_name='kv'",
|
||||
r => r.GetString(0));
|
||||
Assert.Single(pk);
|
||||
Assert.Equal("{\"a\":1,\"b\":2}", pk[0]);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Digest_ChangesWhenSchemaChanges()
|
||||
{
|
||||
using var db = Create();
|
||||
await CreateOrders(db);
|
||||
|
||||
db.RegisterReplicated("orders");
|
||||
var before = db.ReplicatedTables["orders"].Digest;
|
||||
|
||||
await db.ExecuteAsync("ALTER TABLE orders ADD COLUMN note TEXT");
|
||||
db.RegisterReplicated("orders");
|
||||
var after = db.ReplicatedTables["orders"].Digest;
|
||||
|
||||
Assert.NotEqual(before, after);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ReplicatedTables_ExposesRegistry()
|
||||
{
|
||||
using var db = Create();
|
||||
await CreateOrders(db);
|
||||
|
||||
db.RegisterReplicated("orders");
|
||||
|
||||
var table = db.ReplicatedTables["orders"];
|
||||
Assert.Equal("orders", table.Name);
|
||||
Assert.Equal(new[] { "id" }, table.PkColumns);
|
||||
Assert.Equal(new[] { "id", "sku", "qty" }, table.Columns.Select(c => c.Name).ToArray());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user