Files
scadaproj/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/CaptureTriggerTests.cs
T

369 lines
15 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 Update_ChangingPk_TombstonesOldPk_CapturesNewPk()
{
using var db = await NewOrdersDb();
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)");
await db.ExecuteAsync("UPDATE orders SET id = 2 WHERE id = 1");
var oplog = await Oplog(db);
Assert.Equal(3, oplog.Count); // insert + (tombstone old, new row)
var tomb = oplog[1];
Assert.Equal("{\"id\":1}", tomb.PkJson);
Assert.Null(tomb.RowJson);
Assert.Equal(1, tomb.IsTombstone);
var moved = oplog[2];
Assert.Equal("{\"id\":2}", moved.PkJson);
Assert.Equal("{\"id\":2,\"sku\":\"ABC\",\"qty\":5}", moved.RowJson);
Assert.Equal(0, moved.IsTombstone);
Assert.NotEqual(tomb.Hlc, moved.Hlc);
var rv = await db.QueryAsync(
"SELECT pk_json, hlc, is_tombstone FROM __localdb_row_version WHERE table_name='orders' ORDER BY pk_json",
r => (PkJson: r.GetString(0), Hlc: r.GetInt64(1), Tomb: r.GetInt64(2)));
Assert.Equal(2, rv.Count);
var oldRv = rv.Single(v => v.PkJson == "{\"id\":1}");
Assert.Equal(1, oldRv.Tomb);
Assert.Equal(tomb.Hlc, oldRv.Hlc);
var newRv = rv.Single(v => v.PkJson == "{\"id\":2}");
Assert.Equal(0, newRv.Tomb);
Assert.Equal(moved.Hlc, newRv.Hlc);
}
[Fact]
public async Task Update_NotChangingPk_NoSpuriousTombstone()
{
using var db = await NewOrdersDb();
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)");
await db.ExecuteAsync("UPDATE orders SET qty = 9 WHERE id = 1");
var oplog = await Oplog(db);
Assert.Equal(2, oplog.Count);
Assert.All(oplog, o => Assert.Equal(0, o.IsTombstone));
var rv = await db.QueryAsync(
"SELECT pk_json, is_tombstone FROM __localdb_row_version WHERE table_name='orders'",
r => (PkJson: r.GetString(0), Tomb: r.GetInt64(1)));
var v = Assert.Single(rv);
Assert.Equal("{\"id\":1}", v.PkJson);
Assert.Equal(0, v.Tomb);
}
// Load-bearing for the S2 stale-rowid guard: on a pinned transaction connection,
// last_insert_rowid() is a LIVE oplog seq when the non-PK update's tombstone insert does not
// fire; without the AND(pk-changed) guard, S2 would tombstone an unrelated row's row_version.
[Fact]
public async Task Update_NonPkUpdate_InSameTransaction_DoesNotCorruptOtherRowVersions()
{
using var db = await NewOrdersDb();
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'B', 1)");
await using (var tx = await db.BeginTransactionAsync())
{
await tx.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (2, 'A', 2)");
await tx.ExecuteAsync("UPDATE orders SET qty = 9 WHERE id = 1");
await tx.CommitAsync(default);
}
var oplog = await Oplog(db);
Assert.Equal(3, oplog.Count);
Assert.All(oplog, o => Assert.Equal(0, o.IsTombstone));
var aInsertHlc = oplog.Single(o => o.PkJson == "{\"id\":2}").Hlc;
var rv = await db.QueryAsync(
"SELECT pk_json, hlc, is_tombstone FROM __localdb_row_version WHERE table_name='orders'",
r => (PkJson: r.GetString(0), Hlc: r.GetInt64(1), Tomb: r.GetInt64(2)));
Assert.Equal(2, rv.Count);
Assert.All(rv, v => Assert.Equal(0, v.Tomb));
var a = rv.Single(v => v.PkJson == "{\"id\":2}");
Assert.Equal(aInsertHlc, a.Hlc);
}
// The spectator row (id=3) pins per-row guard re-evaluation: for the unchanged-PK row, the
// stale last_insert_rowid() is the spectator's live oplog seq — an unguarded S2 would
// tombstone the spectator's row_version and nothing later would repair it.
[Fact]
public async Task Update_MultiRow_MixedPkChange_PerRowConditionHolds()
{
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.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (2, 'B', 2)");
await tx.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (3, 'C', 3)");
await tx.ExecuteAsync(
"UPDATE orders SET id = CASE WHEN id = 1 THEN 10 ELSE id END, qty = qty + 1 WHERE id IN (1, 2)");
await tx.CommitAsync(default);
}
var oplog = await Oplog(db);
var tomb = Assert.Single(oplog, o => o.IsTombstone == 1);
Assert.Equal("{\"id\":1}", tomb.PkJson);
var rv = await db.QueryAsync(
"SELECT pk_json, hlc, is_tombstone FROM __localdb_row_version WHERE table_name='orders'",
r => (PkJson: r.GetString(0), Hlc: r.GetInt64(1), Tomb: r.GetInt64(2)));
Assert.Equal(4, rv.Count);
var oldPk = rv.Single(v => v.PkJson == "{\"id\":1}");
Assert.Equal(1, oldPk.Tomb);
Assert.Equal(tomb.Hlc, oldPk.Hlc);
foreach (var pkJson in new[] { "{\"id\":2}", "{\"id\":3}", "{\"id\":10}" })
{
var v = rv.Single(x => x.PkJson == pkJson);
Assert.Equal(0, v.Tomb);
var newest = oplog.Where(o => o.PkJson == pkJson).Max(o => o.Hlc);
Assert.Equal(newest, v.Hlc);
}
}
[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());
var rv = await db.QueryAsync(
"SELECT pk_json, hlc FROM __localdb_row_version WHERE table_name='orders'",
r => (PkJson: r.GetString(0), Hlc: r.GetInt64(1)));
Assert.Equal(3, rv.Count);
foreach (var o in oplog)
Assert.Equal(o.Hlc, rv.Single(v => v.PkJson == o.PkJson).Hlc);
}
// Pins the explicit ON CONFLICT upsert form in the row_version capture statements: SQLite
// replaces a trigger-body statement's OR REPLACE with the OUTER statement's conflict handling,
// so a consumer `INSERT ... ON CONFLICT DO UPDATE` used to turn the au trigger's OR REPLACE
// into a plain INSERT -> SQLITE_CONSTRAINT_PRIMARYKEY (1555) on the existing row_version row.
[Fact]
public async Task Upsert_OnConflictDoUpdate_BothBranches_Capture()
{
using var db = await NewOrdersDb();
const string upsert =
"INSERT INTO orders (id, sku, qty) VALUES (7, @sku, @qty) " +
"ON CONFLICT(id) DO UPDATE SET sku = excluded.sku, qty = excluded.qty";
await db.ExecuteAsync(upsert, new { sku = "FIRST", qty = 1 }); // insert branch -> ai
await db.ExecuteAsync(upsert, new { sku = "SECOND", qty = 2 }); // update branch -> au
var oplog = await Oplog(db);
Assert.Equal(2, oplog.Count);
Assert.Equal("{\"id\":7,\"sku\":\"FIRST\",\"qty\":1}", oplog[0].RowJson);
Assert.Equal("{\"id\":7,\"sku\":\"SECOND\",\"qty\":2}", oplog[1].RowJson);
Assert.All(oplog, o => Assert.Equal(0, o.IsTombstone));
var rv = await db.QueryAsync(
"SELECT pk_json, hlc, is_tombstone FROM __localdb_row_version WHERE table_name='orders'",
r => (PkJson: r.GetString(0), Hlc: r.GetInt64(1), Tomb: r.GetInt64(2)));
var v = Assert.Single(rv);
Assert.Equal("{\"id\":7}", v.PkJson);
Assert.Equal(0, v.Tomb);
Assert.Equal(oplog[1].Hlc, v.Hlc);
}
[Fact]
public async Task Capture_CoexistsWithConsumerTrigger()
{
using var db = await NewOrdersDb();
await db.ExecuteAsync("CREATE TABLE audit (id INTEGER PRIMARY KEY AUTOINCREMENT, order_id INTEGER)");
await db.ExecuteAsync(
"CREATE TRIGGER consumer_ai AFTER INSERT ON orders BEGIN " +
"INSERT INTO audit (order_id) VALUES (NEW.id); END");
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (1, 'ABC', 5)");
var audit = await db.QueryAsync("SELECT order_id FROM audit", r => r.GetInt64(0));
Assert.Equal(1L, Assert.Single(audit));
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);
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 Capture_ValuesWithQuotesAndUnicode()
{
using var db = await NewOrdersDb();
const string sku = "O'Brien \"α\" ☃ \\slash";
await db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (@id, @sku, @qty)",
new { id = 1, sku, qty = 5 });
var roundTripped = await db.QueryAsync(
"SELECT json_extract(row_json, '$.sku') FROM __localdb_oplog WHERE table_name='orders'",
r => r.GetString(0));
Assert.Equal(sku, Assert.Single(roundTripped));
}
[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]);
}
}