diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs index 4b3c3d5..9c42ea8 100644 --- a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs @@ -29,7 +29,10 @@ internal sealed class SqliteLocalDb : ILocalDb, IDisposable private readonly string _connectionString; // Held open for the object lifetime: it anchors the WAL journal and is the connection // Dispose flushes the HLC high-water mark through, so the flush cannot fail on reopen. + // Guarded by _masterLock: SqliteConnection is not thread-safe, and registration/dispose + // may race from different threads. private readonly SqliteConnection _master; + private readonly Lock _masterLock = new(); private readonly HybridLogicalClock _clock; private readonly string _nodeId; private bool _disposed; @@ -62,6 +65,7 @@ internal sealed class SqliteLocalDb : ILocalDb, IDisposable "COALESCE((SELECT MAX(hlc) FROM __localdb_oplog),0)," + "COALESCE((SELECT MAX(hlc) FROM __localdb_row_version),0))"); _clock = new HybridLogicalClock(initial); + _master.CreateFunction("zb_hlc_next", () => _clock.Next()); (_nodeId, _) = LocalDbSchema.ReadMeta(_master); } @@ -135,48 +139,58 @@ internal sealed class SqliteLocalDb : ILocalDb, IDisposable public IReadOnlyDictionary ReplicatedTables => _replicatedTables; + // Thread-safe: all master-connection work runs under _masterLock. public ReplicatedTable RegisterReplicated(string tableName) { if (string.IsNullOrWhiteSpace(tableName) || !TableNameRegex.IsMatch(tableName)) throw new LocalDbRegistrationException( $"Cannot register table '{tableName}': name must be a plain identifier matching {TableNameRegex}."); - var columns = new List<(string Name, string Type, int Pk)>(); - using (var cmd = _master.CreateCommand()) - { - cmd.CommandText = $"PRAGMA table_info('{tableName}')"; - using var reader = cmd.ExecuteReader(); - while (reader.Read()) - // table_info cols: 0=cid, 1=name, 2=type, 3=notnull, 4=dflt_value, 5=pk (1-based PK ordinal, 0 if none). - columns.Add((reader.GetString(1), reader.IsDBNull(2) ? "" : reader.GetString(2), reader.GetInt32(5))); - } - - if (columns.Count == 0) - throw new LocalDbRegistrationException($"Cannot register table '{tableName}': it does not exist."); - - var pkColumns = columns.Where(c => c.Pk > 0).OrderBy(c => c.Pk).Select(c => c.Name).ToList(); - if (pkColumns.Count == 0) - throw new LocalDbRegistrationException( - $"Cannot register table '{tableName}': replication requires an explicit primary key."); - - var allColumns = columns.Select(c => (c.Name, c.Type)).ToList(); - var digest = ReplicatedTable.ComputeDigest(tableName, pkColumns, allColumns); - var table = new ReplicatedTable(tableName, pkColumns, allColumns, digest); - - var script = TriggerSqlGenerator.BuildInstallScript(table); - using (var tx = _master.BeginTransaction()) + lock (_masterLock) { + var columns = new List<(string Name, string Type, int Pk)>(); using (var cmd = _master.CreateCommand()) { - cmd.Transaction = tx; - cmd.CommandText = script; - cmd.ExecuteNonQuery(); + cmd.CommandText = $"PRAGMA table_info('{tableName}')"; + using var reader = cmd.ExecuteReader(); + while (reader.Read()) + // table_info cols: 0=cid, 1=name, 2=type, 3=notnull, 4=dflt_value, 5=pk (1-based PK ordinal, 0 if none). + columns.Add((reader.GetString(1), reader.IsDBNull(2) ? "" : reader.GetString(2), reader.GetInt32(5))); } - tx.Commit(); - } - _replicatedTables[tableName] = table; - return table; + if (columns.Count == 0) + throw new LocalDbRegistrationException($"Cannot register table '{tableName}': it does not exist."); + + var blob = columns.FirstOrDefault(c => c.Type.Contains("BLOB", StringComparison.OrdinalIgnoreCase)); + if (blob.Name is not null) + throw new LocalDbRegistrationException( + $"Cannot register table '{tableName}': column '{blob.Name}' is declared {blob.Type} — " + + "BLOB columns are not supported in replicated tables (json_object cannot capture them)."); + + var pkColumns = columns.Where(c => c.Pk > 0).OrderBy(c => c.Pk).Select(c => c.Name).ToList(); + if (pkColumns.Count == 0) + throw new LocalDbRegistrationException( + $"Cannot register table '{tableName}': replication requires an explicit primary key."); + + var allColumns = columns.Select(c => (c.Name, c.Type)).ToList(); + var digest = ReplicatedTable.ComputeDigest(tableName, pkColumns, allColumns); + var table = new ReplicatedTable(tableName, pkColumns, allColumns, digest); + + var script = TriggerSqlGenerator.BuildInstallScript(table); + using (var tx = _master.BeginTransaction()) + { + using (var cmd = _master.CreateCommand()) + { + cmd.Transaction = tx; + cmd.CommandText = script; + cmd.ExecuteNonQuery(); + } + tx.Commit(); + } + + _replicatedTables[tableName] = table; + return table; + } } public void Dispose() @@ -184,8 +198,11 @@ internal sealed class SqliteLocalDb : ILocalDb, IDisposable if (_disposed) return; _disposed = true; - LocalDbSchema.FlushHlcHighWater(_master, _clock.Current); - _master.Dispose(); + lock (_masterLock) + { + LocalDbSchema.FlushHlcHighWater(_master, _clock.Current); + _master.Dispose(); + } } private static async Task> ReadAllAsync(SqliteCommand cmd, Func map, CancellationToken ct) diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/ReplicatedTable.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/ReplicatedTable.cs index 6154ee3..1cd35f7 100644 --- a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/ReplicatedTable.cs +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/ReplicatedTable.cs @@ -1,5 +1,6 @@ using System.Security.Cryptography; using System.Text; +using System.Text.RegularExpressions; namespace ZB.MOM.WW.LocalDb.Registration; @@ -15,15 +16,24 @@ public sealed record ReplicatedTable( { /// /// Lowercase-hex SHA-256 of "{name}|pk:{p1,p2}|cols:{c1:TYPE,c2:TYPE,...}" — PK columns - /// in ordinal order, all columns in cid order with uppercased declared types. + /// in ordinal order, all columns in cid order. Every name/type component is length-prefixed + /// ({len}:{value}) so no combination of names/types can collide with a different schema, + /// and declared types have whitespace runs collapsed then uppercased so equivalent declarations + /// digest identically. /// internal static string ComputeDigest( string name, IReadOnlyList pkColumns, IReadOnlyList<(string Name, string Type)> columns) { - var pk = string.Join(",", pkColumns); - var cols = string.Join(",", columns.Select(c => $"{c.Name}:{c.Type.ToUpperInvariant()}")); - var canonical = $"{name}|pk:{pk}|cols:{cols}"; + var pk = string.Join(",", pkColumns.Select(LengthPrefixed)); + var cols = string.Join(",", columns.Select(c => + $"{LengthPrefixed(c.Name)}:{LengthPrefixed(NormalizeType(c.Type))}")); + var canonical = $"{LengthPrefixed(name)}|pk:{pk}|cols:{cols}"; var hash = SHA256.HashData(Encoding.UTF8.GetBytes(canonical)); return Convert.ToHexStringLower(hash); } + + private static string LengthPrefixed(string value) => $"{value.Length}:{value}"; + + private static string NormalizeType(string type) => + Regex.Replace(type.Trim(), @"\s+", " ").ToUpperInvariant(); } diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/TriggerSqlGenerator.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/TriggerSqlGenerator.cs index 39d78ee..2a1f32f 100644 --- a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/TriggerSqlGenerator.cs +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/TriggerSqlGenerator.cs @@ -5,12 +5,19 @@ namespace ZB.MOM.WW.LocalDb.Registration; /// /// Emits the three AFTER triggers (insert/update/delete) that capture every row change on a /// registered table into __localdb_oplog + __localdb_row_version within the -/// consumer's own write transaction. The WHEN applying = 0 guard skips replicated applies; -/// the last_insert_rowid() back-reference binds each row-version row to the SAME -/// zb_hlc_next() stamp as its oplog entry (one stamp per row change). +/// consumer's own write transaction. The WHEN COALESCE(applying, 0) = 0 guard skips +/// replicated applies (COALESCE so a missing guard row captures loudly rather than silently +/// disabling CDC); the last_insert_rowid() back-reference binds each row-version row to +/// the SAME zb_hlc_next() stamp as its oplog entry (one stamp per row change). A +/// PK-changing UPDATE additionally tombstones the OLD key first — delete+insert semantics. /// internal static class TriggerSqlGenerator { + private const string ApplyingGuard = + "COALESCE((SELECT applying FROM __localdb_applying WHERE id = 1), 0) = 0"; + + private const string TombstoneTail = "1, strftime('%Y-%m-%dT%H:%M:%fZ','now')"; + public static string TriggerName(string table, string suffix) => $"__localdb_{table}_{suffix}"; /// DROP IF EXISTS + CREATE for all three triggers, ready to run in one transaction. @@ -18,7 +25,7 @@ internal static class TriggerSqlGenerator { var sb = new StringBuilder(); foreach (var suffix in new[] { "ai", "au", "ad" }) - sb.Append("DROP TRIGGER IF EXISTS \"").Append(TriggerName(table.Name, suffix)).Append("\";\n"); + sb.Append("DROP TRIGGER IF EXISTS ").Append(Ident(TriggerName(table.Name, suffix))).Append(";\n"); sb.Append(BuildInsert(table)).Append('\n'); sb.Append(BuildUpdate(table)).Append('\n'); @@ -26,49 +33,80 @@ internal static class TriggerSqlGenerator return sb.ToString(); } - private static string BuildInsert(ReplicatedTable t) => Capture(t, "ai", "INSERT", "NEW", - rowJson: JsonObject("NEW", t.Columns.Select(c => c.Name)), isTombstone: "0", - rvTail: "0, NULL"); + private static string BuildInsert(ReplicatedTable t) => + Wrap(t, "ai", "INSERT", NewRowCapture(t)); - private static string BuildUpdate(ReplicatedTable t) => Capture(t, "au", "UPDATE", "NEW", - rowJson: JsonObject("NEW", t.Columns.Select(c => c.Name)), isTombstone: "0", - rvTail: "0, NULL"); + private static string BuildUpdate(ReplicatedTable t) => + Wrap(t, "au", "UPDATE", OldPkTombstoneIfChanged(t) + NewRowCapture(t)); - private static string BuildDelete(ReplicatedTable t) => Capture(t, "ad", "DELETE", "OLD", - rowJson: "NULL", isTombstone: "1", - rvTail: "1, strftime('%Y-%m-%dT%H:%M:%fZ','now')"); - - private static string Capture( - ReplicatedTable t, string suffix, string verb, string alias, - string rowJson, string isTombstone, string rvTail) - { - var name = TriggerName(t.Name, suffix); - var pkJson = JsonObject(alias, t.PkColumns); - var tableLit = Literal(t.Name); - return + private static string BuildDelete(ReplicatedTable t) => + Wrap(t, "ad", "DELETE", $""" - CREATE TRIGGER "{name}" AFTER {verb} ON "{t.Name}" - WHEN (SELECT applying FROM __localdb_applying WHERE id = 1) = 0 - BEGIN INSERT INTO __localdb_oplog (table_name, pk_json, row_json, hlc, node_id, is_tombstone) - VALUES ({tableLit}, - {pkJson}, - {rowJson}, + VALUES ({Literal(t.Name)}, + {JsonObject("OLD", t.PkColumns)}, + NULL, zb_hlc_next(), (SELECT node_id FROM __localdb_meta WHERE id = 1), - {isTombstone}); + 1); INSERT OR REPLACE INTO __localdb_row_version (table_name, pk_json, hlc, node_id, is_tombstone, tombstone_utc) - SELECT table_name, pk_json, hlc, node_id, {rvTail} FROM __localdb_oplog WHERE seq = last_insert_rowid(); - END; + SELECT table_name, pk_json, hlc, node_id, {TombstoneTail} FROM __localdb_oplog WHERE seq = last_insert_rowid(); + + """); + + private static string Wrap(ReplicatedTable t, string suffix, string verb, string body) => + $""" + CREATE TRIGGER {Ident(TriggerName(t.Name, suffix))} AFTER {verb} ON {Ident(t.Name)} + WHEN {ApplyingGuard} + BEGIN + {body.TrimEnd('\n')} + END; + """; + + private static string NewRowCapture(ReplicatedTable t) => + $""" + INSERT INTO __localdb_oplog (table_name, pk_json, row_json, hlc, node_id, is_tombstone) + VALUES ({Literal(t.Name)}, + {JsonObject("NEW", t.PkColumns)}, + {JsonObject("NEW", t.Columns.Select(c => c.Name))}, + zb_hlc_next(), + (SELECT node_id FROM __localdb_meta WHERE id = 1), + 0); + INSERT OR REPLACE INTO __localdb_row_version (table_name, pk_json, hlc, node_id, is_tombstone, tombstone_utc) + SELECT table_name, pk_json, hlc, node_id, 0, NULL FROM __localdb_oplog WHERE seq = last_insert_rowid(); + + """; + + // Both statements carry the SAME pk-changed condition: when the tombstone insert did not fire, + // the row-version statement must also be a no-op, or a stale last_insert_rowid would mis-correlate. + private static string OldPkTombstoneIfChanged(ReplicatedTable t) + { + var changed = string.Join(" OR ", t.PkColumns.Select(p => $"OLD.{Ident(p)} IS NOT NEW.{Ident(p)}")); + return + $""" + INSERT INTO __localdb_oplog (table_name, pk_json, row_json, hlc, node_id, is_tombstone) + SELECT {Literal(t.Name)}, + {JsonObject("OLD", t.PkColumns)}, + NULL, + zb_hlc_next(), + (SELECT node_id FROM __localdb_meta WHERE id = 1), + 1 + WHERE {changed}; + INSERT OR REPLACE INTO __localdb_row_version (table_name, pk_json, hlc, node_id, is_tombstone, tombstone_utc) + SELECT table_name, pk_json, hlc, node_id, {TombstoneTail} FROM __localdb_oplog + WHERE seq = last_insert_rowid() AND ({changed}); + """; } - // json_object('col', ALIAS."col", ...) — natively handles NULLs and SQLite types (BLOB columns error, a v1 limitation). + // json_object('col', ALIAS."col", ...) — natively handles NULLs and SQLite scalar types. private static string JsonObject(string alias, IEnumerable columns) { - var args = columns.Select(c => $"{Literal(c)}, {alias}.\"{c}\""); + var args = columns.Select(c => $"{Literal(c)}, {alias}.{Ident(c)}"); return $"json_object({string.Join(", ", args)})"; } + private static string Ident(string name) => "\"" + name.Replace("\"", "\"\"") + "\""; + private static string Literal(string value) => "'" + value.Replace("'", "''") + "'"; } diff --git a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/CaptureTriggerTests.cs b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/CaptureTriggerTests.cs index 98fb94f..5f6afed 100644 --- a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/CaptureTriggerTests.cs +++ b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/CaptureTriggerTests.cs @@ -97,6 +97,63 @@ public sealed class CaptureTriggerTests : IDisposable 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); + } + [Fact] public async Task MultiRowStatement_OneOplogEntryPerRow_DistinctHlcs() { @@ -108,6 +165,53 @@ public sealed class CaptureTriggerTests : IDisposable 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); + } + + [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] diff --git a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/RegistrationTests.cs b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/RegistrationTests.cs index 0c68608..f9b6118 100644 --- a/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/RegistrationTests.cs +++ b/ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/RegistrationTests.cs @@ -58,6 +58,18 @@ public sealed class RegistrationTests : IDisposable Assert.Contains("ghost", ex.Message); } + [Fact] + public async Task Register_BlobColumn_Throws() + { + using var db = Create(); + await db.ExecuteAsync("CREATE TABLE blobby (id INTEGER PRIMARY KEY, data BLOB)"); + + var ex = Assert.Throws(() => db.RegisterReplicated("blobby")); + Assert.Contains("blobby", ex.Message); + Assert.Contains("data", ex.Message); + Assert.Contains("BLOB", ex.Message); + } + [Fact] public async Task Register_IsIdempotent_RegeneratesTriggers() {